fix(tradein/scrapers): ABORT-лог называл серию блоков, хотя рвал прогон по доле #3201
3 changed files with 186 additions and 7 deletions
|
|
@ -78,6 +78,16 @@ class BlockRatioBreaker:
|
|||
"""Для логов ABORT/BLOCKED -- та же цифра, что раньше выводилась в consecutive=%d."""
|
||||
return self._consecutive_blocks
|
||||
|
||||
@property
|
||||
def window_blocks(self) -> int:
|
||||
"""Числитель ratio-критерия: сколько блоков в текущем окне."""
|
||||
return sum(self._window)
|
||||
|
||||
@property
|
||||
def window_len(self) -> int:
|
||||
"""Знаменатель ratio-критерия: сколько попыток уже влезло в окно."""
|
||||
return len(self._window)
|
||||
|
||||
def record_block(self) -> None:
|
||||
self._consecutive_blocks += 1
|
||||
self._window.append(True)
|
||||
|
|
@ -108,7 +118,12 @@ class BlockRatioBreaker:
|
|||
self.streak_histogram[self._consecutive_blocks] += 1
|
||||
self._consecutive_blocks = 0
|
||||
|
||||
def should_abort(self) -> bool:
|
||||
def abort_reason(self) -> str | None:
|
||||
"""Какой критерий требует обрыва прямо сейчас, или None.
|
||||
|
||||
Возвращает "safety_net" / "ratio" / None. Состояние не меняет, поэтому
|
||||
вызывать можно сколько угодно раз -- в том числе повторно, ради текста лога.
|
||||
"""
|
||||
# Safety-net -- ТОЛЬКО когда ratio-критерий физически недостижим (снапшот
|
||||
# короче окна), иначе пачка safety_min блоков в начале длинного прогона
|
||||
# абортила бы его так же, как до правки (#3184 review MAJOR 2).
|
||||
|
|
@ -117,12 +132,40 @@ class BlockRatioBreaker:
|
|||
and self._pure_block_run
|
||||
and self._consecutive_blocks >= self.safety_min
|
||||
):
|
||||
return True
|
||||
return "safety_net"
|
||||
if len(self._window) == self.window_size and self.window_size > 0:
|
||||
ratio = sum(self._window) / self.window_size
|
||||
if ratio >= self.ratio_threshold:
|
||||
return True
|
||||
return False
|
||||
return "ratio"
|
||||
return None
|
||||
|
||||
def should_abort(self) -> bool:
|
||||
return self.abort_reason() is not None
|
||||
|
||||
def abort_explanation(self) -> str:
|
||||
"""Текст для лога ABORT -- ровно та величина, по которой обрыв и произошёл.
|
||||
|
||||
До этого лог печатал "%d consecutive blocks" всегда, а критерий с #3184
|
||||
стал ratio: прогон 5210 (14 блоков из 20, обрыв ровно по порогу 0.7) выдал
|
||||
"ABORT -- 1 consecutive blocks", потому что в момент срабатывания текущая
|
||||
серия равнялась единице. Число верное, величина не та -- читатель лога видит
|
||||
цифру, по которой обрыва быть не могло, и идёт искать несуществующий баг.
|
||||
|
||||
Пустая строка означает "рвать не по чему" -- вызывается только под
|
||||
should_abort(), так что в логи не попадает.
|
||||
"""
|
||||
reason = self.abort_reason()
|
||||
if reason == "ratio":
|
||||
return (
|
||||
f"доля блоков {self.window_blocks}/{self.window_len} в окне "
|
||||
f"(порог {self.ratio_threshold * 100:.0f}%)"
|
||||
)
|
||||
if reason == "safety_net":
|
||||
return (
|
||||
f"{self._consecutive_blocks} блоков подряд без единого успеха "
|
||||
f"(снапшот {self.snapshot_size} короче окна {self.window_size})"
|
||||
)
|
||||
return ""
|
||||
|
||||
def finalize(self) -> dict[str, int]:
|
||||
"""Досчитать хвостовую пачку (прогон оборвался посреди серии блоков, без
|
||||
|
|
|
|||
|
|
@ -37,6 +37,11 @@ block_ratio_window/_threshold) ИЛИ safety-net для коротких про
|
|||
consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое
|
||||
поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в
|
||||
counters["block_streak_histogram"], иначе эффект правки нечем измерить постфактум.
|
||||
Какой из двух критериев сработал -- в counters["abort_reason"] ("ratio" /
|
||||
"safety_net"); ключа нет, если прогон не обрывался. Он же называется в ABORT-логе
|
||||
своей величиной: доля печатает "14/20", safety-net -- длину серии. Раньше лог
|
||||
печатал серию всегда, и прогон 5210 (обрыв по доле 14/20) отчитался как
|
||||
"ABORT -- 1 consecutive blocks".
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -431,6 +436,7 @@ async def run_avito_detail_backfill(
|
|||
)
|
||||
consecutive_failures = 0
|
||||
aborted_by_blocks = False
|
||||
abort_reason: str | None = None
|
||||
do_sleep = False
|
||||
items_since_warm = 0
|
||||
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
|
||||
|
|
@ -632,12 +638,18 @@ async def run_avito_detail_backfill(
|
|||
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate.
|
||||
# #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни
|
||||
# одного успеха, накопилось max_consecutive_blocks) -- см. module docstring.
|
||||
if breaker.should_abort():
|
||||
abort_reason = breaker.abort_reason()
|
||||
if abort_reason is not None:
|
||||
# Печатаем ту величину, по которой обрыв и произошёл. Раньше здесь
|
||||
# всегда стояло "%d consecutive blocks", хотя критерий с #3184 стал
|
||||
# ratio: прогон 5210 (14 блоков из 20, ровно порог) отпечатал
|
||||
# "ABORT -- 1 consecutive blocks" -- текущая серия в тот момент
|
||||
# действительно равнялась единице, но обрыв был не по ней.
|
||||
logger.error(
|
||||
"avito_detail_backfill: run_id=%d ABORT -- %d consecutive blocks, "
|
||||
"avito_detail_backfill: run_id=%d ABORT -- %s, "
|
||||
"частая причина: %s. enriched=%d attempted=%d",
|
||||
run_id,
|
||||
breaker.consecutive_blocks,
|
||||
breaker.abort_explanation(),
|
||||
_top_failure(failure_census) or "причина не определена",
|
||||
counters.enriched,
|
||||
counters.attempted,
|
||||
|
|
@ -739,6 +751,11 @@ async def run_avito_detail_backfill(
|
|||
streak_histogram = breaker.finalize()
|
||||
if streak_histogram:
|
||||
current_counters["block_streak_histogram"] = streak_histogram # type: ignore[assignment]
|
||||
# Какой из двух критериев оборвал прогон -- в counters, а не только в логах:
|
||||
# разбор простоя идёт SQL-запросом по scrape_runs, а не грепом контейнера,
|
||||
# и без этого ключа "banned" опять не отличить по причине (#3178).
|
||||
if abort_reason is not None:
|
||||
current_counters["abort_reason"] = abort_reason # type: ignore[assignment]
|
||||
runs_mod.mark_backfill_finished(
|
||||
db,
|
||||
run_id,
|
||||
|
|
|
|||
|
|
@ -1308,3 +1308,122 @@ async def test_backfill_ratio_boundary_13_of_20_does_not_abort() -> None:
|
|||
runs.mark_backfill_finished.assert_called_once()
|
||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_ratio_abort_log_names_the_ratio_not_the_streak(caplog: Any) -> None:
|
||||
"""ABORT по доле называет долю, а не текущую серию (#3184 follow-up).
|
||||
|
||||
Воспроизводит прод-прогон 5210: 14 блоков из 20, обрыв ровно по порогу 0.7, но
|
||||
ПОСЛЕДНЯЯ серия блоков к этому моменту равна единице (позиция 19 -- успех,
|
||||
позиция 20 -- блок). Лог печатал "ABORT -- 1 consecutive blocks": число верное,
|
||||
величина не та. Читатель видит цифру, по которой обрыва быть не могло.
|
||||
"""
|
||||
total = 20
|
||||
# Раскладка ровно как в прогоне 5210.
|
||||
block_positions = {2, 3, 5, 6, 7, 8, 9, 11, 12, 13, 14, 16, 17, 20}
|
||||
assert len(block_positions) == 14
|
||||
assert 19 not in block_positions and 20 in block_positions, "серия на обрыве = 1"
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
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),
|
||||
caplog.at_level("ERROR"),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=110, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.blocked == 14
|
||||
abort_records = [r.message for r in caplog.records if "ABORT" in r.message]
|
||||
assert abort_records, "ожидался ABORT-лог"
|
||||
msg = abort_records[0]
|
||||
assert "14/20" in msg, msg
|
||||
assert "consecutive" not in msg, msg
|
||||
assert "1 блоков подряд" not in msg, msg
|
||||
counters = runs.mark_backfill_finished.call_args.args[2]
|
||||
assert counters["abort_reason"] == "ratio"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_safety_net_abort_log_names_the_streak(caplog: Any) -> None:
|
||||
"""ABORT по safety-net называет серию -- ту величину, по которой он и сработал.
|
||||
|
||||
Зеркало предыдущего теста: снапшот короче окна, ratio недостижим, рвёт серия.
|
||||
Если бы обе ветки печатали один и тот же текст, тест выше проходил бы и у
|
||||
сломанного сообщения.
|
||||
"""
|
||||
total = 5 # короче окна 20 -- ratio-критерий физически недостижим
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=AvitoBlockedError("ip blocked"))
|
||||
fake_settings = _fake_settings(
|
||||
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),
|
||||
caplog.at_level("ERROR"),
|
||||
):
|
||||
result = await run_avito_detail_backfill(
|
||||
db, run_id=111, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
assert result.blocked == 5
|
||||
abort_records = [r.message for r in caplog.records if "ABORT" in r.message]
|
||||
assert abort_records, "ожидался ABORT-лог"
|
||||
msg = abort_records[0]
|
||||
assert "5 блоков подряд" in msg, msg
|
||||
assert "доля блоков" not in msg, msg
|
||||
counters = runs.mark_backfill_finished.call_args.args[2]
|
||||
assert counters["abort_reason"] == "safety_net"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_backfill_no_abort_leaves_no_abort_reason_in_counters() -> None:
|
||||
"""Прогон без обрыва не пишет abort_reason -- ключ означает 'оборвались', а не
|
||||
'считали критерий'."""
|
||||
total = 20
|
||||
block_positions = {8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20} # 13/20 = 0.65
|
||||
assert len(block_positions) == 13
|
||||
snapshot = _make_snapshot(total)
|
||||
db = _mock_db(snapshot)
|
||||
runs = MagicMock()
|
||||
mock_fetch = AsyncMock(side_effect=_spaced_outcomes(total, block_positions))
|
||||
fake_settings = _fake_settings(
|
||||
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),
|
||||
):
|
||||
await run_avito_detail_backfill(
|
||||
db, run_id=112, params={"batch_size": total, "budget_sec": 3600}
|
||||
)
|
||||
|
||||
counters = runs.mark_backfill_finished.call_args.args[2]
|
||||
assert "abort_reason" not in counters
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue