Files

100 lines
2.9 KiB
Python
Raw Permalink Normal View History

"""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())