#!/usr/bin/env python3 """ Скрипт для парсинга спецдонесений МЧС Беларуси из .docx файлов. Использует Claude API для извлечения структурированных данных. """ import argparse import asyncio import json import logging import os import sys from pathlib import Path from typing import Dict, List, Optional import httpx from docx import Document from sqlalchemy import create_engine, text from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker # Настройка логирования logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s' ) logger = logging.getLogger(__name__) # Список полей для извлечения FIELDS = [ "age_years", "gender", "has_diagnosis", "diagnosis_type", "has_transport", "season", "time_of_day", "elapsed_before_report_h", "last_seen_direction", "terrain_primary", "water_nearby", "road_nearby", "found_alive", "found_distance_km", "found_direction", "found_location_type", "search_duration_hours", "who_found" ] EXTRACTION_PROMPT = """Ты — система извлечения данных из спецдонесений МЧС Беларуси. Извлеки поля и верни ТОЛЬКО JSON. Поля для извлечения: - age_years: возраст в годах (число) - gender: пол (мужской/женский) - has_diagnosis: есть ли диагноз/заболевание (true/false) - diagnosis_type: тип диагноза если есть (строка) - has_transport: был ли транспорт (true/false) - season: сезон (winter/spring/summer/autumn) - time_of_day: время суток (morning/day/evening/night) - elapsed_before_report_h: часов прошло до сообщения (число) - last_seen_direction: направление последнего наблюдения (строка) - terrain_primary: основной тип местности (forest/field/urban/mountain/water) - water_nearby: водоем рядом (true/false) - road_nearby: дорога рядом (true/false) - found_alive: найден живым (true/false) - found_distance_km: расстояние обнаружения в км (число) - found_direction: направление обнаружения (север/юг/восток/запад и т.д.) - found_location_type: тип места обнаружения (строка) - search_duration_hours: длительность поиска в часах (число) - who_found: кто нашел (спасатели/волонтеры/родственники/сам вышел и т.д.) Для каждого поля верни объект: { "value": <значение или null>, "confidence": "high" | "medium" | "low" } Если поле не упомянуто в тексте — value: null, confidence: "low". Верни ТОЛЬКО валидный JSON в формате: { "age_years": {"value": 65, "confidence": "high"}, "gender": {"value": "мужской", "confidence": "high"}, ... } ТЕКСТ ДОНЕСЕНИЯ: {text} Ответь ТОЛЬКО JSON без дополнительного текста.""" async def extract_text_from_docx(file_path: Path) -> str: """Извлекает текст из .docx файла.""" try: doc = Document(file_path) text = "\n".join([paragraph.text for paragraph in doc.paragraphs]) return text.strip() except Exception as e: logger.error(f"Ошибка чтения {file_path}: {e}") return "" async def extract_data_with_claude(text: str, api_key: str) -> Optional[Dict]: """Извлекает структурированные данные из текста с помощью Claude API.""" prompt = EXTRACTION_PROMPT.format(text=text[:8000]) # Ограничиваем длину try: async with httpx.AsyncClient(timeout=120.0) as client: response = await client.post( "https://api.anthropic.com/v1/messages", headers={ "x-api-key": api_key, "anthropic-version": "2023-06-01", "content-type": "application/json" }, json={ "model": "claude-sonnet-4-20250514", "max_tokens": 4096, "messages": [ { "role": "user", "content": prompt } ] } ) if response.status_code != 200: logger.error(f"Claude API error: {response.status_code} - {response.text}") return None result = response.json() content = result["content"][0]["text"] # Убираем markdown блоки если есть if "```json" in content: content = content.split("```json")[1].split("```")[0].strip() elif "```" in content: content = content.split("```")[1].split("```")[0].strip() data = json.loads(content) return data except Exception as e: logger.error(f"Ошибка извлечения данных: {e}") return None def count_high_confidence_fields(data: Dict) -> int: """Подсчитывает количество полей с high confidence.""" count = 0 for field in FIELDS: if field in data and data[field].get("confidence") == "high": count += 1 return count async def save_to_database( file_path: Path, raw_text: str, extracted_data: Dict, db_url: str ) -> bool: """Сохраняет данные в PostgreSQL.""" try: # Создаем async engine engine = create_async_engine(db_url, echo=False) async_session = sessionmaker( engine, class_=AsyncSession, expire_on_commit=False ) async with async_session() as session: # Сохраняем raw документ raw_doc_query = text(""" INSERT INTO raw_documents (filename, raw_text, extracted_json, created_at) VALUES (:filename, :raw_text, :extracted_json, NOW()) RETURNING id """) result = await session.execute( raw_doc_query, { "filename": file_path.name, "raw_text": raw_text, "extracted_json": json.dumps(extracted_data, ensure_ascii=False) } ) doc_id = result.scalar_one() # Подготавливаем данные для cases case_data = { "raw_document_id": doc_id, "age": extracted_data.get("age_years", {}).get("value"), "gender": extracted_data.get("gender", {}).get("value"), "has_diagnosis": extracted_data.get("has_diagnosis", {}).get("value"), "diagnosis_type": extracted_data.get("diagnosis_type", {}).get("value"), "has_transport": extracted_data.get("has_transport", {}).get("value"), "season": extracted_data.get("season", {}).get("value"), "time_of_day": extracted_data.get("time_of_day", {}).get("value"), "elapsed_before_report_h": extracted_data.get("elapsed_before_report_h", {}).get("value"), "last_seen_direction": extracted_data.get("last_seen_direction", {}).get("value"), "terrain_primary": extracted_data.get("terrain_primary", {}).get("value"), "water_nearby": extracted_data.get("water_nearby", {}).get("value"), "road_nearby": extracted_data.get("road_nearby", {}).get("value"), "found_alive": extracted_data.get("found_alive", {}).get("value"), "found_distance_km": extracted_data.get("found_distance_km", {}).get("value"), "found_direction": extracted_data.get("found_direction", {}).get("value"), "found_location_type": extracted_data.get("found_location_type", {}).get("value"), "search_duration_hours": extracted_data.get("search_duration_hours", {}).get("value"), "who_found": extracted_data.get("who_found", {}).get("value") } # Сохраняем case case_query = text(""" INSERT INTO cases ( raw_document_id, age, gender, has_diagnosis, diagnosis_type, has_transport, season, time_of_day, elapsed_before_report_h, last_seen_direction, terrain_primary, water_nearby, road_nearby, found_alive, found_distance_km, found_direction, found_location_type, search_duration_hours, who_found, created_at ) VALUES ( :raw_document_id, :age, :gender, :has_diagnosis, :diagnosis_type, :has_transport, :season, :time_of_day, :elapsed_before_report_h, :last_seen_direction, :terrain_primary, :water_nearby, :road_nearby, :found_alive, :found_distance_km, :found_direction, :found_location_type, :search_duration_hours, :who_found, NOW() ) """) await session.execute(case_query, case_data) await session.commit() logger.info(f"✓ Сохранено в БД: {file_path.name} (doc_id={doc_id})") return True except Exception as e: logger.error(f"Ошибка сохранения в БД: {e}") return False async def process_file( file_path: Path, api_key: str, db_url: Optional[str], dry_run: bool ) -> Dict: """Обрабатывает один файл.""" logger.info(f"Обработка: {file_path.name}") # Извлекаем текст text = await extract_text_from_docx(file_path) if not text: return { "file": file_path.name, "status": "error", "reason": "Не удалось извлечь текст" } # Извлекаем данные через Claude extracted_data = await extract_data_with_claude(text, api_key) if not extracted_data: return { "file": file_path.name, "status": "error", "reason": "Ошибка извлечения данных" } # Подсчитываем high confidence поля high_conf_count = count_high_confidence_fields(extracted_data) logger.info(f" Полей с high confidence: {high_conf_count}/{len(FIELDS)}") # Проверяем минимальный порог if high_conf_count < 6: logger.warning(f" ⚠ Пропущено: недостаточно полей с high confidence ({high_conf_count} < 6)") return { "file": file_path.name, "status": "skipped", "reason": f"Недостаточно полей с high confidence ({high_conf_count} < 6)", "high_confidence_count": high_conf_count } # Сохраняем в БД если не dry-run if not dry_run and db_url: success = await save_to_database(file_path, text, extracted_data, db_url) status = "success" if success else "db_error" else: logger.info(f" [DRY RUN] Данные не сохранены в БД") status = "dry_run" return { "file": file_path.name, "status": status, "high_confidence_count": high_conf_count, "extracted_data": extracted_data if dry_run else None } async def main(): parser = argparse.ArgumentParser( description="Парсинг спецдонесений МЧС из .docx файлов" ) parser.add_argument( "--dir", required=True, help="Папка с .docx файлами" ) parser.add_argument( "--db", help="DATABASE_URL для PostgreSQL" ) parser.add_argument( "--dry-run", action="store_true", help="Не записывать в БД, только показать результаты" ) args = parser.parse_args() # Проверяем API ключ api_key = os.getenv("ANTHROPIC_API_KEY") if not api_key: logger.error("ANTHROPIC_API_KEY не установлен в переменных окружения") sys.exit(1) # Проверяем директорию dir_path = Path(args.dir) if not dir_path.exists() or not dir_path.is_dir(): logger.error(f"Директория не найдена: {args.dir}") sys.exit(1) # Получаем список .docx файлов docx_files = list(dir_path.glob("*.docx")) if not docx_files: logger.warning(f"Не найдено .docx файлов в {args.dir}") sys.exit(0) logger.info(f"Найдено файлов: {len(docx_files)}") logger.info(f"Режим: {'DRY RUN' if args.dry_run else 'ЗАПИСЬ В БД'}") # Обрабатываем файлы results = [] for file_path in docx_files: result = await process_file(file_path, api_key, args.db, args.dry_run) results.append(result) # Небольшая пауза между запросами к API await asyncio.sleep(1) # Выводим итоговую статистику logger.info("\n" + "="*60) logger.info("ИТОГОВАЯ СТАТИСТИКА") logger.info("="*60) success_count = sum(1 for r in results if r["status"] == "success") skipped_count = sum(1 for r in results if r["status"] == "skipped") error_count = sum(1 for r in results if r["status"] == "error") dry_run_count = sum(1 for r in results if r["status"] == "dry_run") logger.info(f"Всего файлов: {len(results)}") logger.info(f"Успешно обработано: {success_count}") logger.info(f"Пропущено (< 6 high conf): {skipped_count}") logger.info(f"Ошибки: {error_count}") if dry_run_count > 0: logger.info(f"Dry run: {dry_run_count}") # Сохраняем детальный отчет report_path = Path("parse_report.json") with open(report_path, "w", encoding="utf-8") as f: json.dump(results, f, ensure_ascii=False, indent=2) logger.info(f"\nДетальный отчет сохранен в: {report_path}") if __name__ == "__main__": asyncio.run(main())