fix(tradein/avito): серия отказов обрывается и называет причину, а не выедает бюджет (#2674) (#2739)
All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m3s
Deploy Trade-In / build-backend (push) Successful in 57s
Deploy Trade-In / deploy (push) Successful in 1m31s

This commit is contained in:
bot-backend 2026-08-06 15:29:48 +00:00
parent 4aec49f7fb
commit 0dc6f12630
3 changed files with 175 additions and 4 deletions

View file

@ -578,6 +578,7 @@ def mark_backfill_finished(
*, *,
source: str, source: str,
aborted_by_blocks: bool = False, aborted_by_blocks: bool = False,
fail_hint: str | None = None,
) -> None: ) -> None:
"""Честный финал detail-backfill'а (#2674): нулевой прогон ≠ 'done'. """Честный финал detail-backfill'а (#2674): нулевой прогон ≠ 'done'.
@ -601,11 +602,19 @@ def mark_backfill_finished(
`gone` (404 у avito) считается результатом наравне с `enriched`: прогон, `gone` (404 у avito) считается результатом наравне с `enriched`: прогон,
который подтвердил снятие объявлений, работу сделал. который подтвердил снятие объявлений, работу сделал.
`fail_hint` самая частая причина отказа этого прогона (задача считает её сама,
см. avito_detail_backfill._failure_signature). Дописывается в текст статуса,
потому что «blocked=5, обогащено 0» не отвечает на единственный вопрос, ради
которого статус и читают: отказала площадка или наш тракт (#2686, #2698). Логи
контейнера на этот вопрос отвечать не могут они исчезают при пересоздании
контейнера, то есть на первом же деплое после ночного прогона.
""" """
attempted = int(counters.get("attempted") or 0) attempted = int(counters.get("attempted") or 0)
enriched = int(counters.get("enriched") or 0) enriched = int(counters.get("enriched") or 0)
blocked = int(counters.get("blocked") or 0) blocked = int(counters.get("blocked") or 0)
produced = enriched + int(counters.get("gone") or 0) produced = enriched + int(counters.get("gone") or 0)
hint = f"; причина: {fail_hint}" if fail_hint else ""
if attempted == 0: if attempted == 0:
mark_done(db, run_id, counters) mark_done(db, run_id, counters)
@ -614,7 +623,7 @@ def mark_backfill_finished(
if blocked and (aborted_by_blocks or produced == 0): if blocked and (aborted_by_blocks or produced == 0):
reason = ( reason = (
f"backfill-honest-status: {source} остановлен блоками источника — " f"backfill-honest-status: {source} остановлен блоками источника — "
f"blocked={blocked}, обогащено {enriched} из {attempted} попыток (#2674)" f"blocked={blocked}, обогащено {enriched} из {attempted} попыток{hint} (#2674)"
) )
logger.error("%s run_id=%d", reason, run_id) logger.error("%s run_id=%d", reason, run_id)
mark_banned(db, run_id, reason, counters) mark_banned(db, run_id, reason, counters)
@ -624,7 +633,7 @@ def mark_backfill_finished(
reason = ( reason = (
f"backfill-honest-status: {source} без результата — 0 обогащено из " f"backfill-honest-status: {source} без результата — 0 обогащено из "
f"{attempted} попыток (failed={counters.get('failed', 0)}, " f"{attempted} попыток (failed={counters.get('failed', 0)}, "
f"blocked={blocked}) (#2674)" f"blocked={blocked}){hint} (#2674)"
) )
logger.error("%s run_id=%d", reason, run_id) logger.error("%s run_id=%d", reason, run_id)
mark_failed(db, run_id, reason, counters) mark_failed(db, run_id, reason, counters)

View file

@ -13,6 +13,13 @@ session path as the detail-phase of `run_avito_city_sweep`
rotate IP on every block, abort after max_consecutive_blocks. Статус оборванного rotate IP on every block, abort after max_consecutive_blocks. Статус оборванного
блоками прогона 'banned' (#2674, runs.mark_backfill_finished): работу он не блоками прогона 'banned' (#2674, runs.mark_backfill_finished): работу он не
доделал, остаток снапшота уедет в следующую ночь через NULL detail_enriched_at. доделал, остаток снапшота уедет в следующую ночь через NULL detail_enriched_at.
Отказы, не являющиеся блоками, до 2026-08-06 брейкера не имели вовсе: прогоны
3-5 августа делали ~1600 попыток, получали 1600 отказов, ноль обогащений и
выедали весь бюджет (9000 с) вместе с 1600 запросами через единственный прокси.
Теперь такая серия обрывается по max_consecutive_failures, а самая частая причина
отказа пишется в текст статуса прогона (_failure_signature) иначе она живёт
только в логах контейнера, а те исчезают на первом же деплое.
""" """
from __future__ import annotations from __future__ import annotations
@ -20,7 +27,9 @@ from __future__ import annotations
import asyncio import asyncio
import logging import logging
import random import random
import re
import time import time
from collections import Counter
from dataclasses import dataclass, field from dataclasses import dataclass, field
from urllib.parse import urlparse from urllib.parse import urlparse
@ -99,6 +108,38 @@ _OBLAST_AVITO_URL_PATTERNS = tuple(
) )
# Причина отказа карточки без её URL: 1576 отказов одного прогона должны схлопнуться
# в ОДНУ строку, иначе перепись бесполезна.
_URL_IN_MESSAGE_RE = re.compile(r"https?://\S+")
def _failure_signature(exc: BaseException) -> str:
"""Подпись причины отказа: тип исключения + текст без URL.
Зачем (замер 2026-08-06): у прогонов 3 и 4 августа counters говорили
`attempted=1576, failed=1576, blocked=0` и ничего больше. Кто отказал,
площадка или наш тракт, было видно ТОЛЬКО в логах контейнера, а тот
пересоздаётся на каждом деплое и уносит их с собой; в GlitchTip попадают
события уровня ERROR, а поштучные отказы WARNING. Разница между этими
двумя диагнозами разные владельцы задачи (#2686, #2698), поэтому она
обязана переживать перезапуск контейнера, то есть лежать в самом прогоне.
Тип исключения первый разряд диагноза (AvitoBlockedError = площадка
показала 403/firewall; сетевой класс curl_cffi = наш прокси-тракт;
ValueError = ответ пришёл, но не разобран), текст второй.
"""
message = _URL_IN_MESSAGE_RE.sub("<url>", str(exc)).strip()
return f"{type(exc).__name__}: {message}"[:160] if message else type(exc).__name__
def _top_failure(census: Counter[str]) -> str | None:
"""Самая частая причина отказа с её долей; None — отказов не было."""
if not census:
return None
reason, hits = census.most_common(1)[0]
return f"{reason} ({hits} из {sum(census.values())})"
@dataclass @dataclass
class AvitoDetailBackfillResult: class AvitoDetailBackfillResult:
"""Counters for one backfill run.""" """Counters for one backfill run."""
@ -139,6 +180,8 @@ async def run_avito_detail_backfill(
budget_sec: float -- wall-clock budget per run, default 3600s. budget_sec: float -- wall-clock budget per run, default 3600s.
request_delay_sec: float -- delay between listings, default 6.0s. request_delay_sec: float -- delay between listings, default 6.0s.
max_consecutive_blocks: int -- abort threshold, default 5. max_consecutive_blocks: int -- abort threshold, default 5.
max_consecutive_failures: int -- порог обрыва по отказам-не-блокам,
default 25 (см. комментарий у чтения параметра ниже).
Lifecycle: update_heartbeat -> snapshot -> loop with budget guard -> Lifecycle: update_heartbeat -> snapshot -> loop with budget guard ->
mark_backfill_finished (done / banned при блоках / failed при нуле, #2674); mark_backfill_finished (done / banned при блоках / failed при нуле, #2674);
@ -149,6 +192,13 @@ async def run_avito_detail_backfill(
budget_sec = float(params.get("budget_sec", 3600)) budget_sec = float(params.get("budget_sec", 3600))
request_delay_sec = float(params.get("request_delay_sec", 6.0)) request_delay_sec = float(params.get("request_delay_sec", 6.0))
max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5)) max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5))
# Брейкер на отказы-НЕ-блоки. Блоки свой брейкер имели с самого начала, отказы —
# нет, и это стоило трёх ночей подряд: 3-5 августа прогон делал ~1600 попыток,
# получал 1600 отказов, ноль обогащений и выедал весь бюджет 9000 с (плюс 1600
# запросов через единственный прокси, #2638). Порог заметно выше блочного: пачка
# мёртвых карточек (404 → ValueError в curl-режиме) не должна обрывать здоровый
# прогон, а 25 отказов подряд без единого успеха — уже не невезение.
max_consecutive_failures = int(params.get("max_consecutive_failures", 25))
warm_batch = int(params.get("warm_batch", 500)) warm_batch = int(params.get("warm_batch", 500))
research_every = int(params.get("research_every", 50)) research_every = int(params.get("research_every", 50))
block_cooldown_sec = float(params.get("block_cooldown_sec", 30.0)) block_cooldown_sec = float(params.get("block_cooldown_sec", 30.0))
@ -308,9 +358,13 @@ async def run_avito_detail_backfill(
) )
consecutive_blocks = 0 consecutive_blocks = 0
consecutive_failures = 0
aborted_by_blocks = False aborted_by_blocks = False
do_sleep = False do_sleep = False
items_since_warm = 0 items_since_warm = 0
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
# в отличие от логов; см. _failure_signature.
failure_census: Counter[str] = Counter()
for idx, row in enumerate(snapshot): for idx, row in enumerate(snapshot):
# Budget guard # Budget guard
@ -429,8 +483,9 @@ async def run_avito_detail_backfill(
if use_curl: if use_curl:
items_since_warm += 1 items_since_warm += 1
consecutive_blocks = 0 consecutive_blocks = 0
consecutive_failures = 0
except AvitoListingGoneError: except AvitoListingGoneError as gone_exc:
# #2034: мёртвый листинг (404 / removed) — НЕ блок, НЕ failed. # #2034: мёртвый листинг (404 / removed) — НЕ блок, НЕ failed.
# Координатные дыры в lat-null очереди в основном dead-листинги; # Координатные дыры в lat-null очереди в основном dead-листинги;
# browser-mode рендерит их «Ошибка 404» без item-view → раньше это # browser-mode рендерит их «Ошибка 404» без item-view → раньше это
@ -440,6 +495,10 @@ async def run_avito_detail_backfill(
# и не сбрасываем). Метим is_active=FALSE → листинг уходит из scope # и не сбрасываем). Метим is_active=FALSE → листинг уходит из scope
# (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь. # (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
counters.gone += 1 counters.gone += 1
# 404 — честный ответ площадки, значит тракт цел: серия отказов
# прерывается (блочный брейкер 404 не трогает — см. #2034).
consecutive_failures = 0
failure_census[_failure_signature(gone_exc)] += 1
try: try:
with db.begin_nested(): with db.begin_nested():
db.execute( db.execute(
@ -481,6 +540,7 @@ async def run_avito_detail_backfill(
except (AvitoBlockedError, AvitoRateLimitedError) as e: except (AvitoBlockedError, AvitoRateLimitedError) as e:
consecutive_blocks += 1 consecutive_blocks += 1
counters.blocked += 1 counters.blocked += 1
failure_census[_failure_signature(e)] += 1
do_sleep = False do_sleep = False
logger.warning( logger.warning(
"avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s", "avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s",
@ -538,12 +598,14 @@ async def run_avito_detail_backfill(
exc_info=True, exc_info=True,
) )
except TimeoutError: except TimeoutError as e:
# asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias). # asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias).
# Ловим ДО общего Exception (TimeoutError ⊂ OSError ⊂ Exception). Зависший # Ловим ДО общего Exception (TimeoutError ⊂ OSError ⊂ Exception). Зависший
# fetch отменён → листинг failed, переходим к следующему (loop не зависает, # fetch отменён → листинг failed, переходим к следующему (loop не зависает,
# run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем. # run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем.
counters.failed += 1 counters.failed += 1
consecutive_failures += 1
failure_census[_failure_signature(e)] += 1
logger.warning( logger.warning(
"avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip", "avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip",
run_id, run_id,
@ -557,6 +619,8 @@ async def run_avito_detail_backfill(
except Exception as e: except Exception as e:
counters.failed += 1 counters.failed += 1
consecutive_failures += 1
failure_census[_failure_signature(e)] += 1
logger.warning( logger.warning(
"avito_detail_backfill: run_id=%d listing %s failed: %s", "avito_detail_backfill: run_id=%d listing %s failed: %s",
run_id, run_id,
@ -568,6 +632,18 @@ async def run_avito_detail_backfill(
except Exception: except Exception:
pass pass
if consecutive_failures >= max_consecutive_failures:
logger.error(
"avito_detail_backfill: run_id=%d ABORT -- %d отказов подряд без "
"единого успеха, частая причина: %s. enriched=%d attempted=%d",
run_id,
consecutive_failures,
_top_failure(failure_census) or "неизвестна",
counters.enriched,
counters.attempted,
)
break
if counters.attempted % 25 == 0: if counters.attempted % 25 == 0:
current_counters = counters.to_dict() current_counters = counters.to_dict()
runs_mod.update_heartbeat(db, run_id, current_counters) runs_mod.update_heartbeat(db, run_id, current_counters)
@ -580,6 +656,7 @@ async def run_avito_detail_backfill(
current_counters, current_counters,
source="avito_detail_backfill", source="avito_detail_backfill",
aborted_by_blocks=aborted_by_blocks, aborted_by_blocks=aborted_by_blocks,
fail_hint=_top_failure(failure_census),
) )
logger.info( logger.info(
"avito_detail_backfill: run_id=%d FINISHED -- attempted=%d enriched=%d " "avito_detail_backfill: run_id=%d FINISHED -- attempted=%d enriched=%d "

View file

@ -790,3 +790,88 @@ async def test_backfill_use_curl_block_cooldown_research_no_rebuild() -> None:
mock_scraper.return_value._rotate_ip.assert_not_called() mock_scraper.return_value._rotate_ip.assert_not_called()
runs.mark_backfill_finished.assert_called_once() runs.mark_backfill_finished.assert_called_once()
runs.mark_failed.assert_not_called() runs.mark_failed.assert_not_called()
# ---------------------------------------------------------------------------
# Отказы-не-блоки: брейкер + перепись причин (2026-08-06)
# ---------------------------------------------------------------------------
@pytest.mark.asyncio
async def test_backfill_aborts_on_consecutive_failures_and_names_the_reason() -> None:
"""Серия отказов без единого успеха обрывается, а причина попадает в статус прогона.
Прод 3-5 августа: attempted1600, failed1600, blocked=0, enriched=0, весь
бюджет 9000 с и 1600 запросов через единственный прокси и ни слова о том,
ЧТО именно отказало (поштучные отказы логируются WARNING, а логи контейнера
пропадают на первом деплое). Брейкера на отказы-не-блоки не было вовсе.
"""
snapshot = _make_snapshot(200)
db = _mock_db(snapshot)
runs = MagicMock()
# Тот же класс отказа, что видели у соседнего свипа в тот же день.
mock_fetch = AsyncMock(
side_effect=OSError(
"Failed to perform, curl: (56) CONNECT tunnel failed, response 502. "
"See https://curl.se/libcurl/c/libcurl-errors.html"
)
)
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SLEEP, new_callable=AsyncMock),
):
result = await run_avito_detail_backfill(
db,
run_id=77,
params={
"batch_size": 200,
"budget_sec": 3600,
"max_consecutive_failures": 25,
},
)
assert result.attempted == 25, "серия отказов обязана обрываться, а не выедать бюджет"
assert result.failed == 25
assert result.blocked == 0
hint = runs.mark_backfill_finished.call_args.kwargs["fail_hint"]
assert hint is not None
assert "OSError" in hint # тип исключения = кому принадлежит отказ
assert "CONNECT tunnel failed" in hint
assert "25 из 25" in hint # доля, а не единичный пример
assert "https://curl.se" not in hint # URL вырезан, иначе 1600 «разных» причин
@pytest.mark.asyncio
async def test_backfill_success_resets_failure_streak() -> None:
"""Успех между отказами обнуляет серию — здоровый прогон брейкер не трогает."""
snapshot = _make_snapshot(5)
db = _mock_db(snapshot)
runs = MagicMock()
boom = ValueError("avito detail HTTP 500 for https://www.avito.ru/x")
# 2 отказа, успех, 2 отказа — при пороге 3 ни одна серия его не достигает.
mock_fetch = AsyncMock(side_effect=[boom, boom, MagicMock(), boom, boom])
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=True)
with (
patch(_SETTINGS, fake_settings),
patch(_SESSION, return_value=AsyncMock()),
patch(_SCRAPER),
patch(_RUNS, runs),
patch(_FETCH, mock_fetch),
patch(_SAVE, return_value=True),
patch(_SLEEP, new_callable=AsyncMock),
):
result = await run_avito_detail_backfill(
db,
run_id=78,
params={"batch_size": 5, "budget_sec": 3600, "max_consecutive_failures": 3},
)
assert result.attempted == 5
assert result.enriched == 1
assert result.failed == 4