gendesign/tradein-mvp/packages/scraper-kit/src/scraper_kit/snapshot_writer.py
bot-backend 91423e0b53
Some checks failed
Deploy Trade-In / changes (push) Successful in 24s
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 3m41s
Deploy Trade-In / build-backend (push) Successful in 1m56s
Deploy Trade-In / deploy (push) Has been cancelled
fix(tradein/scrapers): метка наблюдения — время строки, а не старта транзакции (#2731) (#2742)
2026-08-06 16:19:56 +00:00

129 lines
8.2 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""listings_snapshots writer — point-in-time observation per listing per day.
Таблица listings_snapshots (016_listings_snapshots.sql):
PRIMARY KEY (listing_id, snapshot_date) — максимум 1 snapshot в сутки.
ON CONFLICT DO UPDATE — берём последний за день (перезаписываем при повторном run'е).
Используется тремя путями:
1. SERP scrape (save_listings в base.py) — записывает цену + позицию в выдаче
каждый раз когда listing появляется в поиске (status='active').
2. Detail backfill (save_detail_enrichment в cian_detail.py) — записывает цену
из detail-страницы для listings которые раньше не имели snapshot'а (status='active').
3. Снятие объявления (#2674): avito_detail_backfill при 404 с площадки пишет
status='closed'. Второй писатель 'closed' — deactivate_stale_listings, он
набирает тысячи строк за прогон и пишет их одним set-based statement'ом
(data-modifying CTE в app/tasks/deactivate_stale_avito.py), не через этот
per-row хелпер.
observed_at пишется `statement_timestamp()`, а НЕ `NOW()` (#2731). `NOW()` в PostgreSQL —
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции, а save_listings
коммитит ОДИН раз в конце всего batch'а — то есть все снимки прогона получали ОДНУ метку.
Прод-замер 2026-08-06: run 3303 — 219 строк, 1 различная метка, при семи минутах работы;
run 3229 — 2643 строки на 55 меток (метка на save_listings-вызов, не на строку).
Цена — не фильтры свежести (14 суток против минут прогона), а разрешение во времени:
темп сбора и порядок внутри прогона по данным не восстановить (отсюда же #2701).
Почему statement_timestamp(), а не clock_timestamp() как в #2702/#2718: этот INSERT — один
statement на строку (вызов в цикле save_listings), поэтому statement_timestamp() двигается
ПОСТРОЧНО и при этом СТАБИЛЕН внутри строки. clock_timestamp() дал бы разные значения даже
двум колонкам одного statement'а (прод-проверка: `clock_timestamp() = clock_timestamp()` → f,
`statement_timestamp() = statement_timestamp()` → t) — это сломало бы равенство
listings.scraped_at = last_seen_at, на которое опирается предикат миграции 161.
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
from datetime import date
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
def upsert_listing_snapshot(
db: Session,
*,
listing_id: int,
price_rub: int,
price_per_m2: int | None = None,
run_id: int | None = None,
snapshot_date: date | None = None,
status: str | None = "active",
) -> None:
"""Записать / обновить snapshot для listing за текущий (или указанный) день.
Идемпотентно: повторный вызов с теми же listing_id + snapshot_date перезаписывает
строку (берём последние данные за день — ON CONFLICT DO UPDATE).
Args:
db: открытая SQLAlchemy session. Commit делает CALLER.
listing_id: PK из таблицы listings.
price_rub: текущая цена в рублях.
price_per_m2: цена за кв.м (вычисляется в caller'е, None если нет).
run_id: FK scrape_runs.id — с каким run'ом связан snapshot (None если backfill
запускается вне run'а).
snapshot_date: дата наблюдения; если None — используется CURRENT_DATE (по БД).
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, 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),
:status,
statement_timestamp()
)
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
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),
status = COALESCE(EXCLUDED.status, listings_snapshots.status),
observed_at = EXCLUDED.observed_at
"""),
{
"lid": listing_id,
"snap_date": snapshot_date,
"run_id": run_id,
"price": price_rub,
"ppm2": price_per_m2,
"status": status,
},
)
logger.debug(
"upsert_listing_snapshot: listing_id=%s date=%s price=%s run_id=%s",
listing_id,
snapshot_date or "CURRENT_DATE",
price_rub,
run_id,
)