fix(tradein): писатели наконец пишут то, что обещает схема — фото подсказок, статус «снято», события объявлений (#2674) #2682
9 changed files with 766 additions and 80 deletions
|
|
@ -345,26 +345,37 @@ def save_imv_result(db: Session, house_id: int, params: dict, result: IMVEvaluat
|
|||
)
|
||||
|
||||
# 3. Suggestions
|
||||
# #2674: до этого фикса в INSERT не входили image_link + area_m2/rooms/floor/
|
||||
# total_floors — колонки есть с миграции 064, но писатель их не заполнял
|
||||
# (25 055 строк на проде с NULL во всех пяти). Ссылка на фото приходит в
|
||||
# suggestions.items[].imageLink, метрики квартиры парсятся из title.
|
||||
for sug in result.suggestions:
|
||||
db.execute(
|
||||
text("""
|
||||
INSERT INTO house_suggestions (
|
||||
house_id, ext_item_id, title, address, price_rub,
|
||||
area_m2, rooms, floor, total_floors,
|
||||
exposure_days, publish_date,
|
||||
item_link, metro_name, metro_distance,
|
||||
item_link, image_link, metro_name, metro_distance,
|
||||
has_good_price_badge, raw_payload, fetched_at
|
||||
) VALUES (
|
||||
:hid, :ext, :title, :addr, :price,
|
||||
CAST(:area AS numeric), :rooms, :floor, :total_floors,
|
||||
:exp, :pdate,
|
||||
:link, :mname, :mdist,
|
||||
:link, :img, :mname, :mdist,
|
||||
:gpb, CAST(:raw AS jsonb), NOW()
|
||||
)
|
||||
ON CONFLICT (house_id, ext_item_id) DO UPDATE SET
|
||||
title = EXCLUDED.title,
|
||||
price_rub = EXCLUDED.price_rub,
|
||||
area_m2 = EXCLUDED.area_m2,
|
||||
rooms = EXCLUDED.rooms,
|
||||
floor = EXCLUDED.floor,
|
||||
total_floors = EXCLUDED.total_floors,
|
||||
exposure_days = EXCLUDED.exposure_days,
|
||||
publish_date = EXCLUDED.publish_date,
|
||||
item_link = EXCLUDED.item_link,
|
||||
image_link = EXCLUDED.image_link,
|
||||
metro_name = EXCLUDED.metro_name,
|
||||
metro_distance = EXCLUDED.metro_distance,
|
||||
has_good_price_badge = EXCLUDED.has_good_price_badge,
|
||||
|
|
@ -377,9 +388,14 @@ def save_imv_result(db: Session, house_id: int, params: dict, result: IMVEvaluat
|
|||
"title": sug.title,
|
||||
"addr": sug.address,
|
||||
"price": sug.price_rub,
|
||||
"area": sug.area_m2,
|
||||
"rooms": sug.rooms,
|
||||
"floor": sug.floor,
|
||||
"total_floors": sug.total_floors,
|
||||
"exp": sug.exposure_days,
|
||||
"pdate": sug.publish_date,
|
||||
"link": sug.item_url,
|
||||
"img": sug.image_link,
|
||||
"mname": sug.metro_name,
|
||||
"mdist": sug.metro_distance,
|
||||
"gpb": sug.has_good_price_badge,
|
||||
|
|
|
|||
|
|
@ -47,6 +47,7 @@ from scraper_kit.providers.avito.detail import (
|
|||
save_detail_enrichment,
|
||||
)
|
||||
from scraper_kit.providers.avito.serp import AvitoScraper
|
||||
from scraper_kit.snapshot_writer import upsert_listing_snapshot
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
|
|
@ -241,7 +242,7 @@ async def run_avito_detail_backfill(
|
|||
text(
|
||||
"""
|
||||
WITH ekb AS (
|
||||
SELECT id, source_url, 'ekb' AS city_scope
|
||||
SELECT id, source_url, price_rub, 'ekb' AS city_scope
|
||||
FROM listings
|
||||
WHERE source = 'avito'
|
||||
AND detail_enriched_at IS NULL
|
||||
|
|
@ -254,7 +255,7 @@ async def run_avito_detail_backfill(
|
|||
LIMIT CAST(:batch_size AS int)
|
||||
),
|
||||
oblast AS (
|
||||
SELECT id, source_url, 'oblast' AS city_scope
|
||||
SELECT id, source_url, price_rub, 'oblast' AS city_scope
|
||||
FROM listings
|
||||
WHERE source = 'avito'
|
||||
AND detail_enriched_at IS NULL
|
||||
|
|
@ -264,9 +265,9 @@ async def run_avito_detail_backfill(
|
|||
ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST
|
||||
LIMIT CAST(:oblast_batch_size AS int)
|
||||
)
|
||||
SELECT id, source_url, city_scope FROM ekb
|
||||
SELECT id, source_url, price_rub, city_scope FROM ekb
|
||||
UNION ALL
|
||||
SELECT id, source_url, city_scope FROM oblast
|
||||
SELECT id, source_url, price_rub, city_scope FROM oblast
|
||||
"""
|
||||
),
|
||||
{
|
||||
|
|
@ -442,6 +443,21 @@ async def run_avito_detail_backfill(
|
|||
text("UPDATE listings SET is_active = FALSE WHERE id = :id"),
|
||||
{"id": row["id"]},
|
||||
)
|
||||
# #2674: 404 с площадки — самый достоверный сигнал снятия,
|
||||
# фиксируем его в дневной истории (listings_snapshots.status
|
||||
# был константой 'active' у всех строк, 394 299). Тот же
|
||||
# SAVEPOINT, что и UPDATE флага: снимок без флага (или
|
||||
# наоборот) невозможен. price_rub из snapshot-SELECT —
|
||||
# .get() консервативен ради mock-снапшотов старых тестов.
|
||||
gone_price = row.get("price_rub")
|
||||
if gone_price is not None:
|
||||
upsert_listing_snapshot(
|
||||
db,
|
||||
listing_id=row["id"],
|
||||
price_rub=gone_price,
|
||||
run_id=run_id,
|
||||
status="closed",
|
||||
)
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"avito_detail_backfill: run_id=%d failed to mark listing %s "
|
||||
|
|
|
|||
|
|
@ -9,6 +9,9 @@
|
|||
TTL=30. novostroyki (9659 активных первичных строк) и NULL-сегмент не трогаем.
|
||||
- avito: все сегменты (segments=None), TTL=10 дней -- поведение без изменений.
|
||||
- Строки НЕ удаляются -- история нужна для бэктеста (#667).
|
||||
- #2674: деактивация в той же транзакции пишет снимок listings_snapshots со статусом
|
||||
'stale' за текущую дату -- «мы N суток не видели». Жёсткое 'closed' (площадка
|
||||
ответила 404) пишет только avito_detail_backfill: смешивать факт с догадкой дорого.
|
||||
|
||||
Задача синхронная (DB-only, никаких внешних HTTP-вызовов) -- запускается kit-scheduler'ом
|
||||
через product_handlers._job_deactivate_stale (wildcard-handler deactivate_stale_*),
|
||||
|
|
@ -41,21 +44,62 @@ logger = logging.getLogger(__name__)
|
|||
# только реальный скрейп).
|
||||
_ALLOWED_STALENESS_COLUMNS = frozenset({"last_seen_at", "scraped_at"})
|
||||
|
||||
# ── Снимок «протухло» в дневной истории (#2674) ───────────────────────────────
|
||||
# listings_snapshots.status до этого фикса был константой 'active' у всех строк
|
||||
# (394 299 на момент находки) — оба места вызова upsert_listing_snapshot передавали
|
||||
# литерал 'active', и это честно: там объявление ДЕЙСТВИТЕЛЬНО видели. А деактивация
|
||||
# TTL-задачей не оставляла в истории вообще никакого следа. Из-за этого дата снятия
|
||||
# объявления (лучший доступный сигнал «скорее всего продано») не запрашивалась из
|
||||
# истории, а восстанавливалась на глаз: последний показ + предполагаемый срок жизни.
|
||||
#
|
||||
# ПОЧЕМУ 'stale', А НЕ 'closed'. Эта задача НЕ знает, что объявление снято, — она
|
||||
# знает только, что МЫ его N суток не видели, а это разные факты, когда TTL короче
|
||||
# простоя обхода. Замер: прогон по домклику 02.08 снял 6131 объявление за раз (TTL
|
||||
# 14 суток против 12 суток простоя обхода) — с общим статусом это были бы 6131
|
||||
# фальшивая «дата продажи» одной датой. Продукт про цены, смешивать факт с догадкой
|
||||
# дорого. Поэтому:
|
||||
# 'closed' — только путь 404: площадка ответила «нет» (avito_detail_backfill);
|
||||
# 'stale' — этот путь: «мы N суток не смотрели».
|
||||
# Дата всё равно фиксируется, но читатель отличает одно от другого. Ограничения
|
||||
# CHECK на колонке нет (проверено на проде), миграция не нужна — только COMMENT.
|
||||
#
|
||||
# Пишем снимок в ТОЙ ЖЕ транзакции, что и UPDATE флага: деактивация без снимка (или
|
||||
# наоборот) невозможна по построению — один statement, data-modifying CTE.
|
||||
# 1:1 по строкам: `stale` возвращает уникальные listings.id (PK), каждая даёт ровно
|
||||
# одну затронутую строку listings_snapshots (INSERT либо DO UPDATE — оба считаются
|
||||
# в rowcount), поэтому rowcount statement'а по-прежнему равен числу деактивированных.
|
||||
# price_rub берём из listings (NOT NULL в схеме) — это последняя известная цена.
|
||||
# ON CONFLICT: если снимок за сегодня уже есть (объявление видели активным утром,
|
||||
# а вечером сработал TTL) — только переводим статус в 'stale', цену не переписываем.
|
||||
_STALE_SNAPSHOT_TAIL = """
|
||||
INSERT INTO listings_snapshots
|
||||
(listing_id, snapshot_date, run_id, price_rub, status, observed_at)
|
||||
SELECT id, CURRENT_DATE, CAST(:run_id AS bigint), price_rub, 'stale', NOW()
|
||||
FROM stale
|
||||
ON CONFLICT (listing_id, snapshot_date) DO UPDATE SET
|
||||
status = 'stale',
|
||||
observed_at = EXCLUDED.observed_at
|
||||
"""
|
||||
|
||||
|
||||
def _build_all_segments_sql(staleness_column: str) -> Any:
|
||||
"""UPDATE без фильтра по сегменту: все сегменты для данного source.
|
||||
|
||||
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings,
|
||||
поэтому f-string-подстановка имени колонки безопасна. Значения (:listing_source,
|
||||
:ttl_days) остаются param-binding — psycopg v3 safe (никаких :param::type).
|
||||
:ttl_days, :run_id) остаются param-binding — psycopg v3 safe (никаких :param::type).
|
||||
"""
|
||||
return text(
|
||||
f"""
|
||||
UPDATE listings
|
||||
SET is_active = false
|
||||
WHERE source = :listing_source
|
||||
AND is_active = true
|
||||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||||
WITH stale AS (
|
||||
UPDATE listings
|
||||
SET is_active = false
|
||||
WHERE source = :listing_source
|
||||
AND is_active = true
|
||||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||||
RETURNING id, price_rub
|
||||
)
|
||||
{_STALE_SNAPSHOT_TAIL}
|
||||
"""
|
||||
)
|
||||
|
||||
|
|
@ -68,12 +112,16 @@ def _build_segments_sql(staleness_column: str) -> Any:
|
|||
"""
|
||||
return text(
|
||||
f"""
|
||||
UPDATE listings
|
||||
SET is_active = false
|
||||
WHERE source = :listing_source
|
||||
AND is_active = true
|
||||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||||
AND listing_segment = ANY(CAST(:segments AS text[]))
|
||||
WITH stale AS (
|
||||
UPDATE listings
|
||||
SET is_active = false
|
||||
WHERE source = :listing_source
|
||||
AND is_active = true
|
||||
AND {staleness_column} < NOW() - CAST(:ttl_days || ' days' AS interval)
|
||||
AND listing_segment = ANY(CAST(:segments AS text[]))
|
||||
RETURNING id, price_rub
|
||||
)
|
||||
{_STALE_SNAPSHOT_TAIL}
|
||||
"""
|
||||
)
|
||||
|
||||
|
|
@ -111,9 +159,10 @@ def deactivate_stale_listings(
|
|||
свежесть = scraped_at (двигается только реальным скрейпом).
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||||
Один UPDATE в транзакции. Финализирует scrape_runs (mark_done / mark_failed).
|
||||
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
||||
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
|
||||
|
||||
Returns {"deactivated": N} -- количество обновлённых строк.
|
||||
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
||||
|
||||
Raises:
|
||||
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
|
||||
|
|
@ -140,10 +189,15 @@ def deactivate_stale_listings(
|
|||
"listing_source": listing_source,
|
||||
"ttl_days": ttl_days,
|
||||
"segments": segments,
|
||||
"run_id": run_id,
|
||||
}
|
||||
result = db.execute(_build_segments_sql(staleness_column), params)
|
||||
else:
|
||||
params = {"listing_source": listing_source, "ttl_days": ttl_days}
|
||||
params = {
|
||||
"listing_source": listing_source,
|
||||
"ttl_days": ttl_days,
|
||||
"run_id": run_id,
|
||||
}
|
||||
result = db.execute(_build_all_segments_sql(staleness_column), params)
|
||||
|
||||
counters["deactivated"] = result.rowcount or 0
|
||||
|
|
|
|||
|
|
@ -104,13 +104,51 @@ _SNAPSHOT_SQL = text(
|
|||
"""
|
||||
)
|
||||
|
||||
# ── Event diff: price_change ──────────────────────────────────────────────────
|
||||
# Для каждого источника сравниваем сегодняшнюю цену (snapshot_date = CURRENT_DATE) с
|
||||
# самым свежим ПРЕДЫДУЩИМ снимком (snapshot_date < CURRENT_DATE). Если цена изменилась
|
||||
# (обе NOT NULL, old <> 0) — пишем price_change.
|
||||
# ── Event diff: три выводимых типа событий из пяти в схеме ────────────────────
|
||||
# Для каждого источника сравниваем сегодняшний снимок (snapshot_date = CURRENT_DATE) с
|
||||
# самым свежим ПРЕДЫДУЩИМ (snapshot_date < CURRENT_DATE).
|
||||
# today — снимок за сегодня (только что записан _SNAPSHOT_SQL, в той же транзакции).
|
||||
# p — последний снимок строго ДО сегодня, per-row LATERAL point-lookup (#2607).
|
||||
#
|
||||
# #2674: схема (079) знает пять типов событий, писатель умел один — price_change,
|
||||
# 8288 строк. Дописаны два:
|
||||
# edited — payload_hash изменился, а цена нет (изменение цены уже описано
|
||||
# отдельным событием price_change — дублировать его как «редактирование»
|
||||
# значило бы считать одно изменение дважды). Прошлый хеш обязан быть
|
||||
# непустым: md5(NULL) = NULL, и «payload появился впервые» — это не
|
||||
# правка, а первое наблюдение;
|
||||
# first_seen — предыдущего снимка нет вовсе (LEFT JOIN LATERAL даёт p.* = NULL).
|
||||
#
|
||||
# delisted и relisted НЕ ПИШУТСЯ НАМЕРЕННО — они НЕ ВЫВОДИМЫ из наших данных.
|
||||
# is_active в снимке — derived-признак «last_seen_at свежее FRESHNESS_WINDOW_DAYS»,
|
||||
# то есть «мы видели», а не «объявление есть на площадке». При покрытии обхода 10-35%
|
||||
# такой переход рождается тем, что скрейпер СНОВА ДОШЁЛ до источника, а не тем, что
|
||||
# объявление вернулось/ушло. Контрольная группа в наших же данных (14-18.07):
|
||||
# domklik, покрытие 99.9-100%: снятий 1/2/0/2/4 в сутки, возвратов — РОВНО 0 все дни;
|
||||
# yandex, покрытие 34-43%: снятий 343-433 в сутки, возвратов до 155.
|
||||
# Тот же обход, тот же день — разница только в покрытии. Отсюда же всплески:
|
||||
# avito 13.07 (день остановки обхода) — 3023 «снятия» за сутки против контрольной
|
||||
# ставки 1-4, точность события ≈4%; 4705 «возвратов» из 5493 за 12 дней (86%) — это
|
||||
# два дня после возобновления обхода 2-3.08.
|
||||
# Сузить окно свежести НЕ поможет — станет хуже (больше флапаний); окно шире
|
||||
# максимального интервала повторного визита обессмысливает само событие.
|
||||
# Честный ответ схеме — не писать эти два типа, а не наполнять журнал догадками.
|
||||
# Единственный жёсткий сигнал снятия — 404 при поштучном обходе, он пишется в
|
||||
# listings_snapshots.status='closed' (avito_detail_backfill).
|
||||
#
|
||||
# Оставшиеся три события утверждают факты о НАШИХ СОБСТВЕННЫХ строках («появился новый
|
||||
# источник», «хеш изменился при той же цене», «цена другая»), а не о поведении площадки.
|
||||
#
|
||||
# JOIN → LEFT JOIN LATERAL: без LEFT источники без предыдущего снимка отбрасывались
|
||||
# join'ом, поэтому first_seen был недостижим по построению. План #2607 не меняется —
|
||||
# LEFT JOIN LATERAL так же форсирует per-row индексный point-lookup по
|
||||
# idx_lss_source_date, просто не отбрасывает строку при отсутствии предыдущей.
|
||||
#
|
||||
# Ветки разворачиваются CROSS JOIN LATERAL (VALUES ...) — одна строка сравнения даёт
|
||||
# до трёх строк-кандидатов, из которых WHERE e.fires оставляет сработавшие. Это
|
||||
# по-прежнему ОДИН set-based statement (никакого Python-цикла), просто три предиката
|
||||
# вместо одного.
|
||||
#
|
||||
# #2607: раньше `p` был отдельным CTE `DISTINCT ON (listing_source_id) ... FROM
|
||||
# listing_source_snapshots WHERE snapshot_date < CURRENT_DATE` и джойнился обычным JOIN.
|
||||
# Планировщик оценивает `today` в 1 строку (свежевставленные в этой же транзакции строки
|
||||
|
|
@ -127,38 +165,79 @@ _SNAPSHOT_SQL = text(
|
|||
#
|
||||
# Полностью set-based: один INSERT … SELECT по всем источникам, без Python-цикла (LATERAL
|
||||
# — это внутренний план Postgres, не Python-итерация).
|
||||
# change_time = now() детерминирует UNIQUE(listing_source_id, change_time, event_type)
|
||||
# в пределах прогона → ON CONFLICT DO NOTHING делает писатель идемпотентным.
|
||||
#
|
||||
# change_time = date_trunc('day', now()), а НЕ now() (#2674): с now() уникальность
|
||||
# UNIQUE(listing_source_id, change_time, event_type) работала только ВНУТРИ прогона —
|
||||
# второй прогон в те же сутки перезаписывал сегодняшний снимок, предикаты срабатывали
|
||||
# заново с другим временем и давали дубли (2 августа таких прогонов было два).
|
||||
# Суточная гранулярность честнее для суточного же сравнения и включает заявленную
|
||||
# идемпотентность: ON CONFLICT DO NOTHING теперь действительно гасит повтор за день.
|
||||
#
|
||||
# NULLIF(p.price_rub, 0) в diff_percent обязателен: выражения VALUES вычисляются ДО
|
||||
# фильтра `WHERE e.fires`, поэтому предикат "p.price_rub <> 0" от деления на ноль уже
|
||||
# не спасает — без NULLIF первый же источник с нулевой прошлой ценой уронил бы весь
|
||||
# прогон. Результат при этом тот же: строка с NULL-диффом не проходит e.fires.
|
||||
#
|
||||
# Внешний SELECT над data-modifying CTE считает вставленное ПО ТИПАМ (RETURNING отдаёт
|
||||
# только реально вставленные строки, не съеденные ON CONFLICT), сразу в виде ключей
|
||||
# счётчиков `<event_type>_events` — писатель получает готовый dict без Python-агрегации.
|
||||
# Ровно этот счётчик и показал бы четыре нуля из пяти, если бы существовал раньше.
|
||||
_EVENT_DIFF_SQL = text(
|
||||
"""
|
||||
WITH today AS (
|
||||
SELECT listing_source_id, price_rub
|
||||
SELECT listing_source_id, price_rub, payload_hash
|
||||
FROM listing_source_snapshots
|
||||
WHERE snapshot_date = CURRENT_DATE
|
||||
),
|
||||
inserted AS (
|
||||
INSERT INTO listing_source_events (
|
||||
listing_source_id, change_time, event_type, price_rub, diff_percent
|
||||
)
|
||||
SELECT
|
||||
t.listing_source_id,
|
||||
date_trunc('day', now()),
|
||||
e.event_type,
|
||||
t.price_rub,
|
||||
e.diff_percent
|
||||
FROM today t
|
||||
LEFT JOIN LATERAL (
|
||||
SELECT s.snapshot_date, s.price_rub, s.payload_hash
|
||||
FROM listing_source_snapshots s
|
||||
WHERE s.listing_source_id = t.listing_source_id
|
||||
AND s.snapshot_date < CURRENT_DATE
|
||||
ORDER BY s.snapshot_date DESC
|
||||
LIMIT 1
|
||||
) p ON true
|
||||
CROSS JOIN LATERAL (VALUES
|
||||
(
|
||||
'first_seen',
|
||||
NULL::numeric,
|
||||
p.snapshot_date IS NULL
|
||||
),
|
||||
(
|
||||
'price_change',
|
||||
round((t.price_rub - p.price_rub)::numeric
|
||||
/ NULLIF(p.price_rub, 0) * 100, 4),
|
||||
t.price_rub IS NOT NULL
|
||||
AND p.price_rub IS NOT NULL
|
||||
AND p.price_rub <> 0
|
||||
AND t.price_rub <> p.price_rub
|
||||
),
|
||||
(
|
||||
'edited',
|
||||
NULL::numeric,
|
||||
p.payload_hash IS NOT NULL
|
||||
AND t.payload_hash IS DISTINCT FROM p.payload_hash
|
||||
AND t.price_rub IS NOT DISTINCT FROM p.price_rub
|
||||
)
|
||||
) AS e(event_type, diff_percent, fires)
|
||||
WHERE e.fires
|
||||
ON CONFLICT (listing_source_id, change_time, event_type) DO NOTHING
|
||||
RETURNING event_type
|
||||
)
|
||||
INSERT INTO listing_source_events (
|
||||
listing_source_id, change_time, event_type, price_rub, diff_percent
|
||||
)
|
||||
SELECT
|
||||
t.listing_source_id,
|
||||
now(),
|
||||
'price_change',
|
||||
t.price_rub,
|
||||
round((t.price_rub - p.price_rub)::numeric / p.price_rub * 100, 4)
|
||||
FROM today t
|
||||
JOIN LATERAL (
|
||||
SELECT s.price_rub
|
||||
FROM listing_source_snapshots s
|
||||
WHERE s.listing_source_id = t.listing_source_id
|
||||
AND s.snapshot_date < CURRENT_DATE
|
||||
ORDER BY s.snapshot_date DESC
|
||||
LIMIT 1
|
||||
) p ON true
|
||||
WHERE t.price_rub IS NOT NULL
|
||||
AND p.price_rub IS NOT NULL
|
||||
AND p.price_rub <> 0
|
||||
AND t.price_rub <> p.price_rub
|
||||
ON CONFLICT (listing_source_id, change_time, event_type) DO NOTHING
|
||||
SELECT event_type || '_events' AS counter_key, count(*) AS n
|
||||
FROM inserted
|
||||
GROUP BY 1
|
||||
"""
|
||||
)
|
||||
|
||||
|
|
@ -166,12 +245,14 @@ _EVENT_DIFF_SQL = text(
|
|||
def snapshot_listing_sources(
|
||||
db: Session, run_id: int, params: dict[str, Any] | None = None
|
||||
) -> dict[str, int]:
|
||||
"""Записать дневной снимок listing_sources + price_change-события.
|
||||
"""Записать дневной снимок listing_sources + события изменений.
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как import_rosreestr_dkp).
|
||||
Два set-based statement'а в одной транзакции:
|
||||
1. upsert снимка на (listing_source_id, CURRENT_DATE) — last-write-wins.
|
||||
2. diff сегодняшней цены против последнего предыдущего снимка → price_change-события.
|
||||
2. diff сегодняшнего снимка против последнего предыдущего → три события,
|
||||
выводимые из наших данных (#2674). delisted/relisted схема разрешает, но
|
||||
они НЕ выводимы при покрытии обхода 10-35% — см. _EVENT_DIFF_SQL.
|
||||
|
||||
Params (из default_params jsonb в scrape_schedules, #2607):
|
||||
budget_sec: float — SET LOCAL statement_timeout на транзакцию (default 900,
|
||||
|
|
@ -183,11 +264,18 @@ def snapshot_listing_sources(
|
|||
|
||||
Финализирует scrape_runs (mark_done / mark_failed) и пишет counters.
|
||||
|
||||
Returns {"snapshotted": N, "price_change_events": M}.
|
||||
Returns {"snapshotted": N, "<event_type>_events": M} — по счётчику на каждый из
|
||||
трёх пишущихся типов, всегда все три ключа (тип, который за прогон не сработал
|
||||
ни разу, честно показывает 0, а не пропадает из counters).
|
||||
"""
|
||||
params = params or {}
|
||||
budget_sec = _clamp_budget_sec(params.get("budget_sec", DEFAULT_BUDGET_SEC))
|
||||
counters: dict[str, int] = {"snapshotted": 0, "price_change_events": 0}
|
||||
counters: dict[str, int] = {
|
||||
"snapshotted": 0,
|
||||
"price_change_events": 0,
|
||||
"edited_events": 0,
|
||||
"first_seen_events": 0,
|
||||
}
|
||||
try:
|
||||
# statement_timeout НЕ принимает bind-параметр ($1/:name) — синтаксис Postgres SET
|
||||
# запрещает placeholder на этом месте (проверено вживую на проде: "syntax error at
|
||||
|
|
@ -205,17 +293,15 @@ def snapshot_listing_sources(
|
|||
)
|
||||
counters["snapshotted"] = snap_result.rowcount or 0
|
||||
|
||||
event_result = db.execute(_EVENT_DIFF_SQL)
|
||||
counters["price_change_events"] = event_result.rowcount or 0
|
||||
# Statement возвращает уже готовые пары (counter_key, n) по типам событий —
|
||||
# dict(...) без Python-агрегации, набор ключей задан инициализацией counters
|
||||
# выше, так что не сработавшие типы остаются нулями, а не исчезают.
|
||||
event_rows = db.execute(_EVENT_DIFF_SQL).fetchall()
|
||||
counters.update(dict(event_rows))
|
||||
|
||||
db.commit()
|
||||
runs_mod.mark_done(db, run_id, counters)
|
||||
logger.info(
|
||||
"snapshot_listing_sources run_id=%d done: snapshotted=%d price_change_events=%d",
|
||||
run_id,
|
||||
counters["snapshotted"],
|
||||
counters["price_change_events"],
|
||||
)
|
||||
logger.info("snapshot_listing_sources run_id=%d done: %s", run_id, counters)
|
||||
return counters
|
||||
except Exception as exc:
|
||||
logger.exception(
|
||||
|
|
|
|||
|
|
@ -0,0 +1,29 @@
|
|||
-- 213_listings_snapshots_status_vocab.sql
|
||||
-- #2674 — словарь listings_snapshots.status стал трёхзначным: 'stale' ≠ 'closed'.
|
||||
--
|
||||
-- ПРОБЛЕМА: комментарий колонки (016) обещал два значения — 'active' / 'closed' —
|
||||
-- и при этом 'closed' не писал никто и никогда: status был константой 'active' у
|
||||
-- всех 394 704 строк при 55 448 реально неактивных объявлениях. Писатель «снято»
|
||||
-- появился в #2674, но одним значением обойтись нельзя:
|
||||
-- - путь 404 (avito_detail_backfill) ЗНАЕТ, что объявления нет: площадка ответила;
|
||||
-- - путь TTL (deactivate_stale_listings) знает только, что МЫ N суток не смотрели.
|
||||
-- Замер: прогон по домклику 02.08 деактивировал 6131 объявление за раз (TTL 14 суток
|
||||
-- против 12 суток простоя обхода) — под общим статусом это 6131 фальшивая «дата
|
||||
-- продажи» одной датой. Продукт про цены: смешивать факт с догадкой дорого.
|
||||
--
|
||||
-- ДЕЛАЕТ: только обновляет COMMENT — сама колонка `text` без CHECK, DDL не нужен.
|
||||
-- CHECK намеренно НЕ добавляем: 394 704 существующие строки валидны, а жёсткий
|
||||
-- словарь на историческую таблицу — деструктивный риск ради нулевой выгоды.
|
||||
--
|
||||
-- Idempotent: COMMENT ON COLUMN — безусловная перезапись, безопасно повторно.
|
||||
-- Apply after: 212_sber_index_pull_weekly.sql
|
||||
|
||||
BEGIN;
|
||||
|
||||
COMMENT ON COLUMN listings_snapshots.status IS
|
||||
'''active'' = объявление видели в прогоне. '
|
||||
'''closed'' = площадка ответила 404 на поштучном обходе (жёсткий факт снятия). '
|
||||
'''stale'' = TTL-деактивация: мы N суток не смотрели (догадка, НЕ дата продажи). '
|
||||
'NULL = неизвестно. Словарь расширен в #2674 — до него писалось только ''active''.';
|
||||
|
||||
COMMIT;
|
||||
431
tradein-mvp/backend/tests/test_2674_writers_honor_schema.py
Normal file
431
tradein-mvp/backend/tests/test_2674_writers_honor_schema.py
Normal file
|
|
@ -0,0 +1,431 @@
|
|||
"""#2674 — писатели наконец пишут то, что обещает схема.
|
||||
|
||||
Три находки одного класса: колонка есть, писатель есть, тест на писателя зелёный,
|
||||
а данные не появляются. Обычный юнит-тест такое не ловит по построению — он
|
||||
проверяет то, что автор себе представлял. Ловится это сверкой «что схема обещает»
|
||||
с «что писатель реально перечисляет», поэтому тесты ниже читают миграции и
|
||||
сравнивают их с SQL писателя, а не повторяют его же список колонок.
|
||||
|
||||
Числа с прода на 2026-08-05/06 (до фикса):
|
||||
1. house_suggestions — 25 055 строк, image_link/area_m2/rooms/floor/total_floors
|
||||
заполнены у 0 из них (колонки с миграции 064, ~74 дня).
|
||||
2. listings_snapshots.status — 'active' у всех 394 704 строк при 55 448 реально
|
||||
неактивных объявлений; ни 'closed', ни 'stale' не писал никто и никогда.
|
||||
3. listing_source_events — 8288 строк, все price_change; edited/first_seen —
|
||||
ноль за всё время.
|
||||
|
||||
Отдельный класс тестов — гейты на то, что писатель НЕ пишет: delisted/relisted схема
|
||||
разрешает, но при покрытии обхода 10-35% они неотличимы от «скрейпер снова дошёл»
|
||||
(контроль — домклик со 100% покрытием: 0 возвратов за 5 суток), а TTL-путь не имеет
|
||||
права называть протухание снятием. Журнал и история из догадок хуже пустых.
|
||||
|
||||
БД и сеть замоканы — реального Postgres не нужно.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import inspect
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
from scraper_kit.providers.avito.imv import _parse_suggestion
|
||||
|
||||
from app.services import house_imv_backfill as hib
|
||||
from app.tasks import deactivate_stale_avito as deact_mod
|
||||
from app.tasks import listing_source_snapshot as snap_mod
|
||||
|
||||
_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql"
|
||||
_FIXTURES = Path(__file__).resolve().parent / "fixtures"
|
||||
|
||||
|
||||
# ── Общие хелперы: схема vs writer ────────────────────────────────────────────
|
||||
|
||||
|
||||
def _declared_columns(migration: str, table: str) -> set[str]:
|
||||
"""Имена колонок из CREATE TABLE IF NOT EXISTS <table> ( ... ); в миграции."""
|
||||
body = migration.split(f"CREATE TABLE IF NOT EXISTS {table} (")[1].split("\n);")[0]
|
||||
cols: set[str] = set()
|
||||
for line in body.splitlines():
|
||||
m = re.match(r"\s+([a-z_]+)\s+[a-z]", line)
|
||||
if m:
|
||||
cols.add(m.group(1))
|
||||
return cols
|
||||
|
||||
|
||||
def _insert_columns(sql: str, table: str) -> set[str]:
|
||||
"""Имена колонок из INSERT INTO <table> ( ... ) VALUES."""
|
||||
m = re.search(rf"INSERT INTO {table}\s*\(([^)]*)\)", sql, re.S)
|
||||
assert m is not None, f"не найден INSERT INTO {table}"
|
||||
return {c.strip() for c in m.group(1).split(",") if c.strip()}
|
||||
|
||||
|
||||
# ══ 1. Фотографии подсказок Avito IMV ═════════════════════════════════════════
|
||||
|
||||
|
||||
def test_suggestion_parser_keeps_image_link() -> None:
|
||||
"""imageLink из ответа площадки доезжает до модели, а не выбрасывается."""
|
||||
sugg = _parse_suggestion(
|
||||
{
|
||||
"id": 8000753763,
|
||||
"title": "3-к. квартира, 61,6 м², 1/5 эт.",
|
||||
"price": 7400000,
|
||||
"imageLink": "https://80.img.avito.st/image/1/abc",
|
||||
}
|
||||
)
|
||||
assert sugg.image_link == "https://80.img.avito.st/image/1/abc"
|
||||
# В raw_payload ссылку по-прежнему не дублируем — у неё теперь своя колонка.
|
||||
assert "imageLink" not in (sugg.raw_payload or {})
|
||||
|
||||
|
||||
def test_suggestion_parser_derives_metrics_from_title() -> None:
|
||||
"""rooms/area_m2/floor/total_floors парсятся из title тем же путём, что у
|
||||
placementHistory (колонки house_suggestions существуют с миграции 064)."""
|
||||
sugg = _parse_suggestion(
|
||||
{"id": 1, "title": "3-к. квартира, 61,6 м², 1/5 эт.", "price": 7400000}
|
||||
)
|
||||
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (3, 61.6, 1, 5)
|
||||
|
||||
|
||||
def test_suggestion_parser_survives_unparsable_title() -> None:
|
||||
"""Нераспознанный заголовок → None'ы, а не исключение (строка всё равно пишется)."""
|
||||
sugg = _parse_suggestion({"id": 2, "title": "Апартаменты", "price": 1})
|
||||
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (None, None, None, None)
|
||||
|
||||
|
||||
def test_studio_title_keeps_area_and_floors() -> None:
|
||||
"""Студия не роняет разбор целиком: 1991 заголовок из 25 055 (7.9%) — без комнатности.
|
||||
|
||||
Обязательная группа комнатности обнуляла ВСЕ ЧЕТЫРЕ поля, хотя площадь и этажность
|
||||
в заголовке есть. rooms=0 — конвенция kit'а («0 = студия»), а не «неизвестно».
|
||||
"""
|
||||
sugg = _parse_suggestion(
|
||||
{"id": 3, "title": "Квартира-студия, 34,2 м², 9/10 эт.", "price": 3_500_000}
|
||||
)
|
||||
assert (sugg.rooms, sugg.area_m2, sugg.floor, sugg.total_floors) == (0, 34.2, 9, 10)
|
||||
|
||||
|
||||
def test_placement_history_gets_same_title_fix() -> None:
|
||||
"""Тот же регексп чинит второго писателя — house_placement_history (8.8% без площади)."""
|
||||
from scraper_kit.providers.avito.imv import _parse_placement_item
|
||||
|
||||
item = _parse_placement_item({"id": 9, "title": "Квартира-студия, 28 м², 2/17 эт."})
|
||||
assert (item.rooms, item.area_m2, item.floor, item.total_floors) == (0, 28.0, 2, 17)
|
||||
|
||||
|
||||
def test_live_fixture_suggestions_carry_image_link() -> None:
|
||||
"""Живой capture avito_imv_getdata.json: у подсказок реально есть imageLink."""
|
||||
data = json.loads((_FIXTURES / "avito_imv_getdata.json").read_text("utf-8"))
|
||||
items = data["suggestions"]["items"]
|
||||
parsed = [_parse_suggestion(raw) for raw in items]
|
||||
assert parsed, "фикстура без подсказок — тест бессмыслен"
|
||||
assert all(s.image_link for s in parsed)
|
||||
|
||||
|
||||
def test_house_suggestions_insert_covers_every_declared_column() -> None:
|
||||
"""Regression-гейт на весь класс бага: INSERT обязан перечислять КАЖДУЮ колонку
|
||||
house_suggestions из миграции 064 (кроме автоинкрементного id).
|
||||
|
||||
Именно этот тест покраснел бы 74 дня назад: image_link (и заодно area_m2/rooms/
|
||||
floor/total_floors) объявлены схемой, но в запрос вставки не входили — 25 055
|
||||
строк с NULL. Тест не дублирует список колонок писателя, а сверяет его со схемой,
|
||||
поэтому ловит и следующую забытую колонку.
|
||||
"""
|
||||
declared = _declared_columns(
|
||||
(_SQL_DIR / "064_house_imv_phase_c.sql").read_text("utf-8"), "house_suggestions"
|
||||
)
|
||||
written = _insert_columns(inspect.getsource(hib.save_imv_result), "house_suggestions")
|
||||
assert (
|
||||
declared - {"id"} <= written
|
||||
), f"колонки без писателя: {sorted(declared - {'id'} - written)}"
|
||||
|
||||
|
||||
def test_save_imv_result_binds_image_link_and_metrics() -> None:
|
||||
"""save_imv_result передаёт значения подсказки в bind-параметры (не только в SQL)."""
|
||||
sugg = _parse_suggestion(
|
||||
{
|
||||
"id": 777,
|
||||
"title": "2-к. квартира, 42 м², 4/5 эт.",
|
||||
"price": 6300000,
|
||||
"imageLink": "https://img/x.jpg",
|
||||
}
|
||||
)
|
||||
result = MagicMock(
|
||||
cache_key="k",
|
||||
recommended_price=1,
|
||||
lower_price=1,
|
||||
higher_price=1,
|
||||
market_count=1,
|
||||
raw_response=None,
|
||||
placement_history=[],
|
||||
suggestions=[sugg],
|
||||
)
|
||||
params = {
|
||||
"rooms": 2,
|
||||
"area_m2": 42.0,
|
||||
"floor": 4,
|
||||
"floor_at_home": 5,
|
||||
"house_type": "panel",
|
||||
"renovation_type": "cosmetic",
|
||||
"has_balcony": True,
|
||||
"has_loggia": False,
|
||||
}
|
||||
db = MagicMock()
|
||||
hib.save_imv_result(db, house_id=1, params=params, result=result)
|
||||
|
||||
sugg_calls = [
|
||||
c for c in db.execute.call_args_list if "INSERT INTO house_suggestions" in str(c.args[0])
|
||||
]
|
||||
assert len(sugg_calls) == 1
|
||||
bound = sugg_calls[0].args[1]
|
||||
assert bound["img"] == "https://img/x.jpg"
|
||||
assert (bound["rooms"], bound["area"], bound["floor"], bound["total_floors"]) == (2, 42.0, 4, 5)
|
||||
|
||||
|
||||
# ══ 2. «Снято» и «протухло» в дневной истории объявлений ══════════════════════
|
||||
|
||||
_DEACT_SQL_BUILDERS = (deact_mod._build_all_segments_sql, deact_mod._build_segments_sql)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("build", _DEACT_SQL_BUILDERS)
|
||||
@pytest.mark.parametrize("column", sorted(deact_mod._ALLOWED_STALENESS_COLUMNS))
|
||||
def test_deactivation_writes_stale_snapshot_in_same_statement(build: Any, column: str) -> None:
|
||||
"""Деактивация и снимок 'stale' — один statement, значит одна транзакция.
|
||||
|
||||
До #2674 задача только двигала флаг: в listings_snapshots не появлялось ничего,
|
||||
и дата снятия объявления (лучший сигнал «скорее всего продано») восстанавливалась
|
||||
на глаз из последнего показа + предполагаемого срока жизни.
|
||||
"""
|
||||
sql = str(build(column).text)
|
||||
assert "SET is_active = false" in sql
|
||||
assert "RETURNING id, price_rub" in sql
|
||||
assert "INSERT INTO listings_snapshots" in sql
|
||||
assert "'stale'" in sql
|
||||
# Снимок пишется по строкам, которые вернул сам UPDATE, — не отдельной выборкой.
|
||||
assert "FROM stale" in sql
|
||||
# Идемпотентность: повторный прогон в те же сутки не падает на PK.
|
||||
assert "ON CONFLICT (listing_id, snapshot_date) DO UPDATE" in sql
|
||||
# psycopg v3: никаких :param::type.
|
||||
assert not re.search(r":\w+::", sql)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("build", _DEACT_SQL_BUILDERS)
|
||||
def test_ttl_path_never_claims_closed(build: Any) -> None:
|
||||
"""TTL-путь НЕ имеет права писать 'closed' — он не знает, что объявление снято.
|
||||
|
||||
Замер: прогон по домклику 02.08 деактивировал 6131 объявление за раз (TTL 14 суток
|
||||
против 12 суток простоя обхода). Под статусом 'closed' это 6131 фальшивая «дата
|
||||
продажи» одной датой. 'closed' остаётся только за 404 — там ответила площадка.
|
||||
"""
|
||||
assert "'closed'" not in str(build("last_seen_at").text)
|
||||
|
||||
|
||||
class _FakeResult:
|
||||
def __init__(self, rowcount: int) -> None:
|
||||
self.rowcount = rowcount
|
||||
|
||||
|
||||
class _FakeDB:
|
||||
def __init__(self, rowcount: int = 0) -> None:
|
||||
self._rowcount = rowcount
|
||||
self.executed: list[tuple[Any, Any]] = []
|
||||
self.committed = False
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||
self.executed.append((stmt, params))
|
||||
return _FakeResult(self._rowcount)
|
||||
|
||||
def commit(self) -> None:
|
||||
self.committed = True
|
||||
|
||||
def rollback(self) -> None: # pragma: no cover — путь ошибки тут не проверяется
|
||||
pass
|
||||
|
||||
|
||||
def test_deactivate_stale_listings_threads_run_id_into_snapshot(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""run_id доезжает до снимка — провенанс «каким прогоном закрыто» не теряется."""
|
||||
monkeypatch.setattr(deact_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||
monkeypatch.setattr(deact_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
db = _FakeDB(rowcount=7)
|
||||
|
||||
out = deact_mod.deactivate_stale_listings(
|
||||
db, # type: ignore[arg-type]
|
||||
run_id=555,
|
||||
listing_source="avito",
|
||||
ttl_days=10,
|
||||
)
|
||||
|
||||
assert out == {"deactivated": 7}
|
||||
assert len(db.executed) == 1, "деактивация и снимок обязаны быть одним statement'ом"
|
||||
stmt, bound = db.executed[0]
|
||||
assert "INSERT INTO listings_snapshots" in str(stmt)
|
||||
assert bound is not None and bound["run_id"] == 555
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_avito_404_records_closed_snapshot() -> None:
|
||||
"""404 с площадки — самый достоверный сигнал снятия; он тоже попадает в историю."""
|
||||
from scraper_kit.avito_exceptions import AvitoListingGoneError
|
||||
|
||||
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
|
||||
|
||||
db = MagicMock()
|
||||
sel = MagicMock()
|
||||
sel.mappings.return_value.all.return_value = [
|
||||
{"id": 42, "source_url": "/items/42", "price_rub": 5_000_000}
|
||||
]
|
||||
db.execute.return_value = sel
|
||||
fake_settings = MagicMock(scraper_fetch_mode="cffi", avito_detail_backfill_use_curl=False)
|
||||
|
||||
with (
|
||||
patch("app.tasks.avito_detail_backfill.settings", fake_settings),
|
||||
patch("app.tasks.avito_detail_backfill.AsyncSession", return_value=AsyncMock()),
|
||||
patch("app.tasks.avito_detail_backfill.AvitoScraper"),
|
||||
patch("app.tasks.avito_detail_backfill.runs_mod", MagicMock()),
|
||||
patch("app.tasks.avito_detail_backfill.asyncio.sleep", new_callable=AsyncMock),
|
||||
patch(
|
||||
"app.tasks.avito_detail_backfill.fetch_detail",
|
||||
AsyncMock(side_effect=AvitoListingGoneError("404 gone")),
|
||||
),
|
||||
patch("app.tasks.avito_detail_backfill.upsert_listing_snapshot") as snap,
|
||||
):
|
||||
result = await run_avito_detail_backfill(db, run_id=3, params={"budget_sec": 60})
|
||||
|
||||
assert result.gone == 1
|
||||
snap.assert_called_once()
|
||||
assert snap.call_args.kwargs["status"] == "closed"
|
||||
assert snap.call_args.kwargs["listing_id"] == 42
|
||||
assert snap.call_args.kwargs["price_rub"] == 5_000_000
|
||||
|
||||
|
||||
# ══ 3. Журнал событий объявлений ══════════════════════════════════════════════
|
||||
|
||||
|
||||
def _schema_event_types() -> set[str]:
|
||||
"""Пять типов из CHECK-констрейнта миграции 079 — источник правды."""
|
||||
sql = (_SQL_DIR / "079_listing_source_history.sql").read_text("utf-8")
|
||||
check = sql.split("event_type IN (")[1].split(")")[0]
|
||||
return set(re.findall(r"'([a-z_]+)'", check))
|
||||
|
||||
|
||||
# Два типа схемы НЕ ВЫВОДИМЫ из наших данных и намеренно не пишутся (#2674).
|
||||
# is_active в снимке значит «мы видели», а не «есть на площадке», поэтому переход
|
||||
# рождается тем, что скрейпер снова дошёл до источника. Контрольная группа за 14-18.07:
|
||||
# domklik при покрытии 99.9-100% дал возвратов РОВНО 0 и снятий 1-4 в сутки, yandex при
|
||||
# 34-43% — снятий 343-433 в сутки. Тот же обход, тот же день, разница только в покрытии.
|
||||
# Отсюда: avito 13.07 (остановка обхода) 3023 «снятия» за сутки против контрольных 1-4
|
||||
# (точность ≈4%), и 4705 «возвратов» из 5493 за 12 дней — два дня после возобновления.
|
||||
_NOT_DERIVABLE_EVENT_TYPES = {"delisted", "relisted"}
|
||||
|
||||
|
||||
def test_event_writer_covers_every_derivable_schema_event_type() -> None:
|
||||
"""Писатель обязан уметь каждый ВЫВОДИМЫЙ тип из CHECK схемы.
|
||||
|
||||
До #2674 из пяти типов писался один (price_change, 8288 строк). Дописаны два
|
||||
выводимых; два оставшихся — сознательное решение, а не забытая ветка (см.
|
||||
_NOT_DERIVABLE_EVENT_TYPES). Тест сверяет со схемой, а не с копией списка,
|
||||
поэтому покраснеет и на шестом типе, добавленном в CHECK без писателя.
|
||||
"""
|
||||
declared = _schema_event_types()
|
||||
assert len(declared) == 5, f"схема 079 изменилась: {sorted(declared)}"
|
||||
sql = str(snap_mod._EVENT_DIFF_SQL.text)
|
||||
expected = declared - _NOT_DERIVABLE_EVENT_TYPES
|
||||
missing = {t for t in expected if f"'{t}'" not in sql}
|
||||
assert not missing, f"выводимые типы без писателя: {sorted(missing)}"
|
||||
|
||||
|
||||
def test_not_derivable_events_are_never_written() -> None:
|
||||
"""delisted/relisted не пишутся: при покрытии обхода 10-35% они неотличимы от
|
||||
«скрейпер снова дошёл». Контроль — домклик со 100% покрытием: 0 возвратов за 5 суток.
|
||||
|
||||
Гейт против «дописать для полноты»: журнал из догадок хуже пустого журнала.
|
||||
"""
|
||||
sql = str(snap_mod._EVENT_DIFF_SQL.text)
|
||||
written = {t for t in _NOT_DERIVABLE_EVENT_TYPES if f"'{t}'" in sql}
|
||||
assert not written, f"невыводимые типы попали в писатель: {sorted(written)}"
|
||||
# is_active больше не читается вовсе — иначе ветка вернётся незаметно.
|
||||
assert "is_active" not in sql
|
||||
|
||||
|
||||
def test_first_seen_requires_left_join_and_derivations_use_snapshot_fields() -> None:
|
||||
"""Ветки выводятся из полей снимка, first_seen достижим только через LEFT JOIN.
|
||||
|
||||
С обычным JOIN источник без предыдущего снимка отбрасывался джойном — событие
|
||||
«первое появление» было недостижимо по построению.
|
||||
"""
|
||||
sql = str(snap_mod._EVENT_DIFF_SQL.text)
|
||||
assert "LEFT JOIN LATERAL" in sql
|
||||
assert "p.snapshot_date IS NULL" in sql # first_seen
|
||||
assert "t.payload_hash IS DISTINCT FROM p.payload_hash" in sql # edited
|
||||
# Снимок за сегодня обязан отдавать поля, из которых выводятся ветки.
|
||||
assert "SELECT listing_source_id, price_rub, payload_hash" in sql
|
||||
|
||||
|
||||
def test_event_dedup_works_across_runs_not_only_within_one() -> None:
|
||||
"""change_time усечён до суток: UNIQUE(source, change_time, type) должен гасить
|
||||
повторный прогон в те же сутки (2 августа их было два), а не только строки одного."""
|
||||
sql = str(snap_mod._EVENT_DIFF_SQL.text)
|
||||
assert "date_trunc('day', now())" in sql
|
||||
|
||||
|
||||
def test_price_change_division_guarded_by_nullif() -> None:
|
||||
"""VALUES вычисляется ДО фильтра e.fires → без NULLIF прогон падал бы на
|
||||
первом источнике с нулевой прошлой ценой (предикат p.price_rub <> 0 не спасает)."""
|
||||
sql = str(snap_mod._EVENT_DIFF_SQL.text)
|
||||
assert "NULLIF(p.price_rub, 0)" in sql
|
||||
|
||||
|
||||
class _EventFakeDB:
|
||||
"""Session-заглушка: SET LOCAL и snapshot дают rowcount, event-diff — пары счётчиков."""
|
||||
|
||||
def __init__(self, event_rows: list[tuple[str, int]]) -> None:
|
||||
self._event_rows = event_rows
|
||||
self.committed = False
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||||
result = MagicMock()
|
||||
if "listing_source_events" in str(stmt):
|
||||
result.fetchall.return_value = self._event_rows
|
||||
else:
|
||||
result.rowcount = 100
|
||||
return result
|
||||
|
||||
def commit(self) -> None:
|
||||
self.committed = True
|
||||
|
||||
def rollback(self) -> None: # pragma: no cover — путь ошибки тут не проверяется
|
||||
pass
|
||||
|
||||
|
||||
def test_counters_report_every_written_type_including_zeros(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
"""Счётчики прогона показывают все пишущиеся типы; не сработавший честно равен 0.
|
||||
|
||||
Ровно этого счётчика не хватало, чтобы заметить четыре нуля из пяти за 66 дней.
|
||||
Счётчиков НЕвыводимых типов быть не должно — иначе вечный 0 будет читаться как
|
||||
«событий не было», а не как «мы это сознательно не пишем».
|
||||
"""
|
||||
monkeypatch.setattr(snap_mod.runs_mod, "mark_done", lambda *a, **k: None)
|
||||
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
|
||||
db = _EventFakeDB([("first_seen_events", 600), ("edited_events", 17)])
|
||||
out = snap_mod.snapshot_listing_sources(db, run_id=1) # type: ignore[arg-type]
|
||||
|
||||
assert out["first_seen_events"] == 600
|
||||
assert out["edited_events"] == 17
|
||||
assert out["price_change_events"] == 0
|
||||
# Ключ на каждый пишущийся тип — иначе «ноль» неотличим от «типа нет в counters».
|
||||
for event_type in _schema_event_types() - _NOT_DERIVABLE_EVENT_TYPES:
|
||||
assert f"{event_type}_events" in out
|
||||
for event_type in _NOT_DERIVABLE_EVENT_TYPES:
|
||||
assert f"{event_type}_events" not in out
|
||||
|
|
@ -104,9 +104,11 @@ def test_event_diff_is_set_based_lateral_not_python_loop() -> None:
|
|||
|
||||
def test_event_diff_emits_price_change_with_diff_percent() -> None:
|
||||
assert "'price_change'" in _EVENT_DIFF_SQL
|
||||
# diff_percent = (new-old)/old*100.
|
||||
# diff_percent = (new-old)/old*100. NULLIF на знаменателе (#2674): выражения
|
||||
# VALUES вычисляются ДО фильтра e.fires, поэтому предикат "p.price_rub <> 0"
|
||||
# больше не защищает само деление.
|
||||
assert "(t.price_rub - p.price_rub)" in _EVENT_DIFF_SQL
|
||||
assert "/ p.price_rub * 100" in _EVENT_DIFF_SQL
|
||||
assert "/ NULLIF(p.price_rub, 0) * 100" in _EVENT_DIFF_SQL
|
||||
# Only when the price actually changed and old is a usable denominator.
|
||||
assert "t.price_rub <> p.price_rub" in _EVENT_DIFF_SQL
|
||||
assert "p.price_rub <> 0" in _EVENT_DIFF_SQL
|
||||
|
|
@ -200,22 +202,30 @@ def test_migration_079_uses_psycopg_safe_sql() -> None:
|
|||
|
||||
|
||||
class _FakeResult:
|
||||
def __init__(self, rowcount: int) -> None:
|
||||
def __init__(self, rowcount: int, rows: list[tuple[str, int]] | None = None) -> None:
|
||||
self.rowcount = rowcount
|
||||
self._rows = rows or []
|
||||
|
||||
def fetchall(self) -> list[tuple[str, int]]:
|
||||
"""Event-diff statement возвращает пары (counter_key, n) — см. #2674."""
|
||||
return self._rows
|
||||
|
||||
|
||||
class _FakeDB:
|
||||
"""Minimal stand-in for a SQLAlchemy Session — records execute() calls, returns rowcounts."""
|
||||
|
||||
def __init__(self, rowcounts: list[int]) -> None:
|
||||
def __init__(
|
||||
self, rowcounts: list[int], event_rows: list[tuple[str, int]] | None = None
|
||||
) -> None:
|
||||
self._rowcounts = list(rowcounts)
|
||||
self._event_rows = event_rows or []
|
||||
self.executed: list[Any] = []
|
||||
self.committed = False
|
||||
self.rolled_back = False
|
||||
|
||||
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||
self.executed.append((stmt, params))
|
||||
return _FakeResult(self._rowcounts.pop(0))
|
||||
return _FakeResult(self._rowcounts.pop(0), self._event_rows)
|
||||
|
||||
def commit(self) -> None:
|
||||
self.committed = True
|
||||
|
|
@ -225,7 +235,12 @@ class _FakeDB:
|
|||
|
||||
|
||||
def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""snapshot_listing_sources maps the two execute() rowcounts to its counters + marks done."""
|
||||
"""snapshot_listing_sources maps the snapshot rowcount + per-type event rows to counters.
|
||||
|
||||
#2674: event-diff больше не отдаёт один rowcount — внешний SELECT над
|
||||
data-modifying CTE возвращает пары (counter_key, n) по типам событий, а
|
||||
counters инициализированы всеми пятью ключами (не сработавший тип = честный 0).
|
||||
"""
|
||||
marked: dict[str, Any] = {}
|
||||
monkeypatch.setattr(
|
||||
snap_mod.runs_mod,
|
||||
|
|
@ -234,11 +249,16 @@ def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
|
|||
)
|
||||
monkeypatch.setattr(snap_mod.runs_mod, "mark_failed", lambda *a, **k: None)
|
||||
|
||||
# rowcounts: SET LOCAL statement_timeout (ignored), snapshot upsert, event-diff insert.
|
||||
db = _FakeDB(rowcounts=[0, 18355, 42])
|
||||
# rowcounts: SET LOCAL statement_timeout (ignored), snapshot upsert, event-diff.
|
||||
db = _FakeDB(rowcounts=[0, 18355, 0], event_rows=[("price_change_events", 42)])
|
||||
out = snap_mod.snapshot_listing_sources(db, run_id=99) # type: ignore[arg-type]
|
||||
|
||||
assert out == {"snapshotted": 18355, "price_change_events": 42}
|
||||
assert out == {
|
||||
"snapshotted": 18355,
|
||||
"price_change_events": 42,
|
||||
"edited_events": 0,
|
||||
"first_seen_events": 0,
|
||||
}
|
||||
assert db.committed is True
|
||||
assert len(db.executed) == 3
|
||||
# First statement sets the per-transaction wall-clock budget (#2607).
|
||||
|
|
@ -249,7 +269,7 @@ def test_counter_logic_with_fake_db(monkeypatch: pytest.MonkeyPatch) -> None:
|
|||
assert params is not None and params["run_id"] == 99
|
||||
# Run finalised via mark_done with the same counters.
|
||||
assert marked["run_id"] == 99
|
||||
assert marked["counters"] == {"snapshotted": 18355, "price_change_events": 42}
|
||||
assert marked["counters"] == out
|
||||
|
||||
|
||||
# ── budget_sec / statement_timeout (#2607) ─────────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -141,9 +141,18 @@ class IMVSuggestion:
|
|||
title: str | None = None
|
||||
address: str | None = None
|
||||
price_rub: int | None = None
|
||||
# rooms/area_m2/floor/total_floors парсятся из title (#2674) — тот же путь, что
|
||||
# у IMVPlacementHistoryItem; колонки в house_suggestions есть с миграции 064.
|
||||
rooms: int | None = None
|
||||
area_m2: float | None = None
|
||||
floor: int | None = None
|
||||
total_floors: int | None = None
|
||||
exposure_days: int | None = None
|
||||
publish_date: date | None = None
|
||||
item_url: str | None = None
|
||||
# image_link — ссылка на фото лота (raw imageLink). Колонка house_suggestions.image_link
|
||||
# существует с 064, но до #2674 значение выбрасывалось парсером.
|
||||
image_link: str | None = None
|
||||
metro_name: str | None = None
|
||||
metro_distance: str | None = None
|
||||
metro_color: str | None = None
|
||||
|
|
@ -206,10 +215,18 @@ def compute_imv_cache_key(
|
|||
|
||||
|
||||
# Паттерн для заголовка вида "2-к. квартира, 42 м², 4/5 эт."
|
||||
# Комнатность НЕОБЯЗАТЕЛЬНА (#2674): 1991 заголовок из 25 055 (7.9%) — «Квартира-студия,
|
||||
# 34,2 м², 9/10 эт.». Площадь и этажность там есть, но обязательная группа комнатности
|
||||
# роняла match целиком и обнуляла ВСЕ ЧЕТЫРЕ поля. Тот же потолок был виден на соседней
|
||||
# таблице (8.8% строк house_placement_history без площади) — обоих писателей чинит один
|
||||
# регексп. Опциональная группа жадная, поэтому «3-к. квартира…» по-прежнему даёт rooms=3.
|
||||
_TITLE_RE = re.compile(
|
||||
r"^(?P<rooms>\d+)-к[.\s].*?(?P<area>[\d,]+)\s*м².*?(?P<floor>\d+)/(?P<total>\d+)\s*эт",
|
||||
r"^(?:(?P<rooms>\d+)-к[.\s])?.*?(?P<area>[\d,]+)\s*м².*?(?P<floor>\d+)/(?P<total>\d+)\s*эт",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
# Студия = 0 комнат — конвенция kit'а (scraper_kit.base.RawLot.rooms «0 = студия»,
|
||||
# providers/yandex/detail.py). Отличаем её от «комнатность неизвестна» (None).
|
||||
_TITLE_STUDIO_RE = re.compile(r"студи[яюей]", re.IGNORECASE)
|
||||
|
||||
|
||||
def _parse_title(title: str | None) -> dict[str, int | float | None]:
|
||||
|
|
@ -226,7 +243,11 @@ def _parse_title(title: str | None) -> dict[str, int | float | None]:
|
|||
if not m:
|
||||
return result
|
||||
try:
|
||||
result["rooms"] = int(m.group("rooms"))
|
||||
rooms = m.group("rooms")
|
||||
if rooms:
|
||||
result["rooms"] = int(rooms)
|
||||
elif _TITLE_STUDIO_RE.search(title):
|
||||
result["rooms"] = 0 # студия
|
||||
result["area_m2"] = float(m.group("area").replace(",", "."))
|
||||
result["floor"] = int(m.group("floor"))
|
||||
result["total_floors"] = int(m.group("total"))
|
||||
|
|
@ -262,20 +283,28 @@ def _parse_placement_item(raw: dict[str, Any]) -> IMVPlacementHistoryItem:
|
|||
|
||||
def _parse_suggestion(raw: dict[str, Any]) -> IMVSuggestion:
|
||||
"""Парсит один элемент из suggestions.items."""
|
||||
title = raw.get("title")
|
||||
parsed = _parse_title(title)
|
||||
metro = raw.get("metro") or {}
|
||||
colors: list[str] = metro.get("colors") or []
|
||||
|
||||
# Фильтруем imageLink из raw_payload
|
||||
# imageLink выносим в отдельное поле (колонка house_suggestions.image_link) и
|
||||
# убираем из raw_payload, чтобы не хранить ссылку дважды.
|
||||
clean_raw = {k: v for k, v in raw.items() if k not in ("imageLink",)}
|
||||
|
||||
return IMVSuggestion(
|
||||
ext_item_id=str(raw["id"]),
|
||||
title=raw.get("title"),
|
||||
title=title,
|
||||
address=raw.get("address"),
|
||||
price_rub=raw.get("price"),
|
||||
rooms=parsed["rooms"], # type: ignore[arg-type]
|
||||
area_m2=parsed["area_m2"], # type: ignore[arg-type]
|
||||
floor=parsed["floor"], # type: ignore[arg-type]
|
||||
total_floors=parsed["total_floors"], # type: ignore[arg-type]
|
||||
exposure_days=raw.get("exposure"),
|
||||
publish_date=_unix_to_date(raw.get("publishDate")),
|
||||
item_url=raw.get("itemLink"),
|
||||
image_link=raw.get("imageLink"),
|
||||
metro_name=metro.get("name"),
|
||||
metro_distance=metro.get("distance"),
|
||||
metro_color=colors[0] if colors else None,
|
||||
|
|
|
|||
|
|
@ -4,11 +4,16 @@
|
|||
PRIMARY KEY (listing_id, snapshot_date) — максимум 1 snapshot в сутки.
|
||||
ON CONFLICT DO UPDATE — берём последний за день (перезаписываем при повторном run'е).
|
||||
|
||||
Используется двумя путями:
|
||||
Используется тремя путями:
|
||||
1. SERP scrape (save_listings в base.py) — записывает цену + позицию в выдаче
|
||||
каждый раз когда listing появляется в поиске.
|
||||
каждый раз когда listing появляется в поиске (status='active').
|
||||
2. Detail backfill (save_detail_enrichment в cian_detail.py) — записывает цену
|
||||
из detail-страницы для listings которые раньше не имели snapshot'а.
|
||||
из 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 хелпер.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue