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
129 lines
8.2 KiB
Python
129 lines
8.2 KiB
Python
"""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,
|
||
)
|