199 lines
7.1 KiB
Python
199 lines
7.1 KiB
Python
"""
|
|
Article processor for handling the processing of articles from the scraper directory.
|
|
Manages batching, processing, and integration with the fact extraction system.
|
|
"""
|
|
|
|
import os
|
|
import json
|
|
import logging
|
|
import time
|
|
from datetime import datetime
|
|
from typing import List, Tuple
|
|
|
|
from fact_extractor import FactExtractor
|
|
from cache_manager import CacheManager
|
|
from metrics_collector import metrics_collector
|
|
from config import BATCH_SIZE, CACHE_FILE
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
class ArticleProcessor:
|
|
"""Processes articles from the scraper directory and extracts facts."""
|
|
|
|
def __init__(self):
|
|
self.fact_extractor = FactExtractor()
|
|
self.cache_manager = CacheManager(CACHE_FILE)
|
|
self.batch_size = BATCH_SIZE
|
|
|
|
def find_unprocessed_articles(self, scraper_dir: str = "/scraper/articles") -> List[Tuple[str, str]]:
|
|
"""
|
|
Find all unprocessed articles in the scraper directory.
|
|
|
|
Args:
|
|
scraper_dir (str): Path to the scraper articles directory
|
|
|
|
Returns:
|
|
List of tuples (file_path, filename)
|
|
"""
|
|
unprocessed_articles = []
|
|
|
|
try:
|
|
for root, dirs, files in os.walk(scraper_dir):
|
|
for file in files:
|
|
if file.endswith('.json'):
|
|
file_path = os.path.join(root, file)
|
|
if not self.cache_manager.is_processed(file_path):
|
|
unprocessed_articles.append((file_path, file))
|
|
|
|
logger.info(f"Found {len(unprocessed_articles)} unprocessed articles")
|
|
return unprocessed_articles
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error finding unprocessed articles: {e}")
|
|
return []
|
|
|
|
def process_article_file(self, file_path: str, filename: str) -> dict:
|
|
"""
|
|
Process a single article file and extract facts.
|
|
|
|
Args:
|
|
file_path (str): Path to the article file
|
|
filename (str): Name of the article file
|
|
|
|
Returns:
|
|
Dictionary containing the extracted facts or None if failed
|
|
"""
|
|
try:
|
|
with open(file_path, 'r', encoding='utf-8') as f:
|
|
article_data = json.load(f)
|
|
|
|
# Extract facts from the article
|
|
facts = self.fact_extractor.extract_facts_from_article(
|
|
article_data.get('original_content', ''),
|
|
article_data.get('title', filename)
|
|
)
|
|
|
|
# Add metadata
|
|
facts['source'] = article_data.get('source', 'Unknown')
|
|
facts['published'] = article_data.get('published', 'Unknown')
|
|
facts['filename'] = filename
|
|
facts['processed_at'] = datetime.now().isoformat()
|
|
|
|
# Mark as processed in cache
|
|
self.cache_manager.mark_processed(file_path)
|
|
|
|
logger.info(f"Successfully processed article: {filename}")
|
|
return facts
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error processing article {filename}: {e}")
|
|
metrics_collector.increment_articles_failed()
|
|
return None
|
|
|
|
def process_batch(self, articles_batch: List[Tuple[str, str]]) -> Tuple[int, int]:
|
|
"""
|
|
Process a batch of articles.
|
|
|
|
Args:
|
|
articles_batch (List): List of (file_path, filename) tuples
|
|
|
|
Returns:
|
|
Tuple of (successful_count, failed_count)
|
|
"""
|
|
successful = 0
|
|
failed = 0
|
|
|
|
logger.info(f"Processing batch of {len(articles_batch)} articles")
|
|
|
|
for file_path, filename in articles_batch:
|
|
try:
|
|
facts = self.process_article_file(file_path, filename)
|
|
if facts:
|
|
successful += 1
|
|
metrics_collector.increment_articles_processed()
|
|
else:
|
|
failed += 1
|
|
metrics_collector.increment_articles_failed()
|
|
except Exception as e:
|
|
logger.error(f"Error processing batch item {filename}: {e}")
|
|
failed += 1
|
|
metrics_collector.increment_articles_failed()
|
|
|
|
logger.info(f"Batch completed: {successful} successful, {failed} failed")
|
|
return successful, failed
|
|
|
|
def process_all_articles(self, scraper_dir: str = "/scraper/articles") -> dict:
|
|
"""
|
|
Process all unprocessed articles in the scraper directory.
|
|
|
|
Args:
|
|
scraper_dir (str): Path to the scraper articles directory
|
|
|
|
Returns:
|
|
Dictionary with processing statistics
|
|
"""
|
|
metrics_collector.start_processing()
|
|
|
|
start_time = time.time()
|
|
|
|
# Find all unprocessed articles
|
|
unprocessed_articles = self.find_unprocessed_articles(scraper_dir)
|
|
|
|
if not unprocessed_articles:
|
|
logger.info("No unprocessed articles found")
|
|
metrics_collector.stop_processing()
|
|
return {
|
|
"total_processed": 0,
|
|
"total_failed": 0,
|
|
"duration": 0,
|
|
"status": "no_new_articles"
|
|
}
|
|
|
|
logger.info(f"Starting to process {len(unprocessed_articles)} articles in batches of {self.batch_size}")
|
|
|
|
total_processed = 0
|
|
total_failed = 0
|
|
|
|
# Process articles in batches
|
|
for i in range(0, len(unprocessed_articles), self.batch_size):
|
|
batch = unprocessed_articles[i:i + self.batch_size]
|
|
successful, failed = self.process_batch(batch)
|
|
total_processed += successful
|
|
total_failed += failed
|
|
|
|
# Add a small delay between batches to prevent overwhelming the system
|
|
if i + self.batch_size < len(unprocessed_articles):
|
|
time.sleep(0.1)
|
|
|
|
end_time = time.time()
|
|
duration = end_time - start_time
|
|
|
|
metrics_collector.stop_processing()
|
|
metrics_collector.record_processing_time(duration)
|
|
|
|
stats = {
|
|
"total_processed": total_processed,
|
|
"total_failed": total_failed,
|
|
"duration": duration,
|
|
"batch_size": self.batch_size,
|
|
"cache_stats": self.cache_manager.get_cache_stats(),
|
|
"status": "completed"
|
|
}
|
|
|
|
logger.info(f"Processing completed in {duration:.2f} seconds")
|
|
logger.info(f"Total processed: {total_processed}, Total failed: {total_failed}")
|
|
|
|
return stats
|
|
|
|
def process_new_articles(self, scraper_dir: str = "/scraper/articles") -> dict:
|
|
"""
|
|
Process only new articles (those that haven't been processed yet).
|
|
This is designed for real-time processing of new articles.
|
|
|
|
Args:
|
|
scraper_dir (str): Path to the scraper articles directory
|
|
|
|
Returns:
|
|
Dictionary with processing statistics
|
|
"""
|
|
return self.process_all_articles(scraper_dir) |