fix(tradein): такт в сохранении расписания, position_in_serp невыразим (#2674) #2694
7 changed files with 329 additions and 64 deletions
|
|
@ -1592,8 +1592,22 @@ def update_schedule(
|
|||
"""UPDATE existing schedule (create если не существует, через INSERT ON CONFLICT)."""
|
||||
from app.services.scheduler import compute_next_run_at
|
||||
|
||||
# Compute new next_run_at если window изменился — recompute, иначе keep existing
|
||||
next_at = compute_next_run_at(payload.window_start_hour, payload.window_end_hour)
|
||||
# #2674: такт берётся из default_params — ровно как его читает планировщик
|
||||
# (_claim_run/_defer_next_run_at). Без него compute_next_run_at падал на default=1 и
|
||||
# ЛЮБОЕ сохранение сбивало источник на «завтра»: недельный avito_full_load после
|
||||
# правки окна побежал бы через сутки. На суточных источниках баг был невидим —
|
||||
# для них «завтра» и есть правильный ответ.
|
||||
# None-safe так же, как в scheduler: `"interval_days": null` в jsonb → 1, не TypeError.
|
||||
_interval_days = payload.default_params.get("interval_days")
|
||||
interval_days = max(1, int(_interval_days)) if _interval_days is not None else 1
|
||||
|
||||
# Явно заданный оператором момент уважается как есть (в т.ч. в прошлом — «запустить
|
||||
# сейчас»). Иначе считаем от такта.
|
||||
next_at = payload.next_run_at or compute_next_run_at(
|
||||
payload.window_start_hour,
|
||||
payload.window_end_hour,
|
||||
interval_days=interval_days,
|
||||
)
|
||||
|
||||
row = (
|
||||
db.execute(
|
||||
|
|
@ -1607,7 +1621,22 @@ def update_schedule(
|
|||
window_start_hour = EXCLUDED.window_start_hour,
|
||||
window_end_hour = EXCLUDED.window_end_hour,
|
||||
default_params = EXCLUDED.default_params,
|
||||
next_run_at = EXCLUDED.next_run_at,
|
||||
-- #2674: не двигаем уже назначенный запуск, если двигать не за чем.
|
||||
-- Раньше next_run_at перезаписывался ВСЕГДА, поэтому правка соседнего
|
||||
-- поля (enabled, request_delay_sec в params) заново разыгрывала момент
|
||||
-- внутри окна и сдвигала прогон. Сохраняем существующий только когда он
|
||||
-- ещё в будущем И ни окно, ни такт не менялись — тогда пересчёт дал бы
|
||||
-- то же самое окно, только с другим random-смещением.
|
||||
next_run_at = CASE
|
||||
WHEN CAST(:explicit AS boolean) THEN EXCLUDED.next_run_at
|
||||
WHEN scrape_schedules.next_run_at > NOW()
|
||||
AND scrape_schedules.window_start_hour = EXCLUDED.window_start_hour
|
||||
AND scrape_schedules.window_end_hour = EXCLUDED.window_end_hour
|
||||
AND COALESCE(scrape_schedules.default_params ->> 'interval_days', '1')
|
||||
= COALESCE(EXCLUDED.default_params ->> 'interval_days', '1')
|
||||
THEN scrape_schedules.next_run_at
|
||||
ELSE EXCLUDED.next_run_at
|
||||
END,
|
||||
updated_at = NOW()
|
||||
RETURNING id, source, enabled, window_start_hour, window_end_hour,
|
||||
default_params, last_run_id, last_run_at, next_run_at, updated_at
|
||||
|
|
@ -1620,6 +1649,7 @@ def update_schedule(
|
|||
"we": payload.window_end_hour,
|
||||
"params": json.dumps(payload.default_params, ensure_ascii=False),
|
||||
"next_at": next_at,
|
||||
"explicit": payload.next_run_at is not None,
|
||||
},
|
||||
)
|
||||
.mappings()
|
||||
|
|
|
|||
|
|
@ -417,6 +417,11 @@ class ScheduleConfigUpdate(BaseModel):
|
|||
window_start_hour: int = Field(default=2, ge=0, le=23)
|
||||
window_end_hour: int = Field(default=5, ge=0, le=23)
|
||||
default_params: dict[str, Any] = Field(default_factory=dict)
|
||||
# #2674: явная воля оператора по времени следующего запуска. None (умолчание) —
|
||||
# «не трогай, посчитай сам от такта». Заданное значение уважается как есть, включая
|
||||
# прошедшее/now() — это и есть «запустить сейчас» (планировщик берёт строки с
|
||||
# next_run_at <= NOW()), у которого до сих пор не было API и его делали UPDATE'ом.
|
||||
next_run_at: datetime | None = None
|
||||
|
||||
|
||||
# ── House analytics (house_placement_history backfill) ───────────────────────
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ scheduling-путь (`app/scheduler_main.py` безусловно запуска
|
|||
Что осталось в этом модуле — НЕ scheduler-loop, а функции с живыми потребителями вне
|
||||
удалённой machinery:
|
||||
- `compute_next_run_at` — читается admin.py (операторский предпросмотр "next run").
|
||||
С #2674 это re-export kit-версии, а не вторая копия формулы.
|
||||
- `has_running_run` — читается admin.py (UI-индикатор "уже бежит").
|
||||
- `import_rosreestr_dkp` — job-тело, вызываемое kit-handler'ом
|
||||
product_handlers._job_rosreestr_dkp (lazy import).
|
||||
|
|
@ -25,16 +26,24 @@ Zombie-reap, advisory-lock claim и tick-loop теперь целиком в
|
|||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import random
|
||||
from datetime import UTC, datetime, time, timedelta
|
||||
from typing import Any
|
||||
|
||||
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
|
||||
# Обе копии одинаково умели interval_days — но такт доезжал до next_run_at только через
|
||||
# kit (_claim_run/_defer_next_run_at читают default_params["interval_days"]); admin.py
|
||||
# звал эту копию БЕЗ аргумента, получал default=1 и сбивал любой источник на «завтра».
|
||||
# Копия удалена, а не подправлена: пока формула лежит в двух файлах, следующая правка
|
||||
# такта снова разъедется по одному из них. Re-export (а не правка импорта у вызывающих)
|
||||
# сохраняет `from app.services.scheduler import compute_next_run_at` в admin.py и тестах.
|
||||
from scraper_kit.orchestration.scheduler import compute_next_run_at
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.core.shutdown import shutdown_requested
|
||||
from app.services import scrape_runs as runs_mod
|
||||
|
||||
__all__ = ["compute_next_run_at", "has_running_run"]
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# import_rosreestr_dkp: доля per-row INSERT-ошибок (rows_errored / rows_fetched), выше
|
||||
|
|
@ -43,54 +52,6 @@ logger = logging.getLogger(__name__)
|
|||
DKP_IMPORT_ERROR_RATE_THRESHOLD = 0.05
|
||||
|
||||
|
||||
def compute_next_run_at(
|
||||
window_start_hour: int,
|
||||
window_end_hour: int,
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
interval_days: int = 1,
|
||||
) -> datetime:
|
||||
"""Pick random datetime в window [start, end) UTC, через interval_days суток после now.
|
||||
|
||||
interval_days задаёт каденс источника: 1 (default) = daily (back-compat), 7 = weekly.
|
||||
Берётся из schedule.default_params["interval_days"] вызывающим кодом; отсутствие ключа
|
||||
→ 1 → прежнее ежедневное поведение.
|
||||
|
||||
Если window_end_hour <= window_start_hour → cross-midnight window
|
||||
(например 22→3 → окно 22:00-23:59 ИЛИ 00:00-02:59).
|
||||
"""
|
||||
now = now or datetime.now(tz=UTC)
|
||||
interval_days = max(1, int(interval_days))
|
||||
# Целевая дата = now + interval_days суток (interval_days=1 → завтра, как раньше).
|
||||
target = (now + timedelta(days=interval_days)).date()
|
||||
|
||||
if window_end_hour > window_start_hour:
|
||||
# Обычное окно (например 2..5 → 02:00-04:59)
|
||||
start_seconds = window_start_hour * 3600
|
||||
end_seconds = window_end_hour * 3600
|
||||
rand_seconds = random.randint(start_seconds, end_seconds - 1)
|
||||
return datetime.combine(target, time(0, 0), tzinfo=UTC) + timedelta(seconds=rand_seconds)
|
||||
else:
|
||||
# Cross-midnight (22..3 → 22:00-23:59 + 00:00-02:59)
|
||||
# Длина окна = (24-start) + end часов
|
||||
total_seconds = ((24 - window_start_hour) + window_end_hour) * 3600
|
||||
rand_seconds = random.randint(0, total_seconds - 1)
|
||||
# Если rand попадает в первую часть (start..24)
|
||||
first_half = (24 - window_start_hour) * 3600
|
||||
if rand_seconds < first_half:
|
||||
# interval_days=1: текущая дата (если окно ещё не наступило сегодня) или next day.
|
||||
# interval_days>1: всегда целевая дата (стаггер на N суток вперёд).
|
||||
today_ok = interval_days == 1 and now.hour < window_start_hour
|
||||
base_date = now.date() if today_ok else target
|
||||
return datetime.combine(base_date, time(0, 0), tzinfo=UTC) + timedelta(
|
||||
seconds=window_start_hour * 3600 + rand_seconds
|
||||
)
|
||||
else:
|
||||
# Во второй части (0..end), целевого дня
|
||||
offset = rand_seconds - first_half
|
||||
return datetime.combine(target, time(0, 0), tzinfo=UTC) + timedelta(seconds=offset)
|
||||
|
||||
|
||||
def has_running_run(db: Session, source: str) -> bool:
|
||||
"""Есть ли активный run для source (status='running')."""
|
||||
row = db.execute(
|
||||
|
|
|
|||
|
|
@ -0,0 +1,62 @@
|
|||
-- 217_position_in_serp_unexpressible.sql
|
||||
-- listings_snapshots.position_in_serp — диагноз «механизм невыразим», шаг 1 из 2.
|
||||
--
|
||||
-- Dependencies: 016_listings_snapshots.sql (создала колонку).
|
||||
-- Apply after: 216_dead_code_sweep.sql
|
||||
-- Идемпотентно: только COMMENT ON COLUMN (CREATE OR REPLACE-семантика).
|
||||
--
|
||||
-- ── ЧТО НАЙДЕНО ──────────────────────────────────────────────────────────────
|
||||
-- upsert_listing_snapshot принимал position_in_serp, но НИ ОДИН боевой вызывающий
|
||||
-- его не передавал: три call-site'а (scraper_kit/base.py:801 save_listings,
|
||||
-- providers/cian/detail.py:382, app/tasks/avito_detail_backfill.py:454) плюс
|
||||
-- ingest-скрипты — везде аргумент опущен. На проде колонка пуста 0 из 395 240 строк
|
||||
-- (данные с 2026-05-30 по 2026-08-06).
|
||||
--
|
||||
-- ── ПОЧЕМУ ЭТО НЕ «ПОДКЛЮЧИТЬ» ───────────────────────────────────────────────
|
||||
-- Соблазнительный вывод — «парсер выдачи знает индекс карточки, передайте его».
|
||||
-- Он неверен: позиция есть свойство пары (объявление, конкретный прогон выдачи с
|
||||
-- конкретными фильтрами), а выбранная структура этого отношения не выражает.
|
||||
-- PRIMARY KEY (listing_id, snapshot_date) — максимум ОДНА строка на объявление в
|
||||
-- сутки; run_id здесь обычный атрибут, да ещё и под COALESCE в ON CONFLICT.
|
||||
-- Следствия на живых данных:
|
||||
-- * 2026-08-05: 77 прогонов и 5392 total_seen дали 4595 строк снэпшотов; за
|
||||
-- 2026-08-03 в таблице 13 разных run_id на одну дату. Четыре SERP-источника
|
||||
-- (yandex/cian/avito/domclick city_sweep) в одни сутки пишут по одному ключу.
|
||||
-- * Внутри одного city_sweep обход идёт по десяткам гео-якорей радиусом 1500 м с
|
||||
-- перекрытием, и save_listings получает `anchor_lots` — 3 страницы ОДНОГО якоря,
|
||||
-- а не общий ранжированный список. Одно объявление приезжает с разным индексом
|
||||
-- от разных якорей того же прогона.
|
||||
-- * Старый ON CONFLICT писал COALESCE(EXCLUDED.position_in_serp, <старое>) — в
|
||||
-- строке оседал бы индекс последнего писателя дня. Это не «позиция в выдаче», а
|
||||
-- произвольный представитель суток; хуже NULL, потому что читался бы как факт.
|
||||
-- Чтобы позицию можно было хранить честно, нужна отдельная таблица с ключом
|
||||
-- (run_id, listing_id) и сохранёнными фильтрами прогона. Такой задачи сейчас нет —
|
||||
-- ни один потребитель позицию не читает (0 view/matview на проде ссылаются на
|
||||
-- колонку), поэтому колонка удаляется, а не переносится.
|
||||
--
|
||||
-- ── ПОЧЕМУ DROP НЕ ЗДЕСЬ ─────────────────────────────────────────────────────
|
||||
-- Шаг 1 (этот файл + правка кода в том же PR): upsert_listing_snapshot перестаёт
|
||||
-- упоминать колонку; колонка остаётся, комментарий объясняет почему.
|
||||
-- Шаг 2 (отдельный PR, после того как образ с шагом 1 живёт на проде):
|
||||
-- ALTER TABLE listings_snapshots DROP COLUMN IF EXISTS position_in_serp;
|
||||
-- Порядок не косметический. deploy-tradein.yml применяет data/sql/*.sql ДО
|
||||
-- перезапуска контейнеров («(3) Применяем SQL миграции — ДО app»), так что DROP в
|
||||
-- одном деплое с правкой кода оставил бы окно в несколько минут, где ещё живой
|
||||
-- СТАРЫЙ образ выполняет INSERT со списком колонок, включающим удалённую. Запись
|
||||
-- снэпшотов fault-tolerant (SAVEPOINT + warning), поэтому упало бы тихо — ровно тот
|
||||
-- класс потерь, который замечают через недели по дырке в истории цен.
|
||||
|
||||
BEGIN;
|
||||
|
||||
COMMENT ON COLUMN listings_snapshots.position_in_serp IS
|
||||
'МЁРТВАЯ, удаляется следующей миграцией (#2674). Пуста 0/395240 на 2026-08-06: '
|
||||
'писатель принимал аргумент, ни один вызывающий его не передавал. Не подключена, '
|
||||
'потому что механизм невыразим в этой таблице: позиция — свойство пары '
|
||||
'(объявление, конкретный прогон выдачи с конкретными фильтрами), а PK здесь '
|
||||
'(listing_id, snapshot_date) — одна строка на объявление в сутки, при том что в '
|
||||
'одни сутки по этому ключу пишут до 13 прогонов и 4 разных SERP-источника, а '
|
||||
'внутри одного прогона объявление приходит с разным индексом от перекрывающихся '
|
||||
'гео-якорей. Честное хранение требует таблицы с ключом (run_id, listing_id) и '
|
||||
'сохранёнными фильтрами прогона; читателей у позиции нет ни одного.';
|
||||
|
||||
COMMIT;
|
||||
|
|
@ -0,0 +1,171 @@
|
|||
"""Сохранение расписания в админке уважает такт источника (#2674).
|
||||
|
||||
Баг: PUT /admin/scrape/schedules/{source} звал compute_next_run_at БЕЗ interval_days,
|
||||
получал default=1 и ставил next_run_at на завтра — какой бы такт ни стоял в
|
||||
default_params. Недельный avito_full_load после правки соседнего поля побежал бы через
|
||||
сутки. На суточных источниках дефект невидим: для них «завтра» и есть верный ответ,
|
||||
поэтому баг прожил до разбора #2674.
|
||||
|
||||
Планировщик (scraper_kit.orchestration.scheduler._claim_run/_defer_next_run_at) такт
|
||||
читал правильно — расходились именно два входа в одну и ту же формулу.
|
||||
|
||||
Без сети, без БД: endpoint вызывается напрямую с mock-сессией, проверяются bind-params.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
from app.api.v1.admin import update_schedule
|
||||
from app.schemas.trade_in import ScheduleConfigUpdate
|
||||
|
||||
|
||||
def _mock_db() -> MagicMock:
|
||||
db = MagicMock()
|
||||
row = {
|
||||
"id": 1,
|
||||
"source": "avito_full_load",
|
||||
"enabled": True,
|
||||
"window_start_hour": 13,
|
||||
"window_end_hour": 15,
|
||||
"default_params": {},
|
||||
"last_run_id": None,
|
||||
"last_run_at": None,
|
||||
"next_run_at": None,
|
||||
"updated_at": None,
|
||||
}
|
||||
db.execute.return_value.mappings.return_value.fetchone.return_value = row
|
||||
return db
|
||||
|
||||
|
||||
def _saved_params(db: MagicMock) -> dict[str, Any]:
|
||||
"""Bind-params единственного execute в update_schedule."""
|
||||
return db.execute.call_args_list[0].args[1]
|
||||
|
||||
|
||||
def _saved_sql(db: MagicMock) -> str:
|
||||
return str(db.execute.call_args_list[0].args[0])
|
||||
|
||||
|
||||
def test_weekly_source_gets_next_run_in_a_week_not_tomorrow() -> None:
|
||||
"""Такт 7 → сохранение → next_run_at примерно через 7 суток. Падает на старом коде."""
|
||||
db = _mock_db()
|
||||
before = datetime.now(tz=UTC)
|
||||
|
||||
update_schedule(
|
||||
"avito_full_load",
|
||||
ScheduleConfigUpdate(
|
||||
enabled=True,
|
||||
window_start_hour=13,
|
||||
window_end_hour=15,
|
||||
default_params={"interval_days": 7, "concurrency": 1},
|
||||
),
|
||||
db,
|
||||
)
|
||||
|
||||
next_at = _saved_params(db)["next_at"]
|
||||
delta_days = (next_at - before).total_seconds() / 86400
|
||||
# Окно 13:00-15:00 внутри целевых суток → разброс ±1 сутки вокруг ровно 7.
|
||||
assert 6.0 < delta_days < 8.0, f"ожидали ~7 суток, получили {delta_days:.2f}"
|
||||
# Старое поведение (interval_days не передан → default=1) дало бы «завтра».
|
||||
assert delta_days > 2.0, "next_run_at уехал на завтра — такт снова потерян"
|
||||
|
||||
|
||||
def test_daily_source_still_runs_tomorrow() -> None:
|
||||
"""Такт 1 (и его отсутствие) — прежнее поведение, back-compat."""
|
||||
for params in ({}, {"interval_days": 1}):
|
||||
db = _mock_db()
|
||||
before = datetime.now(tz=UTC)
|
||||
update_schedule(
|
||||
"cian_city_sweep",
|
||||
ScheduleConfigUpdate(window_start_hour=2, window_end_hour=5, default_params=params),
|
||||
db,
|
||||
)
|
||||
delta_days = (_saved_params(db)["next_at"] - before).total_seconds() / 86400
|
||||
assert 0.0 < delta_days < 2.0, f"params={params}: ожидали «завтра», got {delta_days:.2f}"
|
||||
|
||||
|
||||
def test_null_interval_days_is_none_safe() -> None:
|
||||
"""`"interval_days": null` в jsonb → такт 1, а не TypeError (как в scheduler)."""
|
||||
db = _mock_db()
|
||||
update_schedule(
|
||||
"cian_city_sweep",
|
||||
ScheduleConfigUpdate(
|
||||
window_start_hour=2, window_end_hour=5, default_params={"interval_days": None}
|
||||
),
|
||||
db,
|
||||
)
|
||||
assert _saved_params(db)["next_at"] is not None
|
||||
|
||||
|
||||
def test_explicit_next_run_at_is_honored() -> None:
|
||||
"""Оператор явно задал момент — уважаем как есть, ничего не пересчитываем."""
|
||||
db = _mock_db()
|
||||
wanted = datetime(2026, 9, 1, 3, 30, tzinfo=UTC)
|
||||
|
||||
update_schedule(
|
||||
"avito_full_load",
|
||||
ScheduleConfigUpdate(
|
||||
window_start_hour=13,
|
||||
window_end_hour=15,
|
||||
default_params={"interval_days": 7},
|
||||
next_run_at=wanted,
|
||||
),
|
||||
db,
|
||||
)
|
||||
|
||||
params = _saved_params(db)
|
||||
assert params["next_at"] == wanted
|
||||
assert params["explicit"] is True
|
||||
|
||||
|
||||
def test_explicit_past_next_run_at_means_run_now() -> None:
|
||||
"""«Запустить сейчас» = момент в прошлом/now: планировщик берёт next_run_at <= NOW()."""
|
||||
db = _mock_db()
|
||||
now = datetime.now(tz=UTC) - timedelta(minutes=1)
|
||||
|
||||
update_schedule(
|
||||
"avito_full_load",
|
||||
ScheduleConfigUpdate(
|
||||
window_start_hour=13, window_end_hour=15, default_params={}, next_run_at=now
|
||||
),
|
||||
db,
|
||||
)
|
||||
|
||||
params = _saved_params(db)
|
||||
assert params["next_at"] == now, "прошедший момент не должен подменяться пересчётом"
|
||||
assert params["explicit"] is True
|
||||
|
||||
|
||||
def test_no_explicit_next_run_at_flags_recompute() -> None:
|
||||
"""Без явного момента флаг explicit=false → SQL решает, двигать ли существующий."""
|
||||
db = _mock_db()
|
||||
update_schedule(
|
||||
"cian_city_sweep",
|
||||
ScheduleConfigUpdate(window_start_hour=2, window_end_hour=5, default_params={}),
|
||||
db,
|
||||
)
|
||||
assert _saved_params(db)["explicit"] is False
|
||||
|
||||
|
||||
def test_upsert_preserves_future_run_when_window_and_interval_unchanged() -> None:
|
||||
"""Shape-check: ON CONFLICT не перезаписывает будущий next_run_at без причины.
|
||||
|
||||
Само ветвление живёт в SQL (CASE), проверить его исполнение офлайн нечем — здесь
|
||||
сторожим, что ветка не исчезла из запроса при следующей правке.
|
||||
"""
|
||||
db = _mock_db()
|
||||
update_schedule(
|
||||
"cian_city_sweep",
|
||||
ScheduleConfigUpdate(window_start_hour=2, window_end_hour=5, default_params={}),
|
||||
db,
|
||||
)
|
||||
sql = _saved_sql(db)
|
||||
assert "scrape_schedules.next_run_at > NOW()" in sql
|
||||
assert "THEN scrape_schedules.next_run_at" in sql
|
||||
assert "CAST(:explicit AS boolean)" in sql # psycopg v3: CAST, не :explicit::boolean
|
||||
|
|
@ -119,7 +119,7 @@ def test_upsert_snapshot_params_minimal():
|
|||
assert p["run_id"] is None
|
||||
assert p["snap_date"] is None
|
||||
assert p["ppm2"] is None
|
||||
assert p["pos"] is None
|
||||
assert "pos" not in p # #2674: position_in_serp больше не пишется
|
||||
assert p["status"] == "active"
|
||||
|
||||
|
||||
|
|
@ -134,7 +134,6 @@ def test_upsert_snapshot_params_full():
|
|||
price_per_m2=130_000,
|
||||
run_id=99,
|
||||
snapshot_date=snap_date,
|
||||
position_in_serp=3,
|
||||
status="active",
|
||||
)
|
||||
|
||||
|
|
@ -144,10 +143,30 @@ def test_upsert_snapshot_params_full():
|
|||
assert params["ppm2"] == 130_000
|
||||
assert params["run_id"] == 99
|
||||
assert params["snap_date"] == snap_date
|
||||
assert params["pos"] == 3
|
||||
assert params["status"] == "active"
|
||||
|
||||
|
||||
def test_upsert_snapshot_does_not_write_position_in_serp():
|
||||
"""#2674: колонка не упоминается ни в сигнатуре, ни в SQL.
|
||||
|
||||
Диагноз — «механизм невыразим»: PK (listing_id, snapshot_date) держит одну строку
|
||||
на объявление в СУТКИ, а позиция есть свойство конкретного прогона выдачи с
|
||||
конкретными фильтрами (в одни сутки по ключу пишут до 13 прогонов и 4 разных
|
||||
SERP-источника). Значение осело бы от последнего писателя дня и читалось бы как
|
||||
факт. Тест ловит попытку «подключить проводку» обратно.
|
||||
"""
|
||||
import inspect
|
||||
|
||||
from scraper_kit.snapshot_writer import upsert_listing_snapshot as fn
|
||||
|
||||
assert "position_in_serp" not in inspect.signature(fn).parameters
|
||||
|
||||
db = _mock_db_simple()
|
||||
upsert_listing_snapshot(db, listing_id=1, price_rub=1_000_000)
|
||||
sql = str(db.execute.call_args_list[0].args[0])
|
||||
assert "position_in_serp" not in sql
|
||||
|
||||
|
||||
def test_upsert_snapshot_no_commit():
|
||||
"""upsert_listing_snapshot НЕ вызывает db.commit() — commit делает caller."""
|
||||
db = _mock_db_simple()
|
||||
|
|
|
|||
|
|
@ -14,7 +14,29 @@
|
|||
набирает тысячи строк за прогон и пишет их одним set-based statement'ом
|
||||
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
|
||||
per-row хелпер.
|
||||
|
||||
position_in_serp (#2674): параметр УДАЛЁН, колонка осталась и остаётся NULL.
|
||||
Это не оборванная проводка — задуманное отношение НЕВЫРАЗИМО в этой таблице.
|
||||
Позиция — свойство пары (объявление, конкретный прогон выдачи с конкретными
|
||||
фильтрами), а PRIMARY KEY здесь (listing_id, snapshot_date): максимум одна строка
|
||||
на объявление в СУТКИ, run_id лишь атрибут под COALESCE. Что это ломает:
|
||||
- за 2026-08-05 в listings_snapshots попали 77 прогонов (13 разных run_id за
|
||||
2026-08-03), в т.ч. четыре SERP-источника сразу (yandex/cian/avito/domclick
|
||||
city_sweep) — все схлопываются в одну строку на объявление;
|
||||
- внутри ОДНОГО прогона city_sweep обходит десятки гео-якорей радиусом 1500 м с
|
||||
перекрытием, и save_listings получает `anchor_lots` (3 страницы одного якоря),
|
||||
а не единый ранжированный список: одно объявление приходит с разным индексом от
|
||||
разных якорей;
|
||||
- ON CONFLICT писал COALESCE(EXCLUDED, existing), т.е. в строке оседал индекс того
|
||||
писателя, кто пришёл последним. Число было бы не «позицией», а произвольным
|
||||
представителем дня — хуже, чем NULL, потому что читалось бы как факт.
|
||||
Чтобы такое хранить, нужна таблица с ключом (run_id, listing_id) и сохранёнными
|
||||
фильтрами прогона. Колонка `listings_snapshots.position_in_serp` дропается отдельной
|
||||
миграцией ПОСЛЕ того как эта правка (код перестал её упоминать) доедет до прода —
|
||||
миграции на деплое применяются ДО перезапуска контейнеров, и одновременный DROP убил
|
||||
бы INSERT ещё живого старого образа. См. 217_position_in_serp_unexpressible.sql.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
|
@ -34,7 +56,6 @@ def upsert_listing_snapshot(
|
|||
price_per_m2: int | None = None,
|
||||
run_id: int | None = None,
|
||||
snapshot_date: date | None = None,
|
||||
position_in_serp: int | None = None,
|
||||
status: str | None = "active",
|
||||
) -> None:
|
||||
"""Записать / обновить snapshot для listing за текущий (или указанный) день.
|
||||
|
|
@ -50,21 +71,21 @@ def upsert_listing_snapshot(
|
|||
run_id: FK scrape_runs.id — с каким run'ом связан snapshot (None если backfill
|
||||
запускается вне run'а).
|
||||
snapshot_date: дата наблюдения; если None — используется CURRENT_DATE (по БД).
|
||||
position_in_serp: позиция в SERP (1..60+), только при SERP-scrape.
|
||||
status: 'active' / 'closed' / None. Default 'active' при обычном scrape.
|
||||
|
||||
Не пишет `listings_snapshots.position_in_serp` — см. модульный docstring выше.
|
||||
"""
|
||||
db.execute(
|
||||
text("""
|
||||
INSERT INTO listings_snapshots
|
||||
(listing_id, snapshot_date, run_id,
|
||||
price_rub, price_per_m2, position_in_serp, status, observed_at)
|
||||
price_rub, price_per_m2, status, observed_at)
|
||||
VALUES (
|
||||
CAST(:lid AS bigint),
|
||||
COALESCE(CAST(:snap_date AS date), CURRENT_DATE),
|
||||
CAST(:run_id AS bigint),
|
||||
CAST(:price AS bigint),
|
||||
CAST(:ppm2 AS int),
|
||||
CAST(:pos AS int),
|
||||
:status,
|
||||
NOW()
|
||||
)
|
||||
|
|
@ -72,9 +93,6 @@ def upsert_listing_snapshot(
|
|||
run_id = COALESCE(EXCLUDED.run_id, listings_snapshots.run_id),
|
||||
price_rub = EXCLUDED.price_rub,
|
||||
price_per_m2 = COALESCE(EXCLUDED.price_per_m2, listings_snapshots.price_per_m2),
|
||||
position_in_serp = COALESCE(
|
||||
EXCLUDED.position_in_serp, listings_snapshots.position_in_serp
|
||||
),
|
||||
status = COALESCE(EXCLUDED.status, listings_snapshots.status),
|
||||
observed_at = EXCLUDED.observed_at
|
||||
"""),
|
||||
|
|
@ -84,7 +102,6 @@ def upsert_listing_snapshot(
|
|||
"run_id": run_id,
|
||||
"price": price_rub,
|
||||
"ppm2": price_per_m2,
|
||||
"pos": position_in_serp,
|
||||
"status": status,
|
||||
},
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue