382 lines
15 KiB
Python
Executable File
382 lines
15 KiB
Python
Executable File
#!/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())
|