140 lines
5.5 KiB
Python
140 lines
5.5 KiB
Python
import json
|
|
import logging
|
|
from datetime import datetime, timezone
|
|
from typing import Optional
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
from .core.config import settings
|
|
from .parsing.mapper import CarMapper
|
|
from .parsing.parser import VehicleParser
|
|
from .storage.db import get_db_session
|
|
from .storage.models import VehicleRawSnapshot, VehicleParseResult, Car
|
|
from .storage.schemas import CarRecord
|
|
|
|
logger = logging.getLogger("iaai_scraper.enrichment_service")
|
|
|
|
|
|
class EnrichmentService:
|
|
"""Сервис для парсинга сырых данных и обогащения автомобилей."""
|
|
|
|
def __init__(self):
|
|
self.db: Session = get_db_session()
|
|
self.parser = VehicleParser()
|
|
self.mapper = CarMapper()
|
|
|
|
def enrich_vehicle(self, snapshot: VehicleRawSnapshot) -> Optional[VehicleParseResult]:
|
|
"""
|
|
Парсит snapshot и сохраняет результат.
|
|
|
|
Args:
|
|
snapshot: Raw snapshot для парсинга
|
|
|
|
Returns:
|
|
VehicleParseResult если парсинг успешен
|
|
"""
|
|
logger.info(f"Enriching snapshot {snapshot.id} for candidate {snapshot.candidate_id}")
|
|
|
|
if not snapshot.success or not snapshot.raw_data:
|
|
logger.warning(f"Snapshot {snapshot.id} has no data to parse")
|
|
return self._save_parse_result(snapshot.id, False, None, "No raw data available")
|
|
|
|
try:
|
|
# Парсим данные
|
|
parsed_data = self._parse_raw_data(snapshot.raw_data)
|
|
if not parsed_data:
|
|
return self._save_parse_result(snapshot.id, False, None, "Parsing failed")
|
|
|
|
# Маппим в CarRecord
|
|
car_record = self._map_to_car_record(parsed_data)
|
|
if not car_record:
|
|
return self._save_parse_result(snapshot.id, False, None, "Mapping failed")
|
|
|
|
# Сохраняем в cars таблицу
|
|
self._save_car(car_record)
|
|
|
|
# Сохраняем успешный результат парсинга
|
|
return self._save_parse_result(snapshot.id, True, json.dumps(parsed_data))
|
|
|
|
except Exception as e:
|
|
error_msg = f"Unexpected error during enrichment: {str(e)}"
|
|
logger.error(f"Error enriching snapshot {snapshot.id}: {error_msg}")
|
|
return self._save_parse_result(snapshot.id, False, None, error_msg)
|
|
|
|
def _parse_raw_data(self, raw_data: str) -> Optional[dict]:
|
|
"""Парсит сырые данные в словарь."""
|
|
try:
|
|
# Если это JSON от browser capture
|
|
if raw_data.strip().startswith('{'):
|
|
data = json.loads(raw_data)
|
|
# Извлекаем vehicle_summary из scrape result
|
|
return data.get("vehicle_summary", {})
|
|
|
|
# Если HTML, используем VehicleParser
|
|
parsed = self.parser.parse_vehicle_page(raw_data, "dummy_url")
|
|
return parsed.get("vehicle_summary", {})
|
|
|
|
except json.JSONDecodeError:
|
|
# Если не JSON, пробуем как HTML
|
|
try:
|
|
parsed = self.parser.parse_vehicle_page(raw_data, "dummy_url")
|
|
return parsed.get("vehicle_summary", {})
|
|
except Exception:
|
|
return None
|
|
|
|
def _map_to_car_record(self, parsed_data: dict) -> Optional[CarRecord]:
|
|
"""Маппит parsed data в CarRecord."""
|
|
try:
|
|
# Используем CarMapper
|
|
payload_insights = {} # TODO: extract from raw data if needed
|
|
return self.mapper.map_to_car_record("dummy_url", parsed_data, payload_insights)
|
|
except Exception as e:
|
|
logger.error(f"Mapping failed: {e}")
|
|
return None
|
|
|
|
def _save_car(self, car_record: CarRecord):
|
|
"""Сохраняет CarRecord в базу данных."""
|
|
# Используем PersistenceService
|
|
from .storage.db import PersistenceService
|
|
from .core.config import settings
|
|
persistence = PersistenceService(settings)
|
|
persistence.create_tables()
|
|
upsert_result = persistence.upsert_car(car_record)
|
|
logger.info(f"Car upserted: {upsert_result}")
|
|
|
|
def _save_parse_result(self, snapshot_id: int, success: bool,
|
|
parsed_data: Optional[str], error_message: Optional[str] = None) -> VehicleParseResult:
|
|
"""Сохраняет результат парсинга."""
|
|
result = VehicleParseResult(
|
|
snapshot_id=snapshot_id,
|
|
success=success,
|
|
parsed_data=parsed_data,
|
|
error_message=error_message
|
|
)
|
|
self.db.add(result)
|
|
self.db.commit()
|
|
self.db.refresh(result)
|
|
return result
|
|
|
|
def process_unparsed_snapshots(self, limit: int = 10) -> int:
|
|
"""
|
|
Обрабатывает snapshots без результатов парсинга.
|
|
|
|
Returns:
|
|
Количество успешно обогащенных snapshots
|
|
"""
|
|
# Находим snapshots без parse results
|
|
snapshots = self.db.query(VehicleRawSnapshot).outerjoin(
|
|
VehicleParseResult, VehicleRawSnapshot.id == VehicleParseResult.snapshot_id
|
|
).filter(
|
|
VehicleRawSnapshot.success == True,
|
|
VehicleParseResult.id.is_(None)
|
|
).limit(limit).all()
|
|
|
|
enriched_count = 0
|
|
for snapshot in snapshots:
|
|
result = self.enrich_vehicle(snapshot)
|
|
if result and result.success:
|
|
enriched_count += 1
|
|
|
|
return enriched_count |