Add sentiment_engine with CryptoSentimentCalibrator fixes - improved keyword lists, lowered FinBERT threshold, added neutral handling
This commit is contained in:
11
sentiment_engine/prefect_flows/__init__.py
Normal file
11
sentiment_engine/prefect_flows/__init__.py
Normal file
@@ -0,0 +1,11 @@
|
||||
"""Prefect flows for scheduled connectors"""
|
||||
|
||||
from .connectors.rss_ingest import rss_ingest_flow
|
||||
from .connectors.api_ingest import api_ingest_flow
|
||||
from .connectors.web_crawl import web_crawl_flow
|
||||
|
||||
__all__ = [
|
||||
"rss_ingest_flow",
|
||||
"api_ingest_flow",
|
||||
"web_crawl_flow",
|
||||
]
|
||||
100
sentiment_engine/prefect_flows/connectors/api_ingest.py
Normal file
100
sentiment_engine/prefect_flows/connectors/api_ingest.py
Normal file
@@ -0,0 +1,100 @@
|
||||
"""API ingestion Prefect flow"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Dict, List, Optional
|
||||
|
||||
import aiohttp
|
||||
from prefect import flow, task
|
||||
from prefect.task_runners import ConcurrentTaskRunner
|
||||
|
||||
from sentiment_engine.utils.config import get_settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@task(retries=2, retry_delay_seconds=60)
|
||||
async def fetch_api_endpoint(
|
||||
url: str,
|
||||
headers: Dict[str, str] = None,
|
||||
params: Dict = None
|
||||
) -> List[dict]:
|
||||
"""Fetch a single API endpoint"""
|
||||
try:
|
||||
async with aiohttp.ClientSession() as session:
|
||||
async with session.get(url, headers=headers, params=params, timeout=30) as resp:
|
||||
if resp.status != 200:
|
||||
logger.warning(f"API {url} returned {resp.status}")
|
||||
return []
|
||||
data = await resp.json()
|
||||
|
||||
# Normalize to list of items
|
||||
items = data if isinstance(data, list) else [data]
|
||||
return items
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching {url}: {e}")
|
||||
raise
|
||||
|
||||
|
||||
@flow(
|
||||
name="api_ingest",
|
||||
task_runner=ConcurrentTaskRunner(max_workers=5),
|
||||
log_prints=True
|
||||
)
|
||||
async def api_ingest_flow():
|
||||
"""Main API ingestion flow for FRED, EDGAR, etc."""
|
||||
settings = get_settings()
|
||||
|
||||
endpoints = [
|
||||
{
|
||||
"name": "fred_vix",
|
||||
"url": "https://api.stlouisfed.org/fred/series/observations",
|
||||
"params": {
|
||||
"series_id": "VIXCLS",
|
||||
"api_key": "${FRED_API_KEY}",
|
||||
"file_type": "json",
|
||||
"limit": 1,
|
||||
"sort_order": "desc"
|
||||
},
|
||||
"source_type": "regulatory"
|
||||
},
|
||||
{
|
||||
"name": "fred_dxy",
|
||||
"url": "https://api.stlouisfed.org/fred/series/observations",
|
||||
"params": {
|
||||
"series_id": "DTWEXBGS",
|
||||
"api_key": "${FRED_API_KEY}",
|
||||
"file_type": "json",
|
||||
"limit": 1,
|
||||
"sort_order": "desc"
|
||||
},
|
||||
"source_type": "regulatory"
|
||||
},
|
||||
# Add more FRED series, EDGAR, etc.
|
||||
]
|
||||
|
||||
results = await asyncio.gather(
|
||||
*[fetch_api_endpoint(ep["url"], params=ep.get("params")) for ep in endpoints],
|
||||
return_exceptions=True
|
||||
)
|
||||
|
||||
all_items = []
|
||||
for i, result in enumerate(results):
|
||||
ep = endpoints[i]
|
||||
if isinstance(result, Exception):
|
||||
logger.error(f"Endpoint {ep['name']} failed: {result}")
|
||||
else:
|
||||
for item in result:
|
||||
all_items.append({
|
||||
"source_id": f"api:{ep['name']}",
|
||||
"source_type": ep["source_type"],
|
||||
"raw_text": str(item),
|
||||
"metadata": {"endpoint": ep["name"], "raw": item}
|
||||
})
|
||||
|
||||
logger.info(f"Fetched {len(all_items)} items from {len(endpoints)} API endpoints")
|
||||
return all_items
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(api_ingest_flow())
|
||||
99
sentiment_engine/prefect_flows/connectors/rss_ingest.py
Normal file
99
sentiment_engine/prefect_flows/connectors/rss_ingest.py
Normal file
@@ -0,0 +1,99 @@
|
||||
"""RSS ingestion Prefect flow"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import List
|
||||
|
||||
import feedparser
|
||||
from prefect import flow, task
|
||||
from prefect.task_runners import ConcurrentTaskRunner
|
||||
|
||||
from sentiment_engine.ingestion.rss import RSSConnector
|
||||
from sentiment_engine.schemas.config import RSSConnectorConfig
|
||||
from sentiment_engine.utils.config import get_settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@task(retries=3, retry_delay_seconds=30)
|
||||
async def fetch_rss_feed(feed_url: str, config: RSSConnectorConfig) -> List[dict]:
|
||||
"""Fetch and parse a single RSS feed"""
|
||||
try:
|
||||
feed = feedparser.parse(feed_url)
|
||||
items = []
|
||||
|
||||
for entry in feed.entries[:config.max_items_per_feed]:
|
||||
title = getattr(entry, "title", "").strip()
|
||||
summary = getattr(entry, "summary", getattr(entry, "description", "")).strip()
|
||||
raw_text = f"{title}\n\n{summary}"
|
||||
|
||||
if not raw_text.strip():
|
||||
continue
|
||||
|
||||
items.append({
|
||||
"source_id": f"rss:{feed_url}",
|
||||
"source_type": "news",
|
||||
"raw_text": raw_text,
|
||||
"title": title,
|
||||
"url": getattr(entry, "link", ""),
|
||||
"author": getattr(entry, "author", ""),
|
||||
"publish_ts": getattr(entry, "published_parsed", None),
|
||||
"metadata": {"feed_url": feed_url}
|
||||
})
|
||||
|
||||
return items
|
||||
except Exception as e:
|
||||
logger.error(f"Error fetching {feed_url}: {e}")
|
||||
raise
|
||||
|
||||
|
||||
@flow(
|
||||
name="rss_ingest",
|
||||
task_runner=ConcurrentTaskRunner(max_workers=10),
|
||||
log_prints=True
|
||||
)
|
||||
async def rss_ingest_flow(feed_urls: List[str] = None):
|
||||
"""Main RSS ingestion flow"""
|
||||
settings = get_settings()
|
||||
|
||||
if feed_urls is None:
|
||||
# Default crypto news feeds
|
||||
feed_urls = [
|
||||
"https://www.coindesk.com/arc/outboundfeeds/rss/",
|
||||
"https://cointelegraph.com/rss",
|
||||
"https://www.theblock.co/rss",
|
||||
"https://decrypt.co/feed",
|
||||
"https://messari.io/feed",
|
||||
"https://cryptoslate.com/feed/",
|
||||
"https://bitcoinmagazine.com/feed/",
|
||||
]
|
||||
|
||||
config = RSSConnectorConfig(
|
||||
name="prefect_rss",
|
||||
source_type="news",
|
||||
feed_urls=feed_urls,
|
||||
max_items_per_feed=50
|
||||
)
|
||||
|
||||
# Fetch all feeds concurrently
|
||||
results = await asyncio.gather(
|
||||
*[fetch_rss_feed(url, config) for url in feed_urls],
|
||||
return_exceptions=True
|
||||
)
|
||||
|
||||
all_items = []
|
||||
for i, result in enumerate(results):
|
||||
if isinstance(result, Exception):
|
||||
logger.error(f"Feed {feed_urls[i]} failed: {result}")
|
||||
else:
|
||||
all_items.extend(result)
|
||||
|
||||
logger.info(f"Fetched {len(all_items)} items from {len(feed_urls)} feeds")
|
||||
|
||||
# In production, publish to NATS
|
||||
# For now, return items
|
||||
return all_items
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(rss_ingest_flow())
|
||||
128
sentiment_engine/prefect_flows/connectors/web_crawl.py
Normal file
128
sentiment_engine/prefect_flows/connectors/web_crawl.py
Normal file
@@ -0,0 +1,128 @@
|
||||
"""Web crawl Prefect flow"""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import subprocess
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
from typing import List
|
||||
|
||||
from prefect import flow, task
|
||||
|
||||
from sentiment_engine.utils.config import get_settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@task(retries=1, retry_delay_seconds=300)
|
||||
async def run_hister_crawl(
|
||||
seed_urls: List[str],
|
||||
allowed_domains: List[str],
|
||||
max_depth: int = 2,
|
||||
job_timeout: int = 3600
|
||||
) -> List[dict]:
|
||||
"""Run Hister crawl job"""
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
seed_file = Path(tmpdir) / "seeds.txt"
|
||||
seed_file.write_text("\n".join(seed_urls))
|
||||
output_file = Path(tmpdir) / "output.jsonl"
|
||||
|
||||
cmd = [
|
||||
"hister", "crawl",
|
||||
"--input", str(seed_file),
|
||||
"--job-id", f"prefect-crawl-{asyncio.current_task().get_name()}",
|
||||
"--depth", str(max_depth),
|
||||
"--delay", "1.0",
|
||||
"--output", str(output_file),
|
||||
"--format", "jsonl"
|
||||
]
|
||||
|
||||
if allowed_domains:
|
||||
cmd.extend(["--allowed-domain", ",".join(allowed_domains)])
|
||||
|
||||
try:
|
||||
proc = await asyncio.create_subprocess_exec(
|
||||
*cmd,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE
|
||||
)
|
||||
|
||||
stdout, stderr = await asyncio.wait_for(
|
||||
proc.communicate(), timeout=job_timeout
|
||||
)
|
||||
|
||||
if proc.returncode != 0:
|
||||
logger.error(f"Hister failed: {stderr.decode()}")
|
||||
return []
|
||||
|
||||
# Parse output
|
||||
items = []
|
||||
if output_file.exists():
|
||||
import json
|
||||
with open(output_file) as f:
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(line)
|
||||
items.append(data)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
|
||||
return items
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
logger.error(f"Hister job timed out after {job_timeout}s")
|
||||
return []
|
||||
except FileNotFoundError:
|
||||
logger.error("Hister not installed")
|
||||
return []
|
||||
|
||||
|
||||
@flow(
|
||||
name="web_crawl",
|
||||
log_prints=True
|
||||
)
|
||||
async def web_crawl_flow(
|
||||
seed_urls: List[str] = None,
|
||||
allowed_domains: List[str] = None,
|
||||
max_depth: int = 2
|
||||
):
|
||||
"""Web crawl flow for sites without RSS/API"""
|
||||
if seed_urls is None:
|
||||
seed_urls = [
|
||||
"https://www.coindesk.com",
|
||||
"https://cointelegraph.com",
|
||||
"https://www.theblock.co",
|
||||
"https://decrypt.co",
|
||||
"https://cryptoslate.com",
|
||||
]
|
||||
|
||||
if allowed_domains is None:
|
||||
allowed_domains = [
|
||||
"coindesk.com", "cointelegraph.com", "theblock.co",
|
||||
"decrypt.co", "cryptoslate.com", "bitcoinmagazine.com"
|
||||
]
|
||||
|
||||
items = await run_hister_crawl(seed_urls, allowed_domains, max_depth)
|
||||
|
||||
logger.info(f"Crawled {len(items)} pages")
|
||||
|
||||
# Convert to normalized items
|
||||
normalized = []
|
||||
for item in items:
|
||||
normalized.append({
|
||||
"source_id": f"web:{item.get('url', '').split('/')[2] if item.get('url') else 'unknown'}",
|
||||
"source_type": "news",
|
||||
"raw_text": f"{item.get('title', '')}\n\n{item.get('content', item.get('text', ''))}",
|
||||
"title": item.get("title"),
|
||||
"url": item.get("url"),
|
||||
"metadata": {"crawler": "hister", "job": "prefect"}
|
||||
})
|
||||
|
||||
return normalized
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(web_crawl_flow())
|
||||
Reference in New Issue
Block a user