Спецификация реализации · живой поиск · tazzz.ru

Спецификация: ёмкость без закупок

Как именно реализовать программу Р1–Р5 из аудита: точные файлы и функции, схемы данных, флаги, порядок выкатки, критерии отката и тесты. Формат — по канону инженерных design doc'ов: контекст → цели и не-цели → инварианты → дизайн по потокам работ → альтернативы → выкатка → наблюдаемость. Каждый поток независим, за своим флагом, откатывается одной переменной.

Версия 1.0 · 27 августа 2026 База ветка full-update, sha 940d5879 Основание аудит emkost-bez-zakupok (все числа — оттуда) Статус ждёт утверждения владельца
01 · Контекст

Контекст, цели, не-цели

Контекст. Аудит 27.08 (разбор) показал: жёсткий потолок ~240 поисков/час складывается из раздутого времени детального слота (50–65% — внутренняя работа: запись в БД 1,8–4,8 с, инлайновый LLM p90 31 с), повторной покупки уже известных машин (39–54% машин живых поисков уже в БД) и неуправляемого пула прокси (5 уникальных IP, один мёртв и невидим, che168 — SPOF серверного IP).

Цели

Не-цели

Порядок потоков в документе = порядок внедрения: от нулевого риска (W0) к единственному потоку с реальным техническим риском (W6). Потоки независимы: любой можно пропустить или откатить, не трогая остальные.

02 · Инварианты

Инварианты выдачи

Требования, зафиксированные владельцем: юзер не должен заметить ничего, кроме скорости. Проверяются на каждой фазе выкатки.

  1. И1 · Состав карточки не меняется. Полный набор полей (цена, «под ключ», ТТХ, комплектация, осмотр, фото) — тот же. Ни одна карточка не эмитится в ленту без «под ключ», если сегодня она эмитится с ним (см. W2: слот освобождается, но эмит ждёт цену).
  2. И2 · Цена и пробег всегда живые. В любом режиме источником цены/пробега «Свежих» является листинг текущей сессии, никогда — сохранённое значение (W5 перекладывает их из листинга поверх карточки из БД).
  3. И3 · Статусы и бейджи как сегодня. Карточка, отданная дельта-путём W5, несёт card_status='full', detail_complete=true и НЕ несёт is_cached — фронт не отличает её от перепарсенной. Единственное допустимое расхождение: фото и текст описания могли устареть до фонового обновления; кнопка «Обновить» остаётся ручным форсом.
  4. И4 · Порядок появления карточек — как сегодня: по позиции листинга; basic → full; hold-basic оси (services/card_visibility.should_hold_basic_emit) не трогаем.
  5. И5 · «Из БД», агент, кампании, refresh-кнопка ведут себя бит-в-бит как раньше: все врезки гейтятся флагами, которые эти пути не выставляют.
  6. И6 · При выключенных флагах поведение бит-в-бит прежнее. Каждый поток — за своим env-флагом с дефолтом «выключено» (паттерн ADMISSION_CONTROL/DETAIL_DISPATCH_MODE, уже канон кодовой базы).
03 · Потоки

Дизайн по потокам работ

Для каждого потока: мотивация числом из аудита → точные изменения → флаги и откат → тесты → метрики приёмки.

W0

Однострочники

~2 часа · риск ноль

Мотивация. Четыре рукотворных ограничителя снимаются без нового кода: устаревший вес dongchedi (Playwright не используется — 6 срабатываний за 7 дней), мёртвый прокси .165 (25% фетчей dongchedi/encar в таймаут 12–20 с), che168 без единого прокси (SPOF, при том что проба через .247 даёт 200 за 1,8 с), ручной потолок 5.

Изменения

#ГдеЧтоЭффект
0.1services/capacity_service.py:59'dongchedi': 250 → 120 (+ поправить комментарий: JSON-API, не Playwright)допуск dongchedi-миксов ×2: книга весов 2800 МБ ÷ 120 ≈ 23 вместо 11
0.2прод, таблица proxiesUPDATE proxies SET enabled=false WHERE host='79.172.218.165' (обе строки: dongchedi + encar). Вернуть, когда провайдер починит (проверка — проба из W4)−25% фетчей-в-таймаут немедленно
0.3прод, таблица proxies4 INSERT-строки site_source='che168' с существующими хостами .214/.230/.247 (+.165 после починки). Скрапер che168 уже умеет прокси-лист (путь DETAIL_PROXY_MATCHED_SITES содержит che168; _raw_proxy_urls('che168') начнёт отдавать пул). CHE168_API_DIRECT оставить как фолбэк-маршрутche168: 1 маршрут → 4–5, SPOF снят; кап площадки 4 → 4×IP (окончательно капы пересчитает W4)
0.4.env продаMAX_CONCURRENT_SEARCHES=5→8, через 2–3 дня наблюдения → 12 (план Т2.2 отчёта ёмкости)залпы перестают стоять в очереди admission
0.5Postgres (см. W1)ALTER SYSTEM SET wal_compression='on'; SELECT pg_reload_conf(); — без рестартадешевле WAL на HDD, полезно всем писателям
0.6services/capacity_service.py:137'che168': 4 → ANTIBAN_DETAIL_LIMIT_JSON в .env: {"che168": 12} (3 фетча × 4 живых маршрута — то же правило, что у dongchedi) — только ПОСЛЕ 0.3потолок che168 ~80 → ~240 сессий/час

Тесты и приёмка

  • Смоук после 0.2/0.3: один живой поиск dongchedi + один che168; в Loki fetch_ms p50 che168 не хуже базлайна (3,1 с), нет всплеска SourceBlockedError.
  • Приёмка 0.1/0.4: capacity_service.max_concurrent_count() в логах admission; глубина очереди v2 в пиках = 0.
  • Откат: каждая правка — независимая (env/SQL), откатывается тем же способом за минуту.
Важно про 0.3: первые сутки держать che168 на консервативном темпе (кап 8, не 12) и смотреть частоту che168_gate.classify != OK в Loki: замер 2026-06-26 показал, чтоche168 лимитит по IP под бурстом — расширение маршрутов легально, бурст с одного маршрута — нет.
W1

Быстрая запись в БД: асинхронный коммит upsert'а

~1 день · риск низкий

Мотивация. save_db_ms: encar p50 4,8 с, dongchedi 2,2 с, che168 1,8 с — на HDD c synchronous_commit=on каждый коммит ждёт fsync WAL. Данные parsed_cars переспарсиваемы: потеря последних ~0,6 с транзакций при падении сервера не создаёт ни повреждений, ни невосполнимых потерь (канон PG: async commit ≠ fsync=off, целостность БД не страдает).

Дизайн

Асинхронный коммит включается per-транзакция через SET LOCAL — глобальный дефолт БД не трогаем (заявки, юзеры, платёжные сущности остаются durable). Точка врезки одна: DatabaseService._add_parsed_car_once (services/database_service.py:270+) — начало транзакции сохранения карточки. Тот же приём — в services/card_projection.py (писатель проекций очереди III) и в батч-персисте сессии (persist_cars_to_db).

# services/db_tx.py (новый модуль, ~30 строк)
ASYNC_COMMIT_FLAG = 'PARSED_CAR_ASYNC_COMMIT'   # дефолт '0' — поведение бит-в-бит

def apply_async_commit(session) -> None:
    """SET LOCAL synchronous_commit=off на ТЕКУЩУЮ транзакцию.
    Вызывать первым оператором TX. No-op при выключенном флаге или не-Postgres."""
    if os.environ.get(ASYNC_COMMIT_FLAG, '0') != '1':
        return
    if session.bind.dialect.name != 'postgresql':
        return
    session.execute(text("SET LOCAL synchronous_commit TO OFF"))

Вызов — первой строкой в _add_parsed_car_once (до upsert; SQLAlchemy autobegin откроет TX этим оператором, и SET LOCAL умрёт вместе с commit/rollback — протечь на другие транзакции сессии не может). Ошибка apply_async_commit глушится (fail-open: пусть медленно, но сохранит).

Дополнительно в этом потоке (по одному коммиту на правку)

  • raw_json из горячего пути. В add_parsed_car поле raw_json пишется только при изменении md5 относительно уже сохранённого (сегодня — перезапись всегда → лишний TOAST-трафик на каждый refresh той же машины). Отдельный дешёвый guard, не миграция.
  • Не перезаписывать неизменное. В upsert-словарь не включать колонки, значение которых равно текущему (SQLAlchemy: сравнение по загруженной строке уже есть — get_parsed_car_by_site_id вызывается в HP-reuse; расширить до диффа). Это срезает и WAL, и обновления части из 20 индексов.

Флаг, тесты, приёмка, откат

  • Флаг PARSED_CAR_ASYNC_COMMIT=1 в .env (durable-дефолт в compose НЕ ставим до недели наблюдения).
  • Тест backend/test_async_commit.py: (а) SET LOCAL реально применяется (SHOW внутри TX), (б) не протекает в следующую TX той же сессии, (в) no-op на sqlite.
  • Приёмка: save_db_ms в Loki (поле уже пишется каждой задачей): encar p50 < 1 с, che168/dongchedi < 0,5 с в течение 48 ч.
  • Откат: флаг в 0, мгновенно. Критерий отката: любой рост ошибок сохранения или деградация p95 других запросов PG.
W2

LLM вне слота: асинхронное дообогащение мощности

~2 дня · риск средний

Мотивация. Этап «расчёт» dongchedi: p50 2,3 с, p90 31 с — при промахе кэша мощности (HorsepowerCache) внутри детального слота выполняется LLM-запрос через глобальный RPM-лимитер. Слот (дефицит: 24 шт.) стоит и не фетчит.

Дизайн

Разделяем «редкий дорогой» путь от «частого дешёвого». Сегодняшняя цепочка в _run_single_car_detail (tasks/distributed_parser_tasks.py:5193): fetch → translate → ensure_hp_before_cost (:3179, здесь LLM при промахе) → cost → save → emit. Новая логика в ensure_hp_before_cost:

# tasks/distributed_parser_tasks.py, ensure_hp_before_cost()
# HP_ASYNC_ENRICH=1: при промахе ВСЕХ дешёвых слоёв (БД-reuse, каталог,
# HorsepowerCache) LLM НЕ вызываем в слоте:
car_data['hp_pending'] = True        # транзитный флаг, в БД не пишется
# → cost пропускается (has_power_info=False, как сегодня при отсутствии HP)
# → save проходит (карточка сохраняется без total_cost_rub)
# → ЭМИТ ПРИДЕРЖИВАЕТСЯ (инвариант И1): вместо emit_car_full — постановка
#   tasks.enrichment_tasks.resolve_hp_and_cost.delay(car_db_id, session_id,
#   site_car_id) в очередь 'default' (control-воркер, слоты не дефицит)
# tasks/enrichment_tasks.py (новый, ~120 строк)
@shared_task(name='tasks.enrichment_tasks.resolve_hp_and_cost',
             queue='default', time_limit=120, acks_late=True)
def resolve_hp_and_cost(car_db_id, session_id, site_car_id):
    # 1. LLM-подбор л.с. (существующий horsepower_estimator, там же write-back
    #    в HorsepowerCache — второй такой же машине LLM уже не нужен)
    # 2. cost_calculator.calculate_total_cost(...) — тот же вызов, что в слоте
    # 3. UPDATE parsed_cars (engine_power_hp, total_cost_rub, calculation_json)
    # 4. session_service.update_car_to_full + ws_emitter.emit_car_full
    #    (механизм существует: services/ws_emitter.py:120)
    # 5. _autodrive_after_car НЕ здесь: mark_detail_done делает СЛОТ-задача
    #    сразу после save — конвейер не ждёт LLM

Ключевые решения

  • Инвариант И1 (карточка только с «под ключ»): эмит уходит в enrichment-задачу — юзер видит карточку чуть позже, но всегда целиком, как сегодня. Для сессии машина считается обработанной сразу после save (finalize не ждёт LLM); enrichment догоняет эмитом даже в завершённую сессию — фронт применяет update_car_to_full независимо от статуса сессии.
  • Касается только промахов кэша: по выборке аудита это меньшинство машин (кэш по (brand, model, year, trim) — популярные модели прогреты). Прогретые идут старым путём бит-в-бит.
  • Дедлайн: если enrichment не смог (LLM недоступен) — карточка эмитится через 60 с без «под ключ», как сегодня ведут себя машины с неразрешимой мощностью (существующая деградация, не новая).
  • Refresh-путь и priority-клик (is_refresh=True / priority_detail_fetch) флаг не читают — там юзер явно ждёт одну машину, LLM в слоте оправдан.

Флаг, тесты, приёмка, откат

  • Флаг HP_ASYNC_ENRICH=1, дефолт 0. Роутинг задачи — в celery_app.py task_routes.
  • Тест backend/test_hp_async_enrich.py: промах кэша → задача поставлена, слот-задача завершилась без LLM; enrichment дописал hp+cost и заэмитил; дедлайн-ветка.
  • Приёмка: cost_ms p90 dongchedi < 3 с (было 31 с); «время до первой карточки» сессии не ухудшилось; доля карточек без total_cost в выдаче не выросла (SQL-проба по свежим сессиям).
  • Откат: флаг в 0. Критерий: рост карточек без «под ключ» дольше 60 с или очередь default > 100.
W3

Диета запросов encar: параллельно и без лишней страницы

~1–2 дня · риск низкий

Мотивация. fetch_ms encar p50 5,6 с, p90 12 с: деталь — это 2–4 последовательных запроса (v1/readside/vehicle → иногда HTML-страница fem.encar.com ради inspection-ID → отчёт осмотра), см. services/scrapers/encar_scraper.py:1173–1400. Каждый запрос — это ещё и расход анти-бан-бюджета IP: вдвое меньше запросов на машину = вдвое больше машин на тот же кап.

Изменения (все — внутри encar_scraper.py)

  • 3.1 · Убрать fem-страницу из горячего пути. Inspection-ID сегодня извлекается из URL первого фото (_get_inspection_car_id_from_photo, регэксп /pic\d+/(\d+)_\d+\.jpg) с фолбэком на HTML-страницу. Сделать photo-путь единственным синхронным; при неудаче регэкспа — НЕ тянуть fem-страницу в слоте, а вернуть карточку без осмотра и поставить лёгкую задачу дозагрузки осмотра (тот же механизм update_car_to_full, что в W2). Метрика: доля фолбэков (ожидаемо <5% — формат URL стабилен).
  • 3.2 · Осмотр параллельно карточке. Два независимых GET (vehicle и inspection) выполняются в ThreadPoolExecutor(max_workers=2) внутри задачи (curl_cffi отпускает GIL — проверено экспериментом аудита). Слоты proxy_limiter это уже учитывают: оба запроса идут через один прокси-слот задачи.
  • 3.3 · Переиспользование соединения. _get_session() держит одну curl_cffi-сессию (keep-alive/HTTP-2 к api.encar.com) на жизнь скрапера — убрать пересоздание сессии между под-запросами, если есть. Мелочь, но на 4 запросах экономит 3 TLS-рукопожатия через прокси (~0,3–0,9 с).
  • 3.4 · Голый except: с ретраем в DIRECT (строки 1190–1196: при любой ошибке прокси запрос молча повторяется с IP сервера) — убрать: ретрай только на 407, в остальных случаях ошибка идёт наверх и учитывается здоровьем прокси (W4). Сегодня этот except маскирует мёртвый .165 и светит серверный IP на encar.

Флаг, тесты, приёмка, откат

  • Флаг ENCAR_PARALLEL_DETAIL=1, дефолт 0 (гейтит 3.1+3.2; 3.3/3.4 — безусловные фиксы).
  • Тест backend/test_encar_detail_diet.py: мок обоих эндпоинтов → карточка идентична последовательной версии поле-в-поле (снапшот-сравнение); ветка фолбэка ID.
  • Приёмка: fetch_ms encar p50 < 3,5 с; число HTTP-запросов на машину (по логам stage=http с trace_id) ≤ 2 в ≥95% задач; отчёт осмотра присутствует у той же доли карточек, что в базлайне.
  • Откат: флаг в 0.
W4

Флот прокси: реестр здоровья, per-маршрут предохранитель, AIMD-темп

~3–4 дня · риск низкий

Мотивация. Мёртвый .165 неделями в ротации; статистика per-process; предохранитель гасит площадку целиком на 30 мин из-за одного адреса; капы 12/4/16 — догадки. Строим на существующих примитивах — services/proxy_limiter.py уже умеет mark_unhealthy / is_unhealthy / has_healthy (Redis, cooldown) и распределённый семафор слотов; таблица proxies уже имеет счётчики (нулевые) и last_used_at.

Архитектура

задача детали pick_route(site) beat: fleet_health_probe раз в 60 с, лёгкий GET на IP services/proxy_fleet.py health: is_unhealthy (есть) выбор: наименее загруженный темп: AIMD-бакет (site×IP) счётчики → таблица proxies breaker: (site×IP), не site 5 маршрутов 4 IP + DIRECT площадки dcd·encar·che168·mde успех/сбой/429 → record_result → темп и здоровье
proxy_fleet — тонкий слой над существующими proxy_limiter и таблицей proxies. Скраперы не переписываются: меняется только источник списка в _get_proxy_urls.

Дизайн: модуль services/proxy_fleet.py

# Redis-схема (все ключи с TTL):
# fleet:health:{proxy}          HASH {ok_streak, fail_streak, last_probe_ms, alive:0|1}
# fleet:rate:{site}:{proxy}     STRING текущий разрешённый темп, req/min (AIMD)
# fleet:tokens:{site}:{proxy}   токен-бакет (Lua: refill по rate, take 1)
# source_breaker:{site}:{egress}:paused   — маршрутный предохранитель (см. ниже)

def pick_route(site: str) -> str | None:
    """Выбрать egress для фетча: живые (not is_unhealthy) → из них с наибольшим
    остатком токенов; ждать токен до FLEET_TOKEN_WAIT_SEC (дефолт 15), иначе
    следующий маршрут. None = все маршруты мертвы/на паузе (вызывающий код
    трактует как сегодняшний 'нет здорового egress' — _has_healthy_egress)."""

def record_result(site, proxy, outcome):  # ok | net_error | block_429_403 | challenge
    """ok: ok_streak+=1; каждый FLEET_AIMD_UP_STREAK (дефолт 50) подряд —
         rate *= 1.1 (потолок FLEET_RATE_MAX, дефолт 60/мин на IP)
       net_error: fail_streak+=1; 3 подряд → mark_unhealthy(cooldown 300 c)
       block/challenge: rate *= 0.5 (пол FLEET_RATE_MIN, дефолт 6/мин)
         + сбой в МАРШРУТНЫЙ предохранитель (site×proxy)
       Счётчики дублируются в proxies.success_count/failure_count батчем
       раз в 60 с (beat), не на каждый запрос."""

Изменения по файлам

ГдеЧто
services/proxy_fleet.pyновый модуль (~250 строк): pick_route, record_result, Lua токен-бакета, стартовые темпы = сегодняшние капы, пересчитанные в req/мин (dongchedi 3 фетча×IP ≈ 24/мин при 7,4 с/фетч → старт 24)
tasks/fleet_tasks.py + celery_app.py (beat)задача fleet_health_probe раз в 60 с (очередь default): на каждый enabled-прокси лёгкий GET пробной цели площадки (encar: search/car/list/general?count=true; dongchedi: all_brand; che168: www.che168.com HEAD; timeout 5 с) через сам прокси; результат → record_result + alive-флаг; батч-flush счётчиков в PG
tasks/distributed_parser_tasks.py:2247_get_proxy_urls: под флагом PROXY_FLEET=on порядок списка отдаёт fleet (живые, наименее загруженные — первыми) вместо random.shuffle; интерфейс (список URL) не меняется — скраперы не трогаем
services/source_breaker.pyновый необязательный аргумент egress у record_failure/is_paused: ключ source_breaker:{site}:{egress}. Пауза ВСЕЙ площадки (старый ключ) взводится только когда proxy_limiter.has_healthy()==False И все маршрутные ключи на паузе. Вызовы в _record_source_signal передают проксю задачи
services/scrapers/encar_scraper.py и др.в исходах фетча добавить вызов record_result (успех/сеть/блок уже классифицированы: SourceBlockedError / ListingFetchError / OK)

Ключевые решения

  • AIMD стартует с сегодняшних капов и растёт только на длинном успехе — по построению не может стать агрессивнее прода дня 0; сигнал блока режет темп мгновенно вдвое. Это тот же принцип, которым TCP ищет ёмкость канала.
  • Существующие капы (antiban_detail_limit, PROXY_MAX_CONCURRENCY) остаются как верхняя страховка одновременности; AIMD управляет темпом (req/мин), а не одновременностью. Снимать капы — отдельное решение после недели данных AIMD.
  • Fail-open: Redis недоступен → pick_route возвращает случайный enabled-прокси (сегодняшнее поведение), record_result глохнет молча. Флот не может «сломать» парсинг сильнее его отсутствия.
  • DIRECT — тоже маршрут (egress='direct'): для che168 и фолбэков он проходит те же здоровье/темп/предохранитель — серверный IP перестаёт быть бесконтрольным.

Тесты, приёмка, откат

  • backend/test_proxy_fleet.py: AIMD-математика (рост/срез/пол/потолок), выбор маршрута при мёртвых, fail-open без Redis, маршрутный→площадочный каскад предохранителя.
  • Приёмка (72 ч): ноль фетчей через unhealthy-маршрут (лог pick_route); паузы всей площадки — только при ≥2 одновременно павших маршрутах; частота 30-минутных пауз ↓ до ~0; fetch_ms p90 encar/dongchedi без 20-секундных выбросов.
  • Откат: PROXY_FLEET=off — прежний random + прежний per-site предохранитель. Beat-проба безвредна и остаётся.
W5

Дельта-Свежие: деталь только новым машинам

~2–3 дня · риск продуктовый, за флагом

Мотивация. 39–54% машин живых поисков уже в БД с полной карточкой, но no_cache=True гонит все в детальный фетч. Листинг несёт живую цену и пробег → свежесть проверяема без детали; ТТХ/комплектация/осмотр неизменяемы.

Дизайн

Точка врезки одна: ветка кэша в on_page_parsed (tasks/distributed_parser_tasks.py:3435–3560). Сегодня: fresh → no_cache → cache_info={'exists': False}. Новый режим FRESH_MODE=delta (читается там же, где no_cache; выставляется в services/search_dispatch.py:149 вместо/рядом с no_cache только для классических «Свежих» — И5):

# on_page_parsed, per car — режим delta (псевдокод ветки):
cache_info = db.check_car_cache(site, site_car_id)      # кэш ЧИТАЕМ (не как сегодня)
if not cache_info['exists'] or not pipeline_complete(...).ok:
    → ветка «новая»: перевод basic + detail_queue          # как сегодня, без изменений
else:
    stored = cache_info['data']
    fresh_price   = car.get(price_field(site))   # price_cny|price_krw|price_eur|price_rub
    fresh_mileage = car.get('mileage_km')
    if prices_equal(fresh_price, stored) and mileage_close(fresh_mileage, stored):
        # === ДЕЛЬТА-ОТДАЧА: 0 внешних запросов ===
        card = stored
        card['card_status']    = 'full'          # НЕ 'from_cache' — инвариант И3
        card['detail_complete']= True
        card['data_freshness'] = now()           # листинг только что подтвердил живость
        card[price_field]      = fresh_price     # И2: листинг — источник цены/пробега
        card['mileage_km']     = fresh_mileage
        card['image_url']      = car.get('image_url') or card.get('image_url')
        stats['delta'] += 1                      # метрика приёмки
        if cache_age_days > FRESH_DELTA_REVALIDATE_DAYS (деф. 14):
            session_service.add_to_refresh_queue(session_id, site_car_id)  # фон
    else:
        # === ЦЕНА/ПРОБЕГ ИЗМЕНИЛИСЬ: локальный пересчёт, 0 внешних запросов ===
        card = stored + свежие price/mileage
        card['total_cost_rub']… = cost_calculator.calculate_total_cost(card)
            # hp/год/объём уже в stored — расчёт чисто локальный (курсы в кэше 4 ч)
        db.add_parsed_car(card)                  # дешёвый upsert (W1)
        session_service.add_to_refresh_queue(session_id, site_car_id)
            # фон докачает детально: вдруг сменились фото/описание вместе с ценой
    emit как обычную full-карточку (порядок И4 сохранён: позиция листинга)

Правила сравнения

  • prices_equal: строгое равенство в валюте площадки (цена — единственное поле, по которому продавцы «шевелят» объявление; листинг и деталь берут его из одного источника площадки — фальшразличий нет). НО у dongchedi листинговая цена бывает в ванях/фэнях — нормализация как в существующем merge (preserve_fields, :5307–5334) — переиспользовать её хелпер.
  • mileage_close: |Δ| ≤ 1 км (защита от округлений форматирования), иначе «изменилось».
  • Машина в статусах delisted / calc_hold / ordered — «нет в кэше» (check_car_cache уже так решает, :186–197) → полный перепарс. Не меняем.
  • Avito/auto.ru/drom: дельта-режим НЕ применяется в фазе 1 (РФ-детали дёшевы, прокси не тратят; включим при желании отдельным значением флага FRESH_MODE_SITES).

Коалесценция дублей (под-поток W5b, опционально)

15% живых сессий — точные дубли фильтров. В _dispatch_search_v2: перед try_admit посчитать md5(platform + canonical_filters); если сессия с тем же хэшем в статусе running моложе SEARCH_COALESCE_SEC (дефолт 120) — создать «инстант-сессию» юзера из её Redis-карточек (механизм уже есть у «Из БД»: run_db_search создаёт инстант-сессии) и подписать на её оставшиеся эмиты. Фаза 2, отдельный флаг SEARCH_COALESCE=on.

Флаг, тесты, приёмка, откат

  • Флаг FRESH_MODE: refetch (дефолт — сегодняшнее no_cache бит-в-бит) | delta. Кампании/агент/«Из БД»/refresh флаг не читают (И5).
  • Тест backend/test_fresh_delta.py: все четыре ветки (новая / совпало / цена изменилась / дырявый кэш), статусы карточки (И3), источник цены (И2), delisted→перепарс.
  • Теневая проверка до включения (приём «замер без мутаций»): на 20 живых сессиях сравнить дельта-карточку с реально перепарсенной — расхождения полей ≤ фото/описание.
  • Приёмка (неделя): stats['delta'] ≈ 30–50% машин; жалоб на устаревшие карточки нет; число детальных задач на сессию ↓ на ту же долю (Loki: задач process_single_car_detail / сессия).
  • Откат: FRESH_MODE=refetch, мгновенно.
W6

Details-воркер на пуле потоков

~1 неделя · единственный реальный техриск · только через канарейку

Мотивация. 24 prefork-процесса (идло 1,2 ГБ, под нагрузкой до 3,6 ГБ) выполняют задачи, которые почти целиком ждут сеть/БД/LLM. Пул потоков: та же работа в одном процессе ~0,6 ГБ; освобождённые ~3 ГБ уходят в книгу весов admission (одновременных поисков 11–15 → 25–35), слоты деталей 24 → 48–64 без роста памяти. Предпосылки проверены аудитом: curl_cffi отпускает GIL (8 потоков — wall 2,23 с), БД на потокобезопасном scoped_session, переводчик уже живёт в многопоточном gunicorn.

Изменения

ГдеЧто
backend/celery_entrypoint.shподдержка CELERY_POOL: при threads добавить --pool=threads и НЕ передавать --max-memory-per-child (в threads не работает); jemalloc оставить (полезен и одному процессу)
docker-compose.ymlканарейка: сервис celery_worker_details_threads — копия details-сервиса с CELERY_POOL=threads, CELERY_DETAILS_CONCURRENCY=16, DB_POOL_SIZE=8, DB_MAX_OVERFLOW=8 (потоки делят ОДИН engine — суммарно коннектов меньше, чем у prefork 25×4); на время канарейки у основного details-воркера CELERY_DETAILS_CONCURRENCY=12 (сумма слотов ≈ прежняя, обе консюмят detail_fetch, -Ofair уже стоит)
services/capacity_service.pymax_concurrent_count() и global_detail_budget() читают CELERY_DETAILS_CONCURRENCY — на канарейке передать суммарное число слотов через существующий env (12+16=28), после полного перехода — одно число потоков
плановый рестартвместо max-memory-per-child: ежесуточный docker restart контейнера в ночное окно (окно уже существует у кампаний) — страховка от накопительных утечек одного процесса

Чек-лист потокобезопасности (ревью перед канарейкой)

  • Синглтоны: get_translation_service_live (локи есть), _get_cost_calculator, get_capacity_service — проверить инициализацию под конкуренцией (double-checked с локом или прогрев в worker_process_init — коммит 8cb1beff уже сделал прогрев переводчика).
  • Скраперы: инстанцируются per-task (не шарятся) — ок; module-level кэши (cookie_jar_store, jar-спеки) — уже под single-flight локами по построению.
  • Глобальные мутации: os.environ — писателей в горячем пути нет (grep); random — потокобезопасен; signal в задачах — не используется.
  • Playwright в деталях: fallback существует у dongchedi — под threads браузер запрещаем флагом (fallback вернёт ошибку фетча → обычный ретрай через prefork-воркер, пока канарейка частичная; после полного перехода dongchedi-fallback переносится в listing-воркер, где браузер уже живёт).

План канарейки (2 недели)

  1. Дни 1–2: канарейка 16 потоков + основной 12 prefork. Сравнение в Loki по container: total_ms, доля failed, RSS обоих контейнеров.
  2. Дни 3–7: при чистых метриках канарейка 32, основной 6. Нагрузочный прогон стрессом из очереди III (реплей-стенд уже есть).
  3. Неделя 2: полный переход (details=threads 48, prefork-сервис выключен), CELERY_DETAILS_CONCURRENCY=48, поднять MAX_CONCURRENT_SEARCHES к книге весов. Priority-воркер переводится тем же способом после недели стабильности.

Критерии отката (любой — шаг назад)

  • Доля failed задач канарейки > базлайна +1 п.п.; любой дедлок/зависание задач > time_limit; RSS канарейки > 1,5 ГБ; рост p95 total_ms > 20%.
  • Откат = убрать сервис из compose и вернуть CELERY_DETAILS_CONCURRENCY=24. Кода это не откатывает — entrypoint-изменение обратносовместимо.
07 · Альтернативы

Рассмотрено и отвергнуто

АльтернативаПочему нет
Купить прокси/SSD сразуВне мандата («без закупок»); главное — сегодняшняя архитектура выжмет из покупки четверть: без W4 новые прокси так же умирают невидимо, без W5 тратятся на повторы. Программа — фундамент под будущие закупки.
gevent вместо threads (W6)curl_cffi — блокирующие C-вызовы: monkey-patch их не кооперативит, event-loop встанет. Потоки с отпущенным GIL — проверено экспериментом; префорк-семантика кода сохраняется почти целиком.
Разрезать задачу детали на fetch-задачу и process-задачу (конвейер из двух очередей)Даёт тот же эффект, что W1+W2, но требует сериализации car_data между задачами (сотни КБ через брокер), новых очередей и ломает autodrive/finalize-логику. W1+W2 достигают цели двумя локальными врезками.
synchronous_commit=off глобально для БДУбирает durability у заявок/юзеров/платежей ради выигрыша только в parsed_cars. SET LOCAL точечно — тот же выигрыш без расширения зоны риска.
UNLOGGED-таблица для parsed_carsТаблица теряется целиком при crash (480 К машин), миграция тяжёлая. Несопоставимый риск ради того же порядка выигрыша, что async commit.
Отдать «Свежие» целиком из БД + фоновый полный репарс (SWR)Нарушает продуктовую семантику: в выдачу попали бы машины, которых УЖЕ нет в листинге (проданные). Дельта-режим W5 строже: показываются только машины, подтверждённые живым листингом сессии.
Убрать 30-минутную паузу предохранителяПауза — правильная защита от настоящего бана. Проблема не в паузе, а в ложных срабатываниях от мёртвого прокси; W4 чинит причину, каскад остаётся последним рубежом.
08 · Выкатка

Фазы и гейты

Каждая фаза — гейт: метрики фазы зелёные ≥ указанного срока, иначе шаг назад (одна env-переменная). Порядок минимизирует связность: W1–W4 не зависят друг от друга; W5 желательно после W4 (фон-refresh нагружает пул — пусть пул уже управляем); W6 — последним, на уже похудевших задачах.

ФазаСоставСрок наблюденияГейт (все условия)Ожидаемый потолок после
Ф0W0 (однострочники)2–3 дняпредохранитель молчит · очередь v2 пуста в пиках · MemAvailable > 1 ГБ~300/ч
Ф1W1 (async commit) + W3 (encar-диета)48 чsave_db p50 < 1 с · fetch encar p50 < 3,5 с · ошибок сохранения 0~400/ч
Ф2W4 (флот) + W2 (LLM вне слота)72 чноль фетчей через unhealthy · cost_ms p90 < 3 с · пауз площадок 0~450/ч
Ф3W5 (дельта-свежие)неделяdelta-доля 30–50% · расхождений карточек нет · детальных задач/сессию −30%+~600/ч
Ф4W6 (threads-канарейка → переход)2 неделикритерии W6 · MAX_CONCURRENT_SEARCHES к книге весов700+/ч, 25–35 одновр.

Ожидаемые потолки — оценки из аудита (формула: слоты × 3600 ÷ слот-время ÷ машин на сессию, с поправкой на дельта-долю); гейты сформулированы по наблюдаемым метрикам, а не по этим оценкам. Суммарная трудоёмкость: ~2 недели чистой разработки + ~3 недели календарного наблюдения (перекрываются).

Выкатка каждой фазы — по существующему ручному порядку (CI стоит из-за биллинга): rsync ветки, целевые docker compose build/up изменённых сервисов, смоук на бою. Порядок и команды — как в журнале очереди III (выкатка 27.08). Миграций Alembic в программе НЕТ ни одной — только код, env и две команды ALTER SYSTEM/INSERT в PG.
09 · Наблюдаемость

Метрики, дашборд, откат

Почти всё уже измеряется: поэтапные тайминги (fetch_ms/translate_ms/cost_ms/save_db_ms/total_ms + site_source) пишутся каждой задачей в Loki — базлайн аудита снят именно с них, эффект каждой фазы виден теми же запросами. Добавляются только:

Единый экран здоровья (Loki-запросы, сохранить в закладки/Grafana)

МетрикаЗапрос (суть)Зелёное
Слот-время по площадкамquantile_over_time(0.5, … total_ms по site_source)dcd < 6 c · encar < 6 с · che168 < 4 с
Запись в БДp50 save_db_ms< 1 с все площадки
Паузы площадокcount "поставлена на паузу" (admin_events)0/сутки
Здоровье флотаfleet_health_probe: alive по каждому IP≥ 3 из 4 + direct
Дельта-доляsum stats.delta / sum stats.total30–50%
Очередь admissionparsing_queue_v2 глубина в пиках0–2
ПамятьMemAvailable в пиках; RSS details-контейнера> 1 ГБ; threads < 1,5 ГБ

Сводная таблица флагов (аварийная шпаргалка)

ФлагПотокДефолтОткат
PARSED_CAR_ASYNC_COMMITW10=0 мгновенно
HP_ASYNC_ENRICHW20=0 мгновенно
ENCAR_PARALLEL_DETAILW30=0 мгновенно
PROXY_FLEETW4off=off мгновенно
FRESH_MODEW5refetch=refetch мгновенно
SEARCH_COALESCEW5boff=off мгновенно
CELERY_POOLW6preforkубрать канарейку из compose
10 · Вопросы

Открытые вопросы владельцу

  1. Мёртвый .165: написать провайдеру (адрес оплачен) или просто выключить строку? W0.2 предполагает выключить сейчас, вернуть по зелёной пробе.
  2. Порог доверия дельте (W5): предлагаемый ревалидационный срок 14 дней (карточка старше — фоновый refresh даже при совпавшей цене). Устраивает? Можно 7 — дороже, свежее описания/фото.
  3. Ночное окно рестарта threads-воркера (W6): использовать существующее окно кампаний (04:00–06:00 МСК)?
  4. Т2.2 темп подъёма потолка: 5→8 сразу в Ф0 и 8→12 после Ф1 — ок, или консервативнее?
  5. W5b (коалесценция дублей): делать в этой программе или отложить (выигрыш ~15%, сложность средняя)?