Support kind=llm profiles (instruction/schema), optional multi-event posts via #eN URLs, and recover stale running/queued parse jobs after worker crashes. Co-authored-by: Cursor <cursoragent@cursor.com>
8.9 KiB
Архитектура парсинга (ЦП)
Центр парсинга (ЦП) — сервисы cp-workers* с реестром адаптеров источников. Собирают события и отправляют их в ЦА через internal API.
Роль в платформе
flowchart LR
CA[ЦА ca-api]
Redis[(Redis cp:jobs:family)]
CP[ЦП adapters]
Src[Sources]
CA -->|"RPUSH JobPayload"| Redis
Redis -->|"BLPOP"| CP
CP --> Src
CP -->|"POST /internal/ingest"| CA
CP -->|"PATCH /internal/jobs/{id}"| CA
CP -->|"GET /internal/listener/subscriptions"| CA
| Направление | Механизм | Назначение |
|---|---|---|
| ЦА → ЦП | Redis cp:jobs:{family} |
Batch-задания по семейству адаптеров |
| ЦП → источники | Адаптер (telegram / crawl4ai / viina) |
Fetch + extract |
| ЦП → ЦА | POST /internal/ingest |
Запись событий (IngestEventItem) |
| ЦП → ЦА | PATCH /internal/jobs/{id} |
Статус batch-задания |
| ЦП ← ЦА | GET /internal/listener/subscriptions |
Каналы real-time (только Telegram) |
Общие контракты: contracts/jobs.py, contracts/ingest.py, contracts/sources.py, contracts/queues.py.
Адаптеры источников
Каждый source_type реализует SourceAdapter:
async def run(job_id, source_config, *, ctx) -> tuple[list[dict], str | None]
Выход — список dict в форме IngestEventItem. ЦА не знает про Crawl4AI/VIINA.
| source_type | Семейство / очередь | Сервис | Зависимости |
|---|---|---|---|
telegram |
telegram → cp:jobs:telegram |
cp-workers |
Telethon |
crawl4ai |
web → cp:jobs:web |
cp-workers-web |
Crawl4AI + Playwright |
viina |
nlp → cp:jobs:nlp |
cp-workers-nlp |
httpx + BeautifulSoup |
Маршрутизация при enqueue в ЦА: queue_key_for_source.
Воркер слушает WORKER_FAMILIES / очереди своих ENABLED_ADAPTERS. Чужой job → requeue в нужную очередь.
Правила для всех адаптеров:
- стабильный
source_url(дедуп в ЦА); source_typeсобытия = тип адаптера;source_configвалидируется схемами изcontracts/sources.py.
Режимы извлечения (Telegram)
Цель всегда одна: фиксированные поля IngestEventItem / Event в ЦА (title, description, locality, coords→lat/lng, event_date, topic, …). Кастомные пользовательские таблицы — roadmap.
extract_mode |
Что делает | Откуда в UI |
|---|---|---|
heuristic |
Legacy regex-парсер постов | raw job API (не пара) |
llm |
DeepSeek на каждый пост (batch) | профиль kind=llm → pair flatten |
profile |
Статичные правила HeuristicProfile |
профиль kind=heuristic → pair flatten |
LLM-режим (extract_mode: llm)
Для неструктурированных Telegram-постов и Crawl4AI:
- ключ
DEEPSEEK_API_KEYв.env(воркерыcp-workers/cp-workers-web; preview вca-api); - Telegram batch: текст → DeepSeek JSON → один или несколько
IngestEventItem(workers/llm_extract.py); - опционально
extract_schema,instruction,required_fields(post-extract gate поверхis_event); multi_event(opt-in вllm_profile/TelegramSourceConfig): модель возвращает массивevents; каждое событие — отдельный ingest. URL: один event →post.url; несколько →post.url#e1,#e2, … (дедуп ЦА поsource_url);- heuristic / profile по-прежнему 1 пост → 1 событие;
- Crawl4AI: страница →
LLMExtractionStrategy(DeepSeek) с fallback на тот же DeepSeek по markdown; - listener: для
llmпока fallback на heuristic (LLM — batch-only by design).
Reusable LLM-профиль в ЦА: ParserProfile.kind=llm + JSON llm_profile (contracts/llm_profile.py). При enqueue flatten → extract_mode=llm + schema/instruction/required_fields/multi_event.
Profile-режим (extract_mode: profile)
Статичный парсер из конструктора ЦА:
heuristic_profileвsource_config(схемаcontracts/heuristic_profile.py);- интерпретатор:
workers/heuristic_profile.py(те же правила, что preview в ЦА); required_fields: пост без заполненных обязательных полей не ингестится (batch + listener);- генерация правил — один раз в админке (
/admin/parser-profiles/generate); runtime без LLM.
Reusable heuristic-профиль: ParserProfile.kind=heuristic + heuristic_profile. Pair на «Парсеры» → flatten extract_mode=profile.
UI «Профили»: выбор kind (heuristic | llm). Связка канал+профиль на «Парсеры».
flowchart LR
Job[JobPayload]
Reg[AdapterRegistry]
TG[TelegramAdapter]
C4[Crawl4AIAdapter]
VI[ViinaAdapter]
Ingest[IngestEventItem]
Job --> Reg
Reg --> TG --> Ingest
Reg --> C4 --> Ingest
Reg --> VI --> Ingest
Структура каталога
centers/parsing/
├── ARCHITECTURE.md
└── workers/
├── Dockerfile # context = repo root; ARG REQUIREMENTS_FILE
├── requirements.txt # telegram
├── requirements-web.txt # crawl4ai
├── requirements-nlp.txt # viina
├── worker.py
└── workers/
├── adapters/
│ ├── base.py # SourceAdapter, WorkerContext
│ ├── registry.py # ENABLED_ADAPTERS
│ ├── telegram.py
│ ├── crawl4ai_adapter.py
│ └── viina.py
├── converter.py
├── parsers/telegram_events.py
└── sources/ # Telethon session / listener / client
Режимы работы
| Режим | Условие | Поведение |
|---|---|---|
| Listener + batch | ENABLED_ADAPTERS включает telegram и TELEGRAM_LISTENER_ENABLED=true |
Shared Telethon + listener + worker_loop |
| Только batch | listener выключен или нет telegram | Только BLPOP по очередям семейства |
Поток batch-заданий
- ЦА
enqueue_job→ Rediscp:jobs:{family}с{ job_id, source_type, source_config } - Воркер семейства:
BLPOP→handle_job→registry.get(source_type).run(...) POST /internal/ingest+ статус job
Legacy-ключ cp:jobs по-прежнему дренируется telegram-воркером (совместимость).
Telegram real-time
TelegramListener: подписки из ЦА, ingest с listener: true (статус ParseJob не трогается).
extract_mode=profile— те же heuristic-правила, что batch;extract_mode=llm— не вызывает DeepSeek; fallback на legacy heuristic (LLM только в batch).
Переменные окружения
| Переменная | Назначение |
|---|---|
ENABLED_ADAPTERS |
Список адаптеров через запятую (telegram, crawl4ai, viina) |
WORKER_FAMILIES |
Какие семейства очередей слушать (telegram, web, nlp) |
REDIS_URL / CA_API_URL / INTERNAL_TOKEN |
Как раньше |
TELEGRAM_* |
Только для cp-workers |
Как добавить новый источник
- Схема
source_configвcontracts/sources.py+ запись вSOURCE_FAMILY(queues.py) - Класс адаптера в
workers/adapters/+ factory вregistry.py - При тяжёлых deps —
requirements-*.txtи сервис вdocker-compose.yml - Поля формы в UI «Парсеры»
Связанные части ЦА
| Файл ЦА | Роль |
|---|---|
services/jobs.py |
enqueue_job → cp:jobs:{family} |
routers/admin.py |
CRUD + валидация source_config |
services/ingest.py |
Сохранение Event + карта |
UI /parsers |
Выбор telegram / crawl4ai / viina |