fix(tradein/scrapers): ABORT-лог называл серию блоков, хотя рвал прогон по доле #3201
3 changed files with 186 additions and 7 deletions
|
|
@ -78,6 +78,16 @@ class BlockRatioBreaker:
|
||||||
"""Для логов ABORT/BLOCKED -- та же цифра, что раньше выводилась в consecutive=%d."""
|
"""Для логов ABORT/BLOCKED -- та же цифра, что раньше выводилась в consecutive=%d."""
|
||||||
return self._consecutive_blocks
|
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:
|
def record_block(self) -> None:
|
||||||
self._consecutive_blocks += 1
|
self._consecutive_blocks += 1
|
||||||
self._window.append(True)
|
self._window.append(True)
|
||||||
|
|
@ -108,7 +118,12 @@ class BlockRatioBreaker:
|
||||||
self.streak_histogram[self._consecutive_blocks] += 1
|
self.streak_histogram[self._consecutive_blocks] += 1
|
||||||
self._consecutive_blocks = 0
|
self._consecutive_blocks = 0
|
||||||
|
|
||||||
def should_abort(self) -> bool:
|
def abort_reason(self) -> str | None:
|
||||||
|
"""Какой критерий требует обрыва прямо сейчас, или None.
|
||||||
|
|
||||||
|
Возвращает "safety_net" / "ratio" / None. Состояние не меняет, поэтому
|
||||||
|
вызывать можно сколько угодно раз -- в том числе повторно, ради текста лога.
|
||||||
|
"""
|
||||||
# Safety-net -- ТОЛЬКО когда ratio-критерий физически недостижим (снапшот
|
# Safety-net -- ТОЛЬКО когда ratio-критерий физически недостижим (снапшот
|
||||||
# короче окна), иначе пачка safety_min блоков в начале длинного прогона
|
# короче окна), иначе пачка safety_min блоков в начале длинного прогона
|
||||||
# абортила бы его так же, как до правки (#3184 review MAJOR 2).
|
# абортила бы его так же, как до правки (#3184 review MAJOR 2).
|
||||||
|
|
@ -117,12 +132,40 @@ class BlockRatioBreaker:
|
||||||
and self._pure_block_run
|
and self._pure_block_run
|
||||||
and self._consecutive_blocks >= self.safety_min
|
and self._consecutive_blocks >= self.safety_min
|
||||||
):
|
):
|
||||||
return True
|
return "safety_net"
|
||||||
if len(self._window) == self.window_size and self.window_size > 0:
|
if len(self._window) == self.window_size and self.window_size > 0:
|
||||||
ratio = sum(self._window) / self.window_size
|
ratio = sum(self._window) / self.window_size
|
||||||
if ratio >= self.ratio_threshold:
|
if ratio >= self.ratio_threshold:
|
||||||
return True
|
return "ratio"
|
||||||
return False
|
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]:
|
def finalize(self) -> dict[str, int]:
|
||||||
"""Досчитать хвостовую пачку (прогон оборвался посреди серии блоков, без
|
"""Досчитать хвостовую пачку (прогон оборвался посреди серии блоков, без
|
||||||
|
|
|
||||||
|
|
@ -37,6 +37,11 @@ block_ratio_window/_threshold) ИЛИ safety-net для коротких про
|
||||||
consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое
|
consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое
|
||||||
поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в
|
поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в
|
||||||
counters["block_streak_histogram"], иначе эффект правки нечем измерить постфактум.
|
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
|
from __future__ import annotations
|
||||||
|
|
@ -431,6 +436,7 @@ async def run_avito_detail_backfill(
|
||||||
)
|
)
|
||||||
consecutive_failures = 0
|
consecutive_failures = 0
|
||||||
aborted_by_blocks = False
|
aborted_by_blocks = False
|
||||||
|
abort_reason: str | None = None
|
||||||
do_sleep = False
|
do_sleep = False
|
||||||
items_since_warm = 0
|
items_since_warm = 0
|
||||||
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
|
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
|
||||||
|
|
@ -632,12 +638,18 @@ async def run_avito_detail_backfill(
|
||||||
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate.
|
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate.
|
||||||
# #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни
|
# #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни
|
||||||
# одного успеха, накопилось max_consecutive_blocks) -- см. module docstring.
|
# одного успеха, накопилось 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(
|
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",
|
"частая причина: %s. enriched=%d attempted=%d",
|
||||||
run_id,
|
run_id,
|
||||||
breaker.consecutive_blocks,
|
breaker.abort_explanation(),
|
||||||
_top_failure(failure_census) or "причина не определена",
|
_top_failure(failure_census) or "причина не определена",
|
||||||
counters.enriched,
|
counters.enriched,
|
||||||
counters.attempted,
|
counters.attempted,
|
||||||
|
|
@ -739,6 +751,11 @@ async def run_avito_detail_backfill(
|
||||||
streak_histogram = breaker.finalize()
|
streak_histogram = breaker.finalize()
|
||||||
if streak_histogram:
|
if streak_histogram:
|
||||||
current_counters["block_streak_histogram"] = streak_histogram # type: ignore[assignment]
|
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(
|
runs_mod.mark_backfill_finished(
|
||||||
db,
|
db,
|
||||||
run_id,
|
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()
|
runs.mark_backfill_finished.assert_called_once()
|
||||||
assert runs.mark_backfill_finished.call_args.kwargs["aborted_by_blocks"] is False
|
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