Files
vector/scripts/parse_reports.py
2026-06-06 18:31:55 +00:00

382 lines
15 KiB
Python
Executable File
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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())