fix(tradein/scraper): отметки времени прогона перестают замерзать в его же транзакции (#2702)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 3m9s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 3m9s
now() в PostgreSQL — синоним transaction_timestamp(): замерзает на СТАРТЕ транзакции. Финализаторы (mark_done/mark_failed/mark_banned) выполняются той же сессией, что и работа задачи, поэтому если её транзакция всё это время оставалась открытой (коммитить было нечего: читающий батч, ноль сохранений, сохранение чужой сессией), их UPDATE попадал ВНУТРЬ неё и получал время НАЧАЛА работы. Прод 2026-08-06 (487 прогонов с finished_at и counters.duration_sec): у 153 заявленная длительность больше окна finished_at − started_at в >1.5 раза, у 133 окно меньше секунды при работе дольше 10 с. 126 из этих 133 окон лежат в 9-64 мс — подпись механизма, а не разброс: столько проходит от коммита claim'а до первого запроса рабочей транзакции. Прогон 346 (cian_history_backfill): 18 230 с работы, окно 32 мс. Дефект не сплошной ровно потому, что зависел от коммита перед финалом: cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят поштучно (окно = работа), yandex_address_backfill (45 из 50) / newbuilding_enrich / cian_history_backfill — нет. Правка: clock_timestamp() вместо now() во всех финализаторах и в heartbeat, обе копии runs-модуля (kit-копия входит в scraper-allowlist деплоя, app-копия — нет, но исполняется она; расходиться им нельзя). Побочно чинится критерий зависших прогонов: reap_zombies сравнивает записанный heartbeat со «сейчас», и обе стороны обязаны быть настоящим временем. На проде у всех 6 прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat остался на отметке старта (max advance 0.0 с) при нормальной длительности до 5.06 ч. started_at уже писался своей закоммиченной транзакцией (create_run коммитит INSERT до возврата run_id) — требование #2702 п.1 выполнялось и раньше; тест это фиксирует. Историю не переписываем: настоящий finished_at нигде не сохранился, но counters.duration_sec измерялся монотонными часами процесса и верен. Миграция 223 проставляет это COMMENT'ами на колонках, чтобы аналитика брала длительность из счётчика, а не из разности отметок. Refs #2702
This commit is contained in:
parent
9d8114158b
commit
cfed5ace75
6 changed files with 396 additions and 26 deletions
|
|
@ -2,6 +2,31 @@
|
|||
|
||||
Таблица scrape_runs создана в 015_scrape_runs.sql.
|
||||
Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled.
|
||||
|
||||
ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL —
|
||||
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции и не двигается,
|
||||
сколько бы та ни жила. Финализаторы (mark_done/mark_failed/mark_banned) выполняются
|
||||
ТОЙ ЖЕ сессией, что и работа задачи, — и если рабочая транзакция всё это время
|
||||
оставалась открытой (задача ничего не коммитила: нечего было сохранять, батч читающий,
|
||||
сохранение шло чужой сессией), их UPDATE попадал ВНУТРЬ неё, и `finished_at` получал
|
||||
время НАЧАЛА работы, а не её конца.
|
||||
|
||||
Замер на проде 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик
|
||||
counters.duration_sec): у 153 заявленная длительность превышала собственное окно
|
||||
finished_at − started_at более чем в 1.5 раза, у 133 окно было меньше секунды при
|
||||
работе дольше 10 с. 126 из этих 133 окон лежат в диапазоне 9-64 мс — это не разброс,
|
||||
а подпись механизма: столько проходит от коммита claim'а до первого запроса рабочей
|
||||
транзакции. Крайний случай — прогон 346 (cian_history_backfill): 18230 с работы,
|
||||
окно 32 мс.
|
||||
|
||||
Дефект был не сплошной ровно потому, что зависел от того, коммитила ли задача перед
|
||||
финалом: cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят
|
||||
поштучно, у них окно совпадало с работой; yandex_address_backfill (45 из 50 прогонов),
|
||||
newbuilding_enrich, cian_history_backfill — нет.
|
||||
|
||||
Побочно это чинит и `heartbeat_at`: он писался тем же `now()` и по той же причине
|
||||
отставал от реальности на возраст открытой транзакции, а на нём стоит поиск зависших
|
||||
прогонов (reap_zombies, порог 6 ч).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -282,7 +307,10 @@ def _alert_on_run_id(
|
|||
|
||||
|
||||
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||
"""INSERT scrape_runs(source, status='running', params, started_at=NOW()).
|
||||
"""INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()).
|
||||
|
||||
started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей
|
||||
транзакции задачи его уже не достаёт (#2702).
|
||||
|
||||
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
||||
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
||||
|
|
@ -293,7 +321,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
|||
text(
|
||||
"""
|
||||
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
||||
VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW())
|
||||
VALUES (
|
||||
:source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp()
|
||||
)
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
|
|
@ -305,7 +335,7 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
|||
|
||||
|
||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
||||
"""UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
|
||||
|
||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||
|
|
@ -316,7 +346,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET heartbeat_at = NOW(),
|
||||
SET heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -364,7 +394,7 @@ def is_cancelled(db: Session, run_id: int) -> bool:
|
|||
|
||||
|
||||
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||
"""Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count.
|
||||
"""Финализация run: status='done', finished_at + counters + total_seen/new_count.
|
||||
|
||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
||||
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
||||
|
|
@ -374,7 +404,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'done',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -413,7 +444,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'failed',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -472,7 +504,8 @@ def mark_banned(
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'banned',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
ban_kind = :ban_kind,
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
|
|
@ -583,7 +616,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'cancelled', finished_at = NOW()
|
||||
SET status = 'cancelled', finished_at = clock_timestamp()
|
||||
WHERE id = :run_id AND status = 'running'
|
||||
RETURNING id
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -0,0 +1,71 @@
|
|||
-- 223_scrape_runs_time_columns_meaning.sql
|
||||
-- Purpose (#2702): зафиксировать в схеме, что отметки времени прогона до этой
|
||||
-- правки не охватывали его работу, и куда смотреть аналитике вместо разности.
|
||||
--
|
||||
-- Dependencies: 015_scrape_runs.sql (создала started_at/finished_at/heartbeat_at),
|
||||
-- 051_scrape_runs_extend.sql (finished_at/counters).
|
||||
-- Apply after: 221_backfill_house_suggestions_image_link.sql
|
||||
-- Идемпотентно: только COMMENT ON COLUMN (перезаписывает сам себя), данных не трогает.
|
||||
--
|
||||
-- ── ЧТО БЫЛО СЛОМАНО ─────────────────────────────────────────────────────────
|
||||
-- Финализаторы (mark_done / mark_failed / mark_banned) писали finished_at и
|
||||
-- heartbeat_at через now(). В PostgreSQL now() == transaction_timestamp(): он
|
||||
-- замерзает на СТАРТЕ транзакции. Финализатор выполняется той же сессией, что и
|
||||
-- работа задачи; если рабочая транзакция всё это время оставалась открытой (задаче
|
||||
-- нечего было коммитить — читающий батч, ноль сохранений, сохранение чужой сессией),
|
||||
-- UPDATE финализатора попадал ВНУТРЬ неё и получал время НАЧАЛА работы.
|
||||
--
|
||||
-- Прод-замер 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик
|
||||
-- counters.duration_sec):
|
||||
-- 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза;
|
||||
-- 133 — окно меньше секунды при работе дольше 10 с.
|
||||
-- 126 из этих 133 окон лежат в 9-64 мс: это не разброс, а подпись механизма —
|
||||
-- столько проходит от коммита claim'а до первого запроса рабочей транзакции.
|
||||
-- Крайние: прогон 346 (cian_history_backfill) — 18 230 с работы при окне 32 мс;
|
||||
-- 497 (newbuilding_enrich) — 6 124 с при 21 мс; 341 (yandex_address_backfill) —
|
||||
-- 1 460 с при 19 мс.
|
||||
--
|
||||
-- Дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом.
|
||||
-- Средние окно/duration_sec по источникам на том же замере:
|
||||
-- cian_history_backfill 2554 / 4222 ← окно короче работы
|
||||
-- newbuilding_enrich 1302 / 2078 ← короче
|
||||
-- yandex_address_backfill 107 / 1022 ← короче в 10 раз
|
||||
-- avito_detail_backfill 2829 / 2016 ← длиннее (норма)
|
||||
-- house_imv_backfill 1173 / 1173 ← совпадает
|
||||
-- cadastral_geo_match 11 / 10 ← совпадает
|
||||
--
|
||||
-- ── ПОЧЕМУ ИСТОРИЮ НЕ ЧИНИМ ──────────────────────────────────────────────────
|
||||
-- Восстановить настоящий finished_at по строке нельзя: реальное время конца нигде
|
||||
-- не сохранилось. Но counters.duration_sec измерялся монотонными часами процесса
|
||||
-- (time.monotonic / time.time в самих задачах) и транзакцией не затронут — он у
|
||||
-- этих строк верный. Поэтому история не переписывается, а помечается: аналитика
|
||||
-- обязана брать длительность из счётчика, а не из разности отметок.
|
||||
--
|
||||
-- Правка кода (clock_timestamp() вместо now() во всех финализаторах и в heartbeat)
|
||||
-- живёт в app/services/scrape_runs.py + packages/scraper-kit/.../orchestration/runs.py.
|
||||
|
||||
COMMENT ON COLUMN scrape_runs.started_at IS
|
||||
'Момент claim''а прогона. Пишется create_run своей транзакцией (commit сразу '
|
||||
'после INSERT), поэтому откат рабочей транзакции задачи его не затрагивает. '
|
||||
'С #2702 — clock_timestamp(); до него now() (= старт транзакции тика планировщика), '
|
||||
'что давало сдвиг в пределах тика.';
|
||||
|
||||
COMMENT ON COLUMN scrape_runs.finished_at IS
|
||||
'Момент финализации прогона. ВНИМАНИЕ: у строк ДО #2702 (2026-08-06) значение '
|
||||
'недостоверно — писалось now() (= transaction_timestamp) внутри рабочей транзакции '
|
||||
'задачи, поэтому у прогонов, ничего не коммитивших по ходу работы, равно времени '
|
||||
'её НАЧАЛА. На проде так вышло у 133 из 487 прогонов со счётчиком длительности '
|
||||
'(окно < 1 с при работе > 10 с). Длительность таких прогонов брать из '
|
||||
'counters->>''duration_sec'' (монотонные часы процесса, транзакцией не затронуты), '
|
||||
'а НЕ из finished_at − started_at.';
|
||||
|
||||
COMMENT ON COLUMN scrape_runs.heartbeat_at IS
|
||||
'Последний признак жизни прогона; на нём стоит поиск зависших (reap_zombies, порог '
|
||||
'6 ч). У строк ДО #2702 отставал от реальности на возраст открытой рабочей '
|
||||
'транзакции по той же причине, что finished_at, — то есть критерий «завис» решал '
|
||||
'по замороженной отметке. С #2702 пишется clock_timestamp().';
|
||||
|
||||
COMMENT ON COLUMN scrape_runs.counters IS
|
||||
'Счётчики прогона (jsonb). Ключ duration_sec, где он есть, измерен монотонными '
|
||||
'часами процесса и остаётся единственным достоверным источником длительности для '
|
||||
'строк до #2702 (см. комментарий к finished_at).';
|
||||
237
tradein-mvp/backend/tests/test_2702_run_timestamps.py
Normal file
237
tradein-mvp/backend/tests/test_2702_run_timestamps.py
Normal file
|
|
@ -0,0 +1,237 @@
|
|||
"""#2702: отметки времени прогона не охватывают его работу.
|
||||
|
||||
Что было. Финализаторы писали `finished_at`/`heartbeat_at` через `now()`, а `now()`
|
||||
в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции.
|
||||
Выполняются финализаторы той же сессией, что и работа задачи, поэтому если рабочая
|
||||
транзакция всё это время оставалась открытой, их UPDATE попадал ВНУТРЬ неё и получал
|
||||
время НАЧАЛА работы.
|
||||
|
||||
Прод-замер 2026-08-06 (487 прогонов с finished_at и counters.duration_sec):
|
||||
* 153 — заявленная длительность больше окна finished_at − started_at в >1.5 раза;
|
||||
* 133 — окно меньше секунды при работе дольше 10 с, причём 126 из них укладываются
|
||||
в 9-64 мс: столько проходит от коммита claim'а до первого запроса рабочей
|
||||
транзакции — подпись механизма, а не разброс.
|
||||
Крайний случай, воспроизведённый ниже дословно: прогон 346 (cian_history_backfill) —
|
||||
18 230 с работы, окно 32 мс.
|
||||
|
||||
Почему дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом.
|
||||
`cadastral_geo_match` / `house_imv_backfill` / `avito_detail_backfill` коммитят
|
||||
поштучно — у них окно совпадало с работой; `yandex_address_backfill` (45 из 50
|
||||
прогонов), `newbuilding_enrich`, `cian_history_backfill` — нет.
|
||||
|
||||
Фальсификация: на старом коде (`now()`) тесты 1 и 3 дают другой ответ.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import inspect
|
||||
import os
|
||||
import re
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
from scraper_kit.orchestration import runs as kit_runs
|
||||
from scraper_kit.orchestration import scheduler as kit_scheduler
|
||||
|
||||
from app.services import scrape_runs as app_runs
|
||||
|
||||
_MODULES = {"kit": kit_runs, "app": app_runs}
|
||||
|
||||
# Колонки, которые обязаны нести НАСТОЯЩЕЕ время, а не время старта транзакции.
|
||||
_TS_COLS = ("started_at", "finished_at", "heartbeat_at")
|
||||
|
||||
|
||||
def _split_top(items: str) -> list[str]:
|
||||
"""Разбить список SQL-элементов по запятым ВЕРХНЕГО уровня (CAST(:x AS t) — один)."""
|
||||
out: list[str] = []
|
||||
depth = 0
|
||||
cur = ""
|
||||
for ch in items:
|
||||
if ch == "," and depth == 0:
|
||||
out.append(cur.strip())
|
||||
cur = ""
|
||||
continue
|
||||
depth += (ch == "(") - (ch == ")")
|
||||
cur += ch
|
||||
out.append(cur.strip())
|
||||
return out
|
||||
|
||||
|
||||
class _Row:
|
||||
"""Строка ответа: id для create_run, source для alert-хука."""
|
||||
|
||||
id = 1
|
||||
source = "src"
|
||||
status = "done"
|
||||
|
||||
|
||||
class _FakeResult:
|
||||
def fetchone(self) -> _Row:
|
||||
return _Row()
|
||||
|
||||
def first(self) -> _Row:
|
||||
return _Row()
|
||||
|
||||
def fetchall(self) -> list[_Row]:
|
||||
return []
|
||||
|
||||
|
||||
class _FakePg:
|
||||
"""Мини-модель PostgreSQL на две функции времени и ленивую транзакцию.
|
||||
|
||||
`now()` == `transaction_timestamp()` — замерзает на старте транзакции;
|
||||
`clock_timestamp()` — настоящие часы. Транзакция открывается лениво на первом
|
||||
execute (autobegin SQLAlchemy) и закрывается commit/rollback. Больше модель
|
||||
ничего не умеет — этого достаточно, чтобы отличить одно от другого.
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.wall: float = 0.0 # «стенные часы» теста
|
||||
self.tx_start: float | None = None
|
||||
self.row: dict[str, float] = {} # что осело в scrape_runs
|
||||
self.calls: list[str] = []
|
||||
|
||||
def _stamp(self, col: str, func: str) -> None:
|
||||
assert self.tx_start is not None
|
||||
self.row[col] = self.tx_start if func.lower() == "now" else self.wall
|
||||
|
||||
def execute(self, stmt: Any, params: Any = None) -> _FakeResult:
|
||||
if self.tx_start is None:
|
||||
self.tx_start = self.wall
|
||||
self.calls.append("execute")
|
||||
sql = str(stmt)
|
||||
for col in _TS_COLS: # UPDATE ... SET <col> = <func>()
|
||||
m = re.search(rf"\b{col}\s*=\s*(now|clock_timestamp)\s*\(\s*\)", sql, re.I)
|
||||
if m is not None:
|
||||
self._stamp(col, m.group(1))
|
||||
ins = re.search(
|
||||
r"INSERT INTO scrape_runs\s*\((.*?)\).*?VALUES\s*\((.*)\)", sql, re.S | re.I
|
||||
)
|
||||
if ins is not None: # INSERT — отметки времени позиционные, в VALUES
|
||||
for col, val in zip(_split_top(ins.group(1)), _split_top(ins.group(2)), strict=False):
|
||||
f = re.fullmatch(r"(now|clock_timestamp)\s*\(\s*\)", val, re.I)
|
||||
if col in _TS_COLS and f is not None:
|
||||
self._stamp(col, f.group(1))
|
||||
return _FakeResult()
|
||||
|
||||
def commit(self) -> None:
|
||||
self.calls.append("commit")
|
||||
self.tx_start = None
|
||||
|
||||
def rollback(self) -> None:
|
||||
self.calls.append("rollback")
|
||||
self.tx_start = None
|
||||
|
||||
|
||||
# ── 1. Финал прогона датируется концом работы, а не стартом транзакции ────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", list(_MODULES))
|
||||
def test_finished_at_covers_the_work_not_the_transaction_start(name: str) -> None:
|
||||
"""Прогон 346 дословно: 18 230 с работы в одной незакоммиченной транзакции.
|
||||
|
||||
На старом коде finished_at = 0.032 (старт рабочей транзакции) → окно 32 мс при
|
||||
пяти часах работы. Это и есть та строка, ради которой заведена задача.
|
||||
"""
|
||||
mod = _MODULES[name]
|
||||
db = _FakePg()
|
||||
db.row["started_at"] = 0.0 # create_run уже закоммитил claim
|
||||
|
||||
db.wall = 0.020
|
||||
mod.update_heartbeat(db, 1, {"listings_processed": 0}) # heartbeat + commit
|
||||
|
||||
db.wall = 0.032
|
||||
db.execute("SELECT id FROM listings WHERE history IS NULL") # рабочая транзакция
|
||||
|
||||
db.wall = 18230.0 # пять часов работы, ни одного коммита
|
||||
mod.mark_done(db, 1, {"listings_processed": 1200, "duration_sec": 18230})
|
||||
|
||||
assert db.row["finished_at"] == pytest.approx(18230.0)
|
||||
assert db.row["heartbeat_at"] == pytest.approx(18230.0)
|
||||
window = db.row["finished_at"] - db.row["started_at"]
|
||||
assert window == pytest.approx(18230.0), "окно прогона обязано охватывать его работу"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", list(_MODULES))
|
||||
@pytest.mark.parametrize("finalizer", ["mark_failed", "mark_banned"])
|
||||
def test_failed_and_banned_finals_are_wall_clock_too(name: str, finalizer: str) -> None:
|
||||
"""Тот же инвариант для неуспешных финалов.
|
||||
|
||||
Их спасал defensive-rollback в начале (он закрывал рабочую транзакцию), но
|
||||
полагаться на побочный эффект чужой защиты нельзя — проверяем явно.
|
||||
"""
|
||||
mod = _MODULES[name]
|
||||
db = _FakePg()
|
||||
db.wall = 0.019
|
||||
db.execute("SELECT 1") # рабочая транзакция открыта
|
||||
db.wall = 1460.0
|
||||
getattr(mod, finalizer)(db, 1, "boom", {"checked": 0, "duration_sec": 1460})
|
||||
assert db.row["finished_at"] == pytest.approx(1460.0)
|
||||
|
||||
|
||||
# ── 2. started_at переживает откат рабочей транзакции ────────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", list(_MODULES))
|
||||
def test_started_at_survives_rolled_back_work_transaction(name: str) -> None:
|
||||
"""Требование #2702 п.1: отметка старта живёт в СВОЕЙ закоммиченной транзакции.
|
||||
|
||||
create_run коммитит INSERT до возврата run_id, а ни один финализатор прогона
|
||||
started_at не переписывает — поэтому откат рабочей транзакции его не достаёт.
|
||||
Финал при этом обязан быть датирован концом работы (это и падает на старом коде).
|
||||
"""
|
||||
mod = _MODULES[name]
|
||||
db = _FakePg()
|
||||
|
||||
run_id = mod.create_run(db, source="cian_history_backfill", params={})
|
||||
assert run_id == 1
|
||||
assert db.calls == ["execute", "commit"], "INSERT прогона обязан коммититься сразу"
|
||||
|
||||
db.wall = 0.030
|
||||
db.execute("UPDATE listings SET address = 'x'") # рабочая транзакция
|
||||
db.wall = 100.0
|
||||
db.rollback() # работа упала и откатилась
|
||||
|
||||
db.wall = 100.5
|
||||
db.execute("SELECT count(*) FROM listings") # новая рабочая транзакция
|
||||
db.wall = 1460.0
|
||||
mod.mark_done(db, 1, {"checked": 0, "duration_sec": 1460})
|
||||
|
||||
assert db.row["started_at"] == pytest.approx(0.0), "started_at не должен сдвигаться"
|
||||
assert db.row["finished_at"] == pytest.approx(1460.0)
|
||||
|
||||
|
||||
# ── 3. Инвариант источника: никакая отметка времени не пишется now() ─────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", list(_MODULES))
|
||||
def test_no_run_timestamp_is_written_with_now(name: str) -> None:
|
||||
"""`now()` в этих модулях не имеет корректного применения — его быть не должно.
|
||||
|
||||
Проверяем весь исходник, а не отдельные запросы: INSERT в create_run пишет
|
||||
started_at/heartbeat_at позиционно (в VALUES), и построчная проверка его бы
|
||||
пропустила — ровно так дефект и дожил до 3 300 прогонов.
|
||||
"""
|
||||
src = inspect.getsource(_MODULES[name])
|
||||
# Регистрозависимо: SQL в этих модулях пишется в верхнем регистре, а строчное
|
||||
# `now()` встречается в объяснительной прозе docstring'ов — ловим SQL, не текст.
|
||||
assert re.search(r"\bNOW\s*\(\s*\)", src) is None
|
||||
assert "clock_timestamp()" in src
|
||||
|
||||
|
||||
def test_zombie_criterion_compares_real_clocks() -> None:
|
||||
"""#2702 п.2: поиск зависших сравнивает записанный heartbeat со «сейчас».
|
||||
|
||||
Обе стороны сравнения обязаны быть настоящим временем: на проде у всех 6
|
||||
прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и
|
||||
остался на отметке старта (max advance 0.0 с) — критерий решал по замороженной
|
||||
отметке, хотя нормальный прогон этого источника длится до 5.06 ч.
|
||||
"""
|
||||
src = inspect.getsource(kit_scheduler.reap_zombies)
|
||||
stmt = re.search(r"UPDATE scrape_runs.*?RETURNING id", src, re.S)
|
||||
assert stmt is not None
|
||||
assert re.search(r"\bNOW\s*\(\s*\)", stmt.group(0)) is None
|
||||
assert stmt.group(0).count("clock_timestamp()") == 2 # finished_at + порог сравнения
|
||||
|
|
@ -129,7 +129,9 @@ def test_mark_skipped_collapse_refreshes_detail_and_started_at() -> None:
|
|||
update = db.sql_of("UPDATE scrape_runs r")
|
||||
assert update is not None
|
||||
sql, params = update
|
||||
assert "started_at = NOW()" in sql
|
||||
# clock_timestamp(), а не NOW(): отметки времени прогона пишутся настоящими
|
||||
# часами, иначе внутри долгой открытой транзакции они замерзают (#2702).
|
||||
assert "started_at = clock_timestamp()" in sql
|
||||
assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска"
|
||||
assert "first_skip_at" in sql, "начало стрика потеряно"
|
||||
assert params["details"] == "37 дн. назад"
|
||||
|
|
|
|||
|
|
@ -9,6 +9,15 @@ app-копии:
|
|||
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
|
||||
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
|
||||
app-копии эта функция не нужна.
|
||||
|
||||
ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL —
|
||||
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы
|
||||
выполняются той же сессией, что и работа задачи, и если её транзакция оставалась
|
||||
открытой всё время работы (нечего было коммитить), их UPDATE попадал ВНУТРЬ неё —
|
||||
`finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с
|
||||
окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в
|
||||
диапазоне 9-64 мс (= от коммита claim'а до первого запроса рабочей транзакции).
|
||||
Полный разбор — в docstring app-копии `app/services/scrape_runs.py`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -297,7 +306,10 @@ def _alert_on_run_id(
|
|||
|
||||
|
||||
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||
"""INSERT scrape_runs(source, status='running', params, started_at=NOW()).
|
||||
"""INSERT scrape_runs(source, status='running', params, started_at=clock_timestamp()).
|
||||
|
||||
started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей
|
||||
транзакции задачи его уже не достаёт (#2702).
|
||||
|
||||
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
|
||||
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
|
||||
|
|
@ -308,7 +320,9 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
|||
text(
|
||||
"""
|
||||
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
|
||||
VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW())
|
||||
VALUES (
|
||||
:source, 'running', CAST(:params AS jsonb), clock_timestamp(), clock_timestamp()
|
||||
)
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
|
|
@ -357,9 +371,9 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
LIMIT 1
|
||||
)
|
||||
UPDATE scrape_runs r
|
||||
SET heartbeat_at = NOW(),
|
||||
started_at = NOW(),
|
||||
finished_at = NOW(),
|
||||
SET heartbeat_at = clock_timestamp(),
|
||||
started_at = clock_timestamp(),
|
||||
finished_at = clock_timestamp(),
|
||||
counters = COALESCE(r.counters, '{}'::jsonb) || jsonb_build_object(
|
||||
'skips', COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1,
|
||||
'detail', CAST(:details AS text),
|
||||
|
|
@ -385,7 +399,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
source, status, error, counters, started_at, heartbeat_at, finished_at
|
||||
)
|
||||
VALUES (
|
||||
:source, 'skipped', :reason, CAST(:counters AS jsonb), NOW(), NOW(), NOW()
|
||||
:source, 'skipped', :reason, CAST(:counters AS jsonb), clock_timestamp(), clock_timestamp(), clock_timestamp()
|
||||
)
|
||||
RETURNING id
|
||||
"""
|
||||
|
|
@ -409,7 +423,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
|
||||
|
||||
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
||||
"""UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки.
|
||||
|
||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
|
||||
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
|
||||
|
|
@ -420,7 +434,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET heartbeat_at = NOW(),
|
||||
SET heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -447,7 +461,7 @@ def is_cancelled(db: Session, run_id: int) -> bool:
|
|||
|
||||
|
||||
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||
"""Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count.
|
||||
"""Финализация run: status='done', finished_at + counters + total_seen/new_count.
|
||||
|
||||
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
|
||||
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
|
||||
|
|
@ -457,7 +471,8 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'done',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -496,7 +511,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int])
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'failed',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
new_count = COALESCE(CAST(:new_count AS int), new_count)
|
||||
|
|
@ -555,7 +571,8 @@ def mark_banned(
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(),
|
||||
SET status = 'banned',
|
||||
finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(),
|
||||
error = :error, counters = CAST(:counters AS jsonb),
|
||||
ban_kind = :ban_kind,
|
||||
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
|
||||
|
|
@ -586,7 +603,7 @@ def mark_cancelled(db: Session, run_id: int) -> bool:
|
|||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'cancelled', finished_at = NOW()
|
||||
SET status = 'cancelled', finished_at = clock_timestamp()
|
||||
WHERE id = :run_id AND status = 'running'
|
||||
RETURNING id
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -235,16 +235,26 @@ def has_running_run(db: Session, source: str) -> bool:
|
|||
|
||||
|
||||
def reap_zombies(db: Session) -> int:
|
||||
"""Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours."""
|
||||
"""Mark scrape_runs as 'zombie' если heartbeat не обновлялся > ZOMBIE_THRESHOLD_HOURS hours.
|
||||
|
||||
clock_timestamp(), а не now() (#2702): критерий сравнивает ЗАПИСАННЫЙ heartbeat со
|
||||
временем «сейчас», и обе стороны сравнения должны быть настоящим временем. `now()`
|
||||
замерзает на старте транзакции — у писавшей heartbeat стороны это давало отставание
|
||||
на весь возраст открытой рабочей транзакции (см. docstring runs.py), у читающей
|
||||
стороны — на возраст тика. На проде это уже стоило ложных срабатываний: у всех 6
|
||||
прогонов cian_history_backfill, помеченных 'zombie', записанный heartbeat так и
|
||||
остался на отметке старта (max advance 0.0 с) — при том что нормальный прогон этого
|
||||
источника длится до 5.06 ч (прогон 346) и обязан был двигать heartbeat.
|
||||
"""
|
||||
zombie_interval = f"{ZOMBIE_THRESHOLD_HOURS} hours"
|
||||
result = db.execute(
|
||||
text(
|
||||
"""
|
||||
UPDATE scrape_runs
|
||||
SET status = 'zombie', finished_at = NOW()
|
||||
SET status = 'zombie', finished_at = clock_timestamp()
|
||||
WHERE status = 'running'
|
||||
AND (heartbeat_at IS NULL
|
||||
OR heartbeat_at < NOW() - CAST(:interval AS interval))
|
||||
OR heartbeat_at < clock_timestamp() - CAST(:interval AS interval))
|
||||
RETURNING id
|
||||
"""
|
||||
),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue