merge(#3051): main (#3421) в ветку импорта по региону — московская дельта поверх region_code/doc_type
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m12s

#3421 въехал в main параллельно с той же миграцией 288 (deals.doc_type,
параметры region_code/doc_types). Разрешение: 288 — целиком версия main;
наша дельта (FDW-колонки okato/quarter_cad_number/district, выключенный seed
rosreestr_dkp_import_77) переехала в 289. scheduler.py — doc_types из main +
canonical_city-маппинг/raw_payload/per-source чекпоинт. deploy-скрипт —
валидация REGION_CODE и DOC_TYPE (интерполируются в SQL текстом).
This commit is contained in:
bot-backend 2026-09-08 23:19:47 +03:00
commit fcf5887225
7 changed files with 983 additions and 93 deletions

View file

@ -247,7 +247,7 @@ def _dkp_source_for_region(region_code: int) -> str:
66 байт-в-байт прежнее имя ('rosreestr_dkp_import'), под которым годами
писались scrape_runs. Остальные регионы получают суффикс кода тот же
формат, что и у строки scrape_schedules ('rosreestr_dkp_import_77',
seed миграция 288), которую резолвит wildcard 'rosreestr_dkp_import_*'
seed миграция 289), которую резолвит wildcard 'rosreestr_dkp_import_*'
в product_handlers.py. Изоляция чекпоинтов между регионами держится именно
на разных source: _resume_dkp_cursor ищет ПРЕДЫДУЩИЙ прогон с ТЕМ ЖЕ source,
поэтому курсор региона 77 никогда не подхватит last_id региона 66 (и
@ -356,12 +356,14 @@ def import_rosreestr_dkp(
было хардкод region_code=66.
Источник: foreign table gendesign_rosreestr_deals (создана в migration 072,
okato/quarter_cad_number/district добавлены миграцией 288).
okato/quarter_cad_number/district добавлены миграцией 289).
SERVER gendesign_remote настроен в 060_postgres_fdw_extension.sql.
USER MAPPING создаётся при startup в core/fdw.py (tradein_fdw_reader).
Область покрытия: region_code из params (default 66 вся Свердловская область,
не только Екатеринбург). Неизвестный код региона (нет в REGIONS) ValueError,
Область покрытия задаётся параметром, а не литералом (#3051 п.6): region_code
приходит из params, дефолт 66 = вся Свердловская область (не только Екатеринбург
прежний ILIKE-фильтр по подстроке города снят, unlocks +47183 сделок вне ЕКБ уже
сидящих в source foreign table). Неизвестный код региона (нет в REGIONS) ValueError,
прогон падает явно, а не молча импортирует мусор с чужим region_code.
region_code=66 (регион БЕЗ canonical_city в реестре) поведение байт-в-байт
@ -377,6 +379,12 @@ def import_rosreestr_dkp(
:canonical_city (CASE WHEN ... IS NOT NULL), а не отдельный Python if/else на
конкретный код региона.
Типы документов тоже параметр (#3051 п.3): doc_types, дефолт ['ДКП'] = прежнее
поведение (только вторичка #549 / Fix_Rosreestr_Dkp_Filter_May24). Для Москвы
ДДУ идут по ценам котлована и медиану развалят, поэтому смешивать их с ДКП можно
только осознанно и с колонкой deals.doc_type (миграция 288), которая теперь
заполняется на импорте.
Фильтры (совпадают с import-rosreestr.sh + Fix_Rosreestr_Dkp_Filter_May24):
- region_code = :region_code (параметризовано, было хардкод 66)
- city IS NOT NULL AND trim(city) != '' ТОЛЬКО если у региона нет canonical_city
@ -384,12 +392,16 @@ def import_rosreestr_dkp(
- area BETWEEN 18 AND 200
- deal_price BETWEEN 1000000 AND 100000000
- street IS NOT NULL AND trim(street) != ''
- doc_type = 'ДКП' (только вторичка #549 / Fix_Rosreestr_Dkp_Filter_May24)
- doc_type = ANY(:doc_types) (param, default ['ДКП'])
- period_start_date >= since (default '2024-01-01')
dedup_hash: 'ros:dkp:' || id плоский натуральный ключ (инъективный, без коллизий,
human-readable). До #576 здесь был md5('ros:dkp:' || id); миграция 077 конвертировала
существующие строки. source_id хранит исходный rosreestr id (дедуп переустанавливаем).
Префикс ':dkp:' НАМЕРЕННО оставлен неизменным после параметризации doc_types: id
уникален в источнике сам по себе, независимо от типа документа, поэтому ключ и без
того не коллизирует; а вот смена формы ключа осиротила бы все уже загруженные строки
(их пришлось бы конвертировать ещё одной миграцией ровно то, что делала 077).
Rooms: выводятся из площади (Росреестр не отдаёт кол-во комнат).
Batch-процессинг: читаем из FDW батчами по batch_size через cursor-based пагинацию
@ -411,8 +423,11 @@ def import_rosreestr_dkp(
"""
since: str = str(params.get("since", "2024-01-01"))
batch_size: int = int(params.get("batch_size", 2000))
# #3051 п.6: регион — параметр, дефолт 66 сохраняет текущее прод-поведение
# (расписание получает явный region_code в миграции 288).
region_code: int = int(params.get("region_code", 66))
# #3051 п.3: типы документов — параметр, дефолт ['ДКП'] = прежний литерал.
doc_types: list[str] = [str(t) for t in params.get("doc_types") or ["ДКП"]]
region = REGIONS.get(region_code)
if region is None:
raise ValueError(
@ -566,7 +581,7 @@ def import_rosreestr_dkp(
AND area BETWEEN 18 AND 200
AND deal_price BETWEEN 1000000 AND 100000000
AND street IS NOT NULL AND trim(street) <> ''
AND doc_type = 'ДКП'
AND doc_type = ANY(CAST(:doc_types AS text[]))
AND period_start_date >= CAST(:since AS date)
AND id > CAST(:last_id AS bigint)
ORDER BY id
@ -578,6 +593,7 @@ def import_rosreestr_dkp(
"batch_size": batch_size,
"region_code": region_code,
"canonical_city": region.canonical_city,
"doc_types": doc_types,
},
)
.mappings()

View file

@ -1,87 +1,67 @@
-- 288_deals_doc_type.sql
-- deals.doc_type (#3051 п.3) + foreign table gendesign_rosreestr_deals: okato/
-- quarter_cad_number/district + disabled seed schedule rosreestr_dkp_import_77.
-- deals.doc_type — тип документа сделки Росреестра (#3051 п.3, п.6).
--
-- Dependencies: 002_core_tables.sql (deals), 072_scrape_schedules_seed_cian_rosreestr.sql
-- (gendesign_rosreestr_deals, scrape_schedules).
-- Dependencies: 002_core_tables.sql (deals), 015_scrape_runs.sql + 072_scrape_schedules_seed_cian_rosreestr.sql
-- (scrape_schedules, строка source='rosreestr_dkp_import')
-- Apply after: 287_proxy_run_attribution.sql
--
-- WHY:
-- Трек 2 подготовки Mera к Москве — импорт сделок Росреестра параметризуется по
-- региону (66 Свердловская обл. / 77 Москва, code-часть в scheduler.py). deals
-- до сих пор не различал ДКП (вторичка) от ДДУ (застройщик) в самой строке —
-- различие жило только в WHERE-фильтре импортёра. Явная колонка нужна для
-- будущего ДДУ-импорта (#3051 п.3) и для аналитики, которая иначе не может
-- отличить типы сделок в одной таблице.
-- ЧАСТЬ A — WHY:
-- Импорт Росреестра до сих пор ронял тип документа на пол: фильтр doc_type = 'ДКП'
-- стоял литералом в WHERE, а в deals не приезжало ничего. Пока скоуп был один
-- (Свердловская обл., только вторичка) это было безобидно — все строки источника
-- 'rosreestr' по построению ДКП. С расширением на Москву безобидность кончается:
-- в источнике за 2024 по region_code=77 лежит 30 627 ДДУ с медианой 112 743 ₽/м²
-- против 107 005 ДКП с медианой 256 250 ₽/м² — это цены котлована, и смешать их
-- в одной таблице без различимого признака значит развалить любую оценку.
-- Колонка нужна ДО того, как импорт начнёт тянуть больше одного типа.
--
-- Foreign table расширена тремя колонками источника (okato, quarter_cad_number,
-- district) — они нужны import_rosreestr_dkp для raw_payload по регионам, где
-- city источника не используется как есть (Москва: city = муниципальный
-- округ/поселение, не город). Существование колонок в public.rosreestr_deals
-- на gendesign-стороне проверено live (2026-09-08, prod psql).
-- ЧАСТЬ A — WHAT:
-- doc_type text NULLable — источник (foreign table gendesign_rosreestr_deals.doc_type)
-- тоже text и NULL допускает; строгий NOT NULL сломал бы не-росреестровые источники
-- (etazhi / domklik_history), у которых понятия «тип документа» нет вовсе.
-- Индекс НЕ добавляем: селективность низкая (2-3 значения), а все живые выборки
-- по deals идут по region_code/deal_date/geom — doc_type там в лучшем случае
-- довесок к уже отобранному диапазону. Появится запрос, который реально режет
-- по doc_type на большом наборе — заведём частичный индекс тогда, по EXPLAIN.
--
-- Seed-строка rosreestr_dkp_import_77 — ВЫКЛЮЧЕНА (enabled=false): миграция
-- только заводит расписание, включение и первый прогон по Москве — отдельное
-- решение main-сессии после ревью кода-части.
-- ЧАСТЬ B — бэкфилл:
-- Всё, что лежит в deals с source='rosreestr', прошло через WHERE doc_type = 'ДКП'
-- (и в scheduler.import_rosreestr_dkp, и в deploy/import-rosreestr.sh, и в
-- backfill-миграции 077) — других типов там физически быть не может. Поэтому
-- проставить 'ДКП' задним числом корректно, а не эвристика.
-- Остальные источники остаются NULL осознанно.
--
-- ИДЕМПОТЕНТНОСТЬ: ADD COLUMN IF NOT EXISTS × 4, COMMENT ON COLUMN (безусловны,
-- но перезаписывают тот же текст), ON CONFLICT (source) DO NOTHING для seed.
-- Backfill (doc_type='ДКП' WHERE source='rosreestr' AND doc_type IS NULL) —
-- повторный прогон no-op (второй раз IS NULL уже не матчит ни одну строку).
-- ЧАСТЬ C — расписание:
-- scrape_schedules.default_params для 'rosreestr_dkp_import' получает явный
-- region_code=66. Раньше регион был неявным дефолтом в коде; после параметризации
-- (params.get('region_code', 66)) неявность становится ловушкой — прод-строка должна
-- сама говорить, какой регион она тянет. doc_types в default_params НЕ пишем:
-- дефолт ['ДКП'] в коде и есть текущее поведение, а запись его в расписание
-- создала бы второе место, где надо не забыть поменять.
--
-- Идемпотентна, безопасна к повторному запуску.
BEGIN;
-- #2752: блокирующий DDL не должен вставать в очередь за чужой сессией и уводить
-- за собой запросы приложения — лучше упасть по таймауту и повторить деплой.
SET LOCAL lock_timeout = '5s';
-- ── deals.doc_type ────────────────────────────────────────────────────────────
-- A: колонка
ALTER TABLE deals ADD COLUMN IF NOT EXISTS doc_type text;
ALTER TABLE deals
ADD COLUMN IF NOT EXISTS doc_type text;
COMMENT ON COLUMN deals.doc_type IS
'Тип документа сделки Росреестра (ДКП / ДДУ). NULL для источников без этого понятия.';
-- Backfill: весь текущий rosreestr-импорт в deals уже отфильтрован по doc_type='ДКП'
-- на стороне import_rosreestr_dkp (WHERE doc_type = 'ДКП') — строки, попавшие в deals
-- ДО этой миграции, были все ДКП, ДДУ-импорта ещё нет ни одной строки.
-- B: бэкфилл — все rosreestr-строки прошли фильтр ДКП на импорте
UPDATE deals
SET doc_type = 'ДКП'
WHERE source = 'rosreestr'
AND doc_type IS NULL;
COMMENT ON COLUMN deals.doc_type IS
'Тип документа сделки Росреестра: ДКП (договор купли-продажи, вторичка) или '
'ДДУ (договор долевого участия, застройщик) — #3051 п.3. NULL для строк не из '
'rosreestr (etazhi/domklik_history) — у них своя типизация сделки, либо не '
'применимо. Заполняется import_rosreestr_dkp из значения источника (WHERE '
'уже фильтрует doc_type=''ДКП'', колонка носит его явно, а не только в фильтре).';
-- ── gendesign_rosreestr_deals: колонки под региональный raw_payload (#3051) ────
-- ALTER FOREIGN TABLE ADD COLUMN — только локальные метаданные (не трогает
-- реальную remote-таблицу), безопасно как ALTER TABLE ADD COLUMN без DEFAULT.
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS okato text;
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS quarter_cad_number text;
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS district text;
-- ── Seed: rosreestr_dkp_import_77 (выключено) ──────────────────────────────────
INSERT INTO scrape_schedules (
source,
enabled,
window_start_hour,
window_end_hour,
default_params
)
VALUES (
'rosreestr_dkp_import_77',
false,
4,
6,
'{"region_code": 77, "since": "2024-01-01", "batch_size": 2000}'::jsonb
)
ON CONFLICT (source) DO NOTHING;
-- C: явный регион в расписании импорта вместо неявного дефолта в коде
UPDATE scrape_schedules
SET default_params = default_params || '{"region_code": 66}'::jsonb
WHERE source = 'rosreestr_dkp_import';
COMMIT;

View file

@ -0,0 +1,60 @@
-- 289_rosreestr_fdw_msk_columns_seed77.sql
-- gendesign_rosreestr_deals: колонки okato/quarter_cad_number/district + disabled
-- seed-строка scrape_schedules для региона 77 (Москва) — #3051 п.3.
--
-- Dependencies: 072_scrape_schedules_seed_cian_rosreestr.sql (gendesign_rosreestr_deals,
-- scrape_schedules).
-- Apply after: 288_deals_doc_type.sql
--
-- WHY:
-- Трек 2 подготовки Mera к Москве — импорт сделок Росреестра параметризуется по
-- региону (66 Свердловская обл. / 77 Москва, code-часть в scheduler.py). Для
-- региона С canonical_city (Москва) исходные city/okato/quarter_cad_number/district
-- уходят в deals.raw_payload (jsonb), т.к. city источника там — муниципальный
-- округ/поселение, не город, и обычный city/address его не покрывают. Foreign
-- table расширена тремя колонками источника; существование их в
-- public.rosreestr_deals на gendesign-стороне проверено live (2026-09-08, prod psql).
--
-- Seed-строка rosreestr_dkp_import_77 — ВЫКЛЮЧЕНА (enabled=false): миграция
-- только заводит расписание, включение и первый прогон по Москве — отдельное
-- решение main-сессии после ревью кода-части.
--
-- ИДЕМПОТЕНТНОСТЬ: ADD COLUMN IF NOT EXISTS × 3, ON CONFLICT (source) DO NOTHING
-- для seed — повторный прогон no-op.
BEGIN;
SET LOCAL lock_timeout = '5s';
-- ── gendesign_rosreestr_deals: колонки под региональный raw_payload (#3051) ────
-- ALTER FOREIGN TABLE ADD COLUMN — только локальные метаданные (не трогает
-- реальную remote-таблицу), безопасно как ALTER TABLE ADD COLUMN без DEFAULT.
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS okato text;
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS quarter_cad_number text;
ALTER FOREIGN TABLE gendesign_rosreestr_deals
ADD COLUMN IF NOT EXISTS district text;
-- ── Seed: rosreestr_dkp_import_77 (выключено) ──────────────────────────────────
INSERT INTO scrape_schedules (
source,
enabled,
window_start_hour,
window_end_hour,
default_params
)
VALUES (
'rosreestr_dkp_import_77',
false,
4,
6,
'{"region_code": 77, "since": "2024-01-01", "batch_size": 2000}'::jsonb
)
ON CONFLICT (source) DO NOTHING;
COMMIT;

View file

@ -151,17 +151,14 @@ def test_migration_077_converts_md5_to_plain_key() -> None:
def test_migration_077_shared_filter_matches_live_import() -> None:
"""Дедуп-релевантные фильтры src CTE миграции 077 совпадают с живым импортом.
Исключения:
Два исключения:
- city ILIKE (см. test_live_import_dropped_ekb_city_filter ниже): 077 backfill'ил
ЕКБ-строки под EKB-only scope того времени; живой импорт расширен на всю
Свердловскую область (region_code=66, все города), поэтому city-фильтр из
живого импорта СНЯТ намеренно.
- region_code (#3051 п.3): 077 — историческая миграция (уже применена на прод,
хардкод region_code=66 трогать нельзя и не нужно она навсегда про 66). Живой
импорт с #3051 параметризован по региону (region_code = CAST(:region_code AS
int)) литерала 66 в его SQL больше нет, см.
test_live_import_region_code_is_bind_param ниже. Остальные клозы обязаны
совпадать байт-в-байт, иначе 077 конвертировал бы не тот набор строк.
Свердловскую область, поэтому city-фильтр из живого импорта СНЯТ намеренно;
- region_code / doc_type: в живом импорте это ПАРАМЕТРЫ (#3051), их совпадение
с 077 проверяется дефолтами params.get(..., 66) / ["ДКП"] в самом импорте.
Остальные клозы обязаны совпадать байт-в-байт, иначе 077 конвертировал бы не тот
набор строк.
"""
sql = _MIGRATION_077.read_text("utf-8")
assert "region_code = 66" in sql, "migration 077 must keep its historical literal"
@ -170,10 +167,12 @@ def test_migration_077_shared_filter_matches_live_import() -> None:
"area BETWEEN 18 AND 200",
"deal_price BETWEEN 1000000 AND 100000000",
"street IS NOT NULL AND trim(street) <> ''",
"doc_type = 'ДКП'",
):
assert clause in sql, f"missing filter clause in migration: {clause!r}"
assert clause in _IMPORT_SRC, f"missing filter clause in import: {clause!r}"
# Исторические литералы 077 остаются на месте (миграция применена, её не правят).
assert "region_code = 66" in sql
assert "doc_type = 'ДКП'" in sql
def test_live_import_region_code_is_bind_param() -> None:

View file

@ -2,11 +2,16 @@
# Импорт реальных сделок Росреестра из gendesign-БД в tradein deals.
#
# Источник: gendesign-postgres-1 / rosreestr_deals (6.8М строк, партиц.).
# Берём квартиры региона REGION_CODE (realestate_type_code=002001003000) за
# последние ~18 мес. rooms выводим из площади (Росреестр не отдаёт кол-во комнат).
# Берём квартиры (realestate_type_code=002001003000) региона REGION_CODE — по умолчанию
# 66 = вся Свердловская область, все города (не только ЕКБ: city-фильтр снят давно,
# шапка про «ЕКБ квартиры» была неправдой).
# rooms выводим из площади (Росреестр не отдаёт кол-во комнат).
# Координаты — NULL, проставляются отдельно (geocode по street).
#
# Только ДКП (вторичка) — ДДУ застройщиков скёюят median, отделены сознательно. См. PR-A.
# DOC_TYPE по умолчанию ДКП (вторичка): ДДУ идут по ценам котлована и скёюят median,
# отделены сознательно (см. PR-A). С #3051 это параметр, а не литерал — для Москвы
# (REGION_CODE=77) типы разделяются колонкой deals.doc_type (миграция 288), которую
# скрипт теперь заполняет.
#
# #3051 п.3: REGION_CODE параметризован (default 66 — Свердловская обл., byte-for-byte
# прежнее поведение). Этот bash-путь — НЕ region-generic: для region_code=77 (Москва) он
@ -19,6 +24,7 @@
# паритет фильтров (area/price/doc_type) обязательно, паритет city-override — нет.
#
# Запуск на прод-хосте: ./import-rosreestr.sh
# REGION_CODE=77 DOC_TYPE='ДДУ' ./import-rosreestr.sh # Москва, первичка
# Повторяемо: дедуп по dedup_hash, новый запуск подтянет свежие кварталы.
set -euo pipefail
@ -28,11 +34,17 @@ DST_PG="${DST_PG:-tradein-postgres}"
SRC_DB="${SRC_DB:-gendesign}"
SRC_USER="${SRC_USER:-gendesign}"
SINCE="${SINCE:-2024-01-01}"
# #3051: регион и тип документа — параметры со старыми дефолтами (поведение не меняется).
REGION_CODE="${REGION_CODE:-66}"
# Значение подставляется в SQL текстом — допускаем только целое число.
DOC_TYPE="${DOC_TYPE:-ДКП}"
# Оба значения интерполируются в SQL текстом (не bind-параметром) — валидация здесь
# и есть единственная защита от SQL-инъекции. Паттерн — в переменной (не инлайн в
# [[ =~ ]]): голая кавычка в regex-операнде ломает bash-парсинг команды.
DOC_TYPE_RE="^[^;']+\$"
[[ "$REGION_CODE" =~ ^[0-9]+$ ]] || { echo "REGION_CODE должен быть целым числом, получено: '$REGION_CODE'" >&2; exit 1; }
[[ "$DOC_TYPE" =~ $DOC_TYPE_RE ]] || { echo "DOC_TYPE не должен содержать ' или ;, получено: '$DOC_TYPE'" >&2; exit 1; }
echo "[$(date -u +%H:%M:%S)] import-rosreestr: регион $REGION_CODE, квартиры с $SINCE"
echo "[$(date -u +%H:%M:%S)] import-rosreestr: region=$REGION_CODE doc_type=$DOC_TYPE с $SINCE"
# 1. Staging-таблица в tradein.
docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
@ -42,7 +54,8 @@ docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
CREATE TABLE deals_ros_staging (
dedup_hash text PRIMARY KEY, source_id text, address text, region_code int, city text,
rooms int, area_m2 numeric,
floor int, year_built int, price_rub bigint, price_per_m2 int, deal_date date
floor int, year_built int, price_rub bigint, price_per_m2 int, deal_date date,
doc_type text
);
"
@ -64,7 +77,8 @@ docker exec "$SRC_PG" psql -U "$SRC_USER" -d "$SRC_DB" -v ON_ERROR_STOP=on -c "
year_build AS year_built,
round(deal_price)::bigint AS price_rub,
round(price_per_sqm)::int AS price_per_m2,
period_start_date AS deal_date
period_start_date AS deal_date,
doc_type AS doc_type
FROM rosreestr_deals
WHERE region_code = $REGION_CODE
AND city IS NOT NULL AND trim(city) <> ''
@ -72,7 +86,7 @@ docker exec "$SRC_PG" psql -U "$SRC_USER" -d "$SRC_DB" -v ON_ERROR_STOP=on -c "
AND area BETWEEN 18 AND 200
AND deal_price BETWEEN 1000000 AND 100000000
AND street IS NOT NULL AND trim(street) <> ''
AND doc_type = 'ДКП'
AND doc_type = '$DOC_TYPE'
AND period_start_date >= '$SINCE'
) TO STDOUT WITH CSV
" | docker exec -i "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
@ -88,10 +102,10 @@ docker exec "$DST_PG" psql -U tradein -d tradein -v ON_ERROR_STOP=on -c "
INSERT INTO deals (
source, dedup_hash, source_id, address, region_code, city, rooms, area_m2, floor,
year_built, price_rub, price_per_m2, deal_date
year_built, price_rub, price_per_m2, deal_date, doc_type
)
SELECT 'rosreestr', dedup_hash, source_id, address, region_code, city, rooms, area_m2, floor,
year_built, price_rub, price_per_m2, deal_date
year_built, price_rub, price_per_m2, deal_date, doc_type
FROM deals_ros_staging
ON CONFLICT (dedup_hash) DO NOTHING;

View file

@ -0,0 +1,147 @@
# Локальный сбор SERP Авито по Москве и МО (эпик #2989, трек 1)
`collect.py` — ручной скрипт **с машины владельца**. Собирает карточки выдачи Авито
(вторичка, Москва + МО) и заливает их в прод-схему `msk_raw`.
Прод-скрейпер, его расписания, прокси-пул и сайдкар **не задействованы вообще**.
Браузер — уже открытый Chrome владельца (подключение по CDP), парсер — импорт из
`packages/scraper-kit`, заливка — поток в `psql` через `ssh`.
## Предусловия
1. **Chrome владельца** запущен с залогиненным техническим аккаунтом Авито и
remote debugging:
```powershell
& "C:\Program Files\Google\Chrome\Application\chrome.exe" --remote-debugging-port=9222
```
Проверка: `curl http://localhost:9222/json/version` отдаёт JSON с `Browser: Chrome/...`.
Другой адрес — переменная `AVITO_CDP`.
Скрипт **не поднимает свой профиль** (`launch_persistent_context` не используется):
он подключается к существующему браузеру, берёт `browser.contexts[0]`, открывает
**свою** вкладку и в конце закрывает **только её**. Чужие вкладки, контекст и сам
браузер не трогаются — это рабочий Chrome владельца.
2. **ssh-доступ на прод** (`ssh selectel` без пароля) — для заливки в БД.
Прямого подключения к прод-Postgres с локалки нет: туннель `:35432` ведёт в мёртвую
копию на Beget.
3. **Playwright в текущем интерпретаторе**:
```powershell
pip install playwright
python -m playwright install chromium
```
В uv-зависимости проекта playwright **не добавлять**: прод-сайдкар пинит
playwright 1.60 под camoufox, и подъём версии сломает его.
Больше ничего ставить не нужно: только stdlib + playwright + импорт `scraper_kit`
(путь `packages/scraper-kit/src` скрипт добавляет в `sys.path` сам, от `__file__`).
## Запуск (PowerShell)
```powershell
cd D:\prjct\gendesign\tradein-mvp\scripts\local-avito-msk
# 0) сухой прогон: ничего не шлём на прод, карточки пишем в runs\cards-<batch>.csv
python .\collect.py --dry-run --measure 5
# 1) обязательный первый прогон — замер (дефолт, 100 загрузок страниц)
python .\collect.py
# 2) полный проход — только явно
python .\collect.py --full --batch-id msk-serp-20260908
# 3) продолжить прерванный прогон по сохранённому плану коридоров
python .\collect.py --full --resume --batch-id msk-serp-20260908
```
Без аргументов скрипт работает в режиме `--measure 100` и полный проход **не начинает**.
Ключи: `--delay` (пауза между загрузками, дефолт 8.0 с ±20 % джиттера — сознательно
совпадает с прод-расписаниями `request_delay_sec` 710 с), `--batch-size` (карточек в
одной заливке, дефолт 1000), `--target-count` (целевой размер коридора, дефолт 1500),
`--base-url`, `--batch-id`, `--out-dir`, `--ssh-host/--container/--db-user/--db-name`.
## Как режется выдача
Потолок пагинации Авито — 30 страниц по 60 = **1800 объявлений на запрос**. Любой
запрос с `count > 1800` целиком не добирается, поэтому строится план ценовых коридоров:
* читаем счётчик «N объявлений» со страницы 1 (`page-title/count`);
* `count > --target-count` → делим коридор пополам **по геометрической середине**
(`sqrt(lo*hi)`): цены логнормальны, арифметическая середина диапазона 1 млн … 100 млн
даёт вырожденно-пустую верхнюю половину;
* верхняя граница открытого коридора подбирается удвоением от 8 млн ₽;
* предохранители: глубина рекурсии ≤ 12 и минимальная ширина коридора (отношение
границ ≤ 1.05). Если коридор уже узкий, а `count` всё ещё > 1800 — он помечается
`truncated: true` в плане, а в лог и в `msk_raw.batches.notes` пишется, сколько
объявлений заведомо не добрано;
* гео-параметры `radius`/`geoCoords` не используются: сервером они не применяются (#3043).
План лежит в `runs/plan-<batch_id>.json` и обновляется после каждой страницы — отсюда
работает `--resume`.
## Стоп на первом признаке блока
Проверки в фиксированном порядке, первое срабатывание = немедленный стоп
(никаких ретраев и никакого «продолжим со следующего коридора»):
| # | Признак | Причина в `notes` |
|---|---|---|
| 1 | HTTP 403 / 439 | `platform` |
| 2 | HTTP 429 | `ratelimit` |
| 3 | `_is_firewall_page(html)` — «доступ ограничен», «проблема с ip», `firewall-container` | `firewall` |
| 4 | `startpow` / «доступ ограничен: проверка безопасности» в первых 4 КБ | `challenge` |
| 5 | 0 карточек при ненулевом счётчике (DOM-drift или тихий блок) | `empty_page` |
При стопе: недоотправленный батч дозаливается, у батча проставляются `finished_at` и
`notes`, скрипт выходит с кодом **2**.
**Что делать при стопе.** Не перезапускать сразу и не крутить ретраи. `platform` /
`firewall` / `challenge` — техаккаунт или IP помечены: пауза на несколько часов,
проверить вручную в браузере, что выдача открывается и аккаунт жив, при повторе
увеличить `--delay`. `ratelimit` — темп слишком высокий: `--delay 15` и выше.
`empty_page` — сначала посмотреть сохранённую страницу в браузере: если выдача
рисуется, значит уехал DOM и чинить надо парсер в `scraper-kit`, а не скрипт.
После разбора — `--resume` с тем же `--batch-id`, уже собранное не потеряется.
## Куда пишем
Схема `msk_raw` на проде (создана заранее):
* `msk_raw.batches(batch_id PK, kind, query, started_at, finished_at, rows_sent, rows_new, notes, uploaded_at)`
* `msk_raw.avito_cards(id, source_id, observed_at, batch_id → batches, kind, url, price, payload, UNIQUE(source_id,batch_id,kind))`
* `msk_raw.avito_latest` — вью `DISTINCT ON (source_id) … ORDER BY source_id, observed_at DESC`
Форма заливки: поток в
`ssh <host> "docker exec -i tradein-postgres psql -U tradein -d tradein -v ON_ERROR_STOP=1 -f -"`.
Внутри одной транзакции: `INSERT` батча (`ON CONFLICT DO NOTHING` — строка обязана
существовать до карточек из-за FK), `CREATE TEMP TABLE _stg (LIKE msk_raw.avito_cards
INCLUDING DEFAULTS)`, `\copy _stg (...) FROM STDIN WITH (FORMAT csv)`, затем
`INSERT … SELECT` в `avito_cards` с `ON CONFLICT (source_id,batch_id,kind) DO NOTHING`
и `UPDATE batches SET rows_sent/rows_new` (`rows_new` = разница `count(*)` по batch_id
до и после вставки).
Одиночный `psql -c` через ssh ломается на квотинге скобок и кавычек — поэтому только поток.
`payload` = `lot.model_dump(mode="json")`. `source_id` в БД `bigint`, а у `ScrapedLot`
строка: приводится к `int`, нечисловые пропускаются со счётчиком (он попадает в
`notes` как `skipped_non_numeric=N`).
## Как проверить залитое
```powershell
ssh selectel "docker exec -i tradein-postgres psql -U tradein -d tradein -c \"SELECT count(*) FROM msk_raw.avito_latest;\""
ssh selectel "docker exec -i tradein-postgres psql -U tradein -d tradein -c \"SELECT batch_id, rows_sent, rows_new, started_at, finished_at, notes FROM msk_raw.batches ORDER BY uploaded_at DESC LIMIT 5;\""
```
## Про фильтр городов
Парсер конструируется как `AvitoScraper(SimpleNamespace(avito_serp_ekb_only=False),
target_city_slug="moskva")`. `avito_serp_ekb_only=False` **обязателен**: с `True`
`_parse_html` выбрасывает всё, у чего в URL нет `/ekaterinburg/`, включая все
подмосковные слаги — из московской выдачи не осталось бы ничего.

View file

@ -0,0 +1,674 @@
#!/usr/bin/env python3
"""Локальный ручной сборщик SERP Авито по Москве и МО (эпик #2989, трек 1).
Запускается ВРУЧНУЮ с машины владельца. Прод-скрейпер, его расписания и
прокси-пул не задействованы вообще: браузер уже открытый Chrome владельца
(подключение по CDP), парсер импорт из scraper-kit, заливка поток в psql
через ssh. Скрипт ничего не устанавливает и своего профиля не поднимает.
Дефолтный режим --measure 100 (замер): полный проход только по явному --full.
"""
from __future__ import annotations
import argparse
import asyncio
import csv
import io
import json
import math
import os
import random
import re
import subprocess
import sys
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from types import SimpleNamespace
from typing import Any, Iterable
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
# --- импорт парсера из scraper-kit без установки backend -------------------
_KIT_SRC = Path(__file__).resolve().parents[2] / "packages" / "scraper-kit" / "src"
if str(_KIT_SRC) not in sys.path:
sys.path.insert(0, str(_KIT_SRC))
from scraper_kit.providers.avito.serp import ( # noqa: E402
AvitoScraper,
_is_firewall_page,
)
# Вкладка, открытая у владельца: вторичка, Москва + МО.
DEFAULT_BASE_URL = (
"https://www.avito.ru/moskva_i_mo/kvartiry/prodam/vtorichka-ASgBAgICAkSSA8YQ5geMUg"
"?f=ASgBAgICA0SSA8YQ5geMUuDC3c0CgOkP"
)
PAGE_SIZE = 60 # карточек на странице выдачи
MAX_PAGES = 30 # потолок пагинации Авито → 30*60 = 1800 на один запрос
HARD_CAP = PAGE_SIZE * MAX_PAGES
PRICE_FLOOR = 500_000 # нижняя граница осмысленного коридора, ₽
PRICE_PROBE_START = 8_000_000 # старт удвоения при поиске верхней границы
PRICE_CEIL = 2_000_000_000
MIN_WIDTH_RATIO = 1.05 # уже этого коридор не делим (геометрическая ширина)
MAX_DEPTH = 12
_BATCH_ID_RE = re.compile(r"^[A-Za-z0-9._-]+$")
_POW_MARKERS = ("startpow", "доступ ограничен: проверка безопасности")
class Blocked(Exception):
"""Первый признак блока. Ретраев нет — только немедленный стоп."""
def __init__(self, reason: str, detail: str = "") -> None:
super().__init__(f"{reason}: {detail}" if detail else reason)
self.reason = reason
self.detail = detail
class BudgetExhausted(Exception):
"""Потолок --measure выбран: штатный выход, не ошибка."""
# --- план коридоров --------------------------------------------------------
@dataclass
class Corridor:
lo: int | None
hi: int | None
count: int | None = None
truncated: bool = False
pages_done: int = 0
status: str = "pending" # pending | done
missed: int = 0 # заведомо недобрано (count - HARD_CAP), если truncated
def label(self) -> str:
lo = "-" if self.lo is None else f"{self.lo:_}"
hi = "-" if self.hi is None else f"{self.hi:_}"
return f"[{lo} .. {hi}]"
def to_json(self) -> dict[str, Any]:
return {
"lo": self.lo, "hi": self.hi, "count": self.count,
"truncated": self.truncated, "pages_done": self.pages_done,
"status": self.status, "missed": self.missed,
}
@staticmethod
def from_json(d: dict[str, Any]) -> "Corridor":
return Corridor(
lo=d.get("lo"), hi=d.get("hi"), count=d.get("count"),
truncated=bool(d.get("truncated")),
pages_done=int(d.get("pages_done") or 0),
status=d.get("status") or "pending", missed=int(d.get("missed") or 0),
)
def planned_pages(self) -> int:
if not self.count:
return 1
return max(1, min(MAX_PAGES, math.ceil(self.count / PAGE_SIZE)))
@dataclass
class Plan:
base_url: str
target: int
batch_id: str
corridors: list[Corridor] = field(default_factory=list)
created_at: str = ""
def save(self, path: Path) -> None:
path.write_text(
json.dumps(
{
"version": 1, "base_url": self.base_url, "target": self.target,
"batch_id": self.batch_id, "created_at": self.created_at,
"corridors": [c.to_json() for c in self.corridors],
},
ensure_ascii=False, indent=1,
),
encoding="utf-8",
)
@staticmethod
def load(path: Path) -> "Plan":
d = json.loads(path.read_text(encoding="utf-8"))
return Plan(
base_url=d["base_url"], target=int(d["target"]), batch_id=d["batch_id"],
created_at=d.get("created_at", ""),
corridors=[Corridor.from_json(c) for c in d.get("corridors", [])],
)
def build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str:
"""URL коридора: pmin/pmax + пагинация. Гео-параметры не добавляем — #3043."""
parts = urlsplit(base_url)
q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True)
if k not in {"p", "pmin", "pmax"}]
if lo is not None:
q.append(("pmin", str(int(lo))))
if hi is not None:
q.append(("pmax", str(int(hi))))
if page > 1:
q.append(("p", str(page)))
return urlunsplit(
(parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment)
)
def geometric_mid(lo: int | None, hi: int) -> int:
"""Геометрическая середина коридора.
Цены логнормальны: арифметическая середина 1 млн..100 млн (50 млн)
отрезает вырожденно-пустую верхнюю половину. sqrt(lo*hi) делит выборку
заметно ровнее.
"""
low = max(int(lo or PRICE_FLOOR), 1)
mid = int(math.sqrt(low * float(hi)))
return max(low + 1, min(hi - 1, mid))
def width_ratio(lo: int | None, hi: int | None) -> float:
if hi is None:
return float("inf")
return float(hi) / max(float(lo or PRICE_FLOOR), 1.0)
# --- загрузка страницы -----------------------------------------------------
def _guard(html: str, status: int | None) -> None:
"""Порядок проверок фиксирован заданием; первое срабатывание = стоп."""
if status in (403, 439):
raise Blocked("platform", f"HTTP {status}")
if status == 429:
raise Blocked("ratelimit", "HTTP 429")
if _is_firewall_page(html):
raise Blocked("firewall", "firewall-страница на HTTP 200")
head = html[:4096].lower()
if any(m in head for m in _POW_MARKERS):
raise Blocked("challenge", "PoW / проверка безопасности")
class Loader:
"""Одна СВОЯ вкладка в уже открытом Chrome владельца (CDP).
Ни браузер, ни контекст, ни чужие вкладки не закрываются и не трогаются:
это рабочий Chrome с залогиненным техаккаунтом.
"""
def __init__(self, delay: float, page_budget: int | None) -> None:
self._delay = delay
self._budget = page_budget
self.loads = 0
self._page: Any = None
self._pw: Any = None
self._browser: Any = None
self._last_load = 0.0
async def __aenter__(self) -> "Loader":
from playwright.async_api import async_playwright
endpoint = os.environ.get("AVITO_CDP", "http://localhost:9222")
self._pw = await async_playwright().start()
try:
self._browser = await self._pw.chromium.connect_over_cdp(endpoint)
except Exception as exc: # noqa: BLE001 — подсказка важнее типа
await self._pw.stop()
raise SystemExit(
f"Не удалось подключиться по CDP к {endpoint}: {exc}\n"
"Запусти Chrome с залогиненным техаккаунтом Авито и ключом "
"--remote-debugging-port=9222, либо укажи адрес в AVITO_CDP."
) from exc
if not self._browser.contexts:
await self._pw.stop()
raise SystemExit(
"В подключённом Chrome нет ни одного контекста. Открой обычное окно "
"Chrome, запущенное с --remote-debugging-port=9222."
)
ctx = self._browser.contexts[0]
self._page = await ctx.new_page()
return self
async def __aexit__(self, *exc: object) -> None:
if self._page is not None:
try:
await self._page.close() # ТОЛЬКО своя вкладка
except Exception: # noqa: BLE001
pass
if self._pw is not None:
try:
await self._pw.stop()
except Exception: # noqa: BLE001
pass
def budget_left(self) -> bool:
return self._budget is None or self.loads < self._budget
async def _pause(self) -> None:
if self._last_load == 0.0:
return
jitter = self._delay * random.uniform(-0.2, 0.2)
wait = max(0.0, self._delay + jitter - (time.monotonic() - self._last_load))
if wait > 0:
print(f" пауза {wait:.1f} с", flush=True)
await asyncio.sleep(wait)
async def fetch(self, url: str) -> tuple[str, int | None]:
if not self.budget_left():
raise BudgetExhausted()
await self._pause()
resp = await self._page.goto(url, wait_until="domcontentloaded", timeout=90_000)
self.loads += 1
self._last_load = time.monotonic()
status = resp.status if resp is not None else None
try:
await self._page.wait_for_selector('[data-marker="item"]', timeout=7_000)
except Exception: # noqa: BLE001 — пустая/блочная страница разбирается ниже
pass
html = await self._page.content()
_guard(html, status)
return html, status
# --- заливка в msk_raw -----------------------------------------------------
def _sql_str(value: str) -> str:
return "'" + value.replace("'", "''") + "'"
def _csv_rows(rows: Iterable[dict[str, Any]]) -> str:
buf = io.StringIO()
writer = csv.writer(buf, lineterminator="\n")
for r in rows:
writer.writerow([
r["source_id"], r["observed_at"], r["batch_id"], r["kind"],
r["url"], r["price"], r["payload"],
])
return buf.getvalue()
def build_sql(batch_id: str, query: str, rows: list[dict[str, Any]],
started_at: str, kind: str = "serp") -> str:
"""Один поток на `psql -f -`: batch (FK!) → TEMP staging → \\copy → INSERT.
Одиночный `psql -c` через ssh ломается на квотинге скобок и кавычек, поэтому
только поток. rows_new = разница count(*) по batch_id до и после вставки.
"""
bid = _sql_str(batch_id)
return (
"BEGIN;\n"
"INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n"
f"VALUES ({bid}, {_sql_str(kind)}, {_sql_str(query)}, "
f"CAST({_sql_str(started_at)} AS timestamptz))\n"
"ON CONFLICT (batch_id) DO NOTHING;\n"
"CREATE TEMP TABLE _stg (LIKE msk_raw.avito_cards INCLUDING DEFAULTS) "
"ON COMMIT DROP;\n"
"CREATE TEMP TABLE _before ON COMMIT DROP AS\n"
f" SELECT count(*) AS n FROM msk_raw.avito_cards WHERE batch_id = {bid};\n"
"\\copy _stg (source_id,observed_at,batch_id,kind,url,price,payload) "
"FROM STDIN WITH (FORMAT csv)\n"
+ _csv_rows(rows)
+ "\\.\n"
"INSERT INTO msk_raw.avito_cards "
"(source_id,observed_at,batch_id,kind,url,price,payload)\n"
"SELECT source_id,observed_at,batch_id,kind,url,price,payload FROM _stg\n"
"ON CONFLICT (source_id,batch_id,kind) DO NOTHING;\n"
"UPDATE msk_raw.batches b SET\n"
" rows_sent = coalesce(b.rows_sent,0) + (SELECT count(*) FROM _stg),\n"
" rows_new = coalesce(b.rows_new,0) +\n"
f" ((SELECT count(*) FROM msk_raw.avito_cards WHERE batch_id = {bid})\n"
" - (SELECT n FROM _before))\n"
f"WHERE b.batch_id = {bid};\n"
"COMMIT;\n"
)
def build_finalize_sql(batch_id: str, query: str, notes: str) -> str:
bid = _sql_str(batch_id)
return (
"INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n"
f"VALUES ({bid}, 'serp', {_sql_str(query)}, now())\n"
"ON CONFLICT (batch_id) DO NOTHING;\n"
f"UPDATE msk_raw.batches SET finished_at = now(), notes = {_sql_str(notes)}\n"
f"WHERE batch_id = {bid};\n"
)
def run_psql(sql: str, ssh_host: str, container: str, db_user: str, db_name: str) -> None:
cmd = [
"ssh", ssh_host,
f"docker exec -i {container} psql -U {db_user} -d {db_name} "
"-v ON_ERROR_STOP=1 -f -",
]
proc = subprocess.run(cmd, input=sql.encode("utf-8"), capture_output=True)
out = (proc.stdout + proc.stderr).decode("utf-8", "replace").strip()
if proc.returncode != 0:
raise RuntimeError(f"psql через ssh вернул {proc.returncode}:\n{out}")
if out:
print(f" psql: {out}", flush=True)
# --- накопитель карточек ---------------------------------------------------
@dataclass
class Sink:
"""Батчами на прод (ssh+psql) или в локальный CSV при --dry-run."""
batch_id: str
started_at: str
query: str
batch_size: int
dry_run: bool
csv_path: Path
ssh_host: str
container: str
db_user: str
db_name: str
buffer: list[dict[str, Any]] = field(default_factory=list)
sent: int = 0
skipped_non_numeric: int = 0
def add(self, lot: Any) -> None:
raw_id = str(getattr(lot, "source_id", "") or "")
try:
source_id = int(raw_id) # в БД bigint, у ScrapedLot — строка
except (TypeError, ValueError):
self.skipped_non_numeric += 1
return
payload = lot.model_dump(mode="json")
self.buffer.append({
"source_id": source_id,
"observed_at": datetime.now(timezone.utc).isoformat(),
"batch_id": self.batch_id,
"kind": "serp",
"url": payload.get("source_url"),
"price": payload.get("price_rub"),
"payload": json.dumps(payload, ensure_ascii=False),
})
def maybe_flush(self) -> None:
if len(self.buffer) >= self.batch_size:
self.flush()
def flush(self) -> None:
if not self.buffer:
return
rows, self.buffer = self.buffer, []
if self.dry_run:
fresh = not self.csv_path.exists()
with self.csv_path.open("a", encoding="utf-8", newline="") as fh:
if fresh:
fh.write("source_id,observed_at,batch_id,kind,url,price,payload\n")
fh.write(_csv_rows(rows))
print(f" [dry-run] {len(rows)} строк → {self.csv_path}", flush=True)
else:
run_psql(
build_sql(self.batch_id, self.query, rows, self.started_at),
self.ssh_host, self.container, self.db_user, self.db_name,
)
print(f" залито {len(rows)} строк в msk_raw.avito_cards", flush=True)
self.sent += len(rows)
def finalize(self, notes: str) -> None:
self.flush()
if self.dry_run:
print(f" [dry-run] finalize: {notes}", flush=True)
return
run_psql(
build_finalize_sql(self.batch_id, self.query, notes),
self.ssh_host, self.container, self.db_user, self.db_name,
)
# --- сбор ------------------------------------------------------------------
def parse_page(scraper: AvitoScraper, html: str, url: str) -> tuple[int | None, list[Any]]:
count = scraper._extract_total_count(html)
lots = scraper._parse_html(html, "https://www.avito.ru")
if not lots and count:
raise Blocked("empty_page", f"0 карточек при счётчике {count}: {url}")
return count, lots
async def probe(loader: Loader, scraper: AvitoScraper, base_url: str,
lo: int | None, hi: int | None) -> tuple[int | None, list[Any]]:
url = build_url(base_url, 1, lo, hi)
html, _ = await loader.fetch(url)
return parse_page(scraper, html, url)
async def build_plan(loader: Loader, scraper: AvitoScraper, base_url: str, target: int,
cache: dict[tuple[int | None, int | None], list[Any]]
) -> list[Corridor]:
"""Адаптивная бисекция по цене; страница 1 каждого коридора кэшируется."""
corridors: list[Corridor] = []
def emit(lo: int | None, hi: int | None, count: int | None,
truncated: bool, lots: list[Any]) -> None:
missed = max(0, (count or 0) - HARD_CAP) if truncated else 0
c = Corridor(lo=lo, hi=hi, count=count, truncated=truncated, missed=missed)
corridors.append(c)
cache[(lo, hi)] = lots
flag = " TRUNCATED" if truncated else ""
print(f" коридор {c.label()} count={count} "
f"страниц={c.planned_pages()}{flag}", flush=True)
if truncated:
print(f" ВНИМАНИЕ: коридор {c.label()} не влезает в потолок "
f"{HARD_CAP}; заведомо не добрано ~{missed} объявлений", flush=True)
async def find_upper(lo: int | None) -> int:
"""Верхнюю границу открытого коридора ищем удвоением от разумного старта."""
cand = max(int(lo or PRICE_FLOOR) * 2, PRICE_PROBE_START)
while cand < PRICE_CEIL:
cnt, _ = await probe(loader, scraper, base_url, cand, None)
print(f" проба хвоста pmin={cand:_} count={cnt}", flush=True)
if cnt is not None and cnt <= target:
return cand
cand *= 2
return cand
async def split(lo: int | None, hi: int | None, depth: int,
count: int | None, lots: list[Any]) -> None:
if count is None:
url = build_url(base_url, 1, lo, hi)
raise Blocked("empty_page", f"счётчик не прочитался: {url}")
if count <= target:
emit(lo, hi, count, False, lots)
return
if depth >= MAX_DEPTH or width_ratio(lo, hi) <= MIN_WIDTH_RATIO:
# Предохранитель: не молчим — помечаем truncated и считаем недобор.
emit(lo, hi, count, count > HARD_CAP, lots)
return
upper = hi if hi is not None else await find_upper(lo)
if hi is None:
tail_cnt, tail_lots = await probe(loader, scraper, base_url, upper, None)
emit(upper, None, tail_cnt,
bool(tail_cnt and tail_cnt > HARD_CAP), tail_lots)
mid = geometric_mid(lo, upper)
for sub_lo, sub_hi in ((lo, mid), (mid, upper)):
sub_cnt, sub_lots = await probe(loader, scraper, base_url, sub_lo, sub_hi)
print(f" проба {sub_lo or '-'}..{sub_hi} count={sub_cnt}", flush=True)
await split(sub_lo, sub_hi, depth + 1, sub_cnt, sub_lots)
root_cnt, root_lots = await probe(loader, scraper, base_url, None, None)
print(f"Всего по базовому запросу: {root_cnt}", flush=True)
await split(None, None, 0, root_cnt, root_lots)
return corridors
async def collect(args: argparse.Namespace) -> int:
# avito_serp_ekb_only=False обязателен: с True парсер выбрасывает всё, где в
# URL нет /ekaterinburg/ — то есть все подмосковные слаги (serp.py:2154).
scraper = AvitoScraper(
SimpleNamespace(avito_serp_ekb_only=False), # type: ignore[arg-type]
target_city_slug="moskva",
)
out_dir = Path(args.out_dir).resolve()
out_dir.mkdir(parents=True, exist_ok=True)
plan_path = out_dir / f"plan-{args.batch_id}.json"
csv_path = out_dir / f"cards-{args.batch_id}.csv"
started_at = datetime.now(timezone.utc).isoformat()
plan: Plan | None = None
if args.resume:
if not plan_path.exists():
print(f"--resume: плана нет — {plan_path}", file=sys.stderr)
return 1
plan = Plan.load(plan_path)
done = sum(1 for c in plan.corridors if c.status == "done")
print(f"Resume по {plan_path}: коридоров {len(plan.corridors)}, "
f"готово {done}", flush=True)
# Resume: URL берём из сохранённого плана, а не из CLI — коридоры посчитаны
# именно под него. Расхождение = молчаливая заливка чужой выдачи под тем же
# batch_id, поэтому это ошибка, а не тихий приоритет одного из двух.
if plan is not None and plan.base_url != args.base_url:
raise SystemExit(
"--resume: план построен для другого URL."
f" В плане {plan.base_url}, в аргументах {args.base_url}."
" Убери --base-url (возьмётся из плана) либо начни новый batch_id."
)
base_url = plan.base_url if plan is not None else args.base_url
page_budget = None if args.full else args.measure
mode = "FULL" if args.full else f"MEASURE<={page_budget}"
print(f"Режим: {mode}; batch_id={args.batch_id}; delay={args.delay}s; "
f"target={args.target_count}; dry_run={args.dry_run}", flush=True)
sink = Sink(
batch_id=args.batch_id, started_at=started_at, query=base_url,
batch_size=args.batch_size, dry_run=args.dry_run, csv_path=csv_path,
ssh_host=args.ssh_host, container=args.container,
db_user=args.db_user, db_name=args.db_name,
)
cache: dict[tuple[int | None, int | None], list[Any]] = {}
total = 0
stop_reason = ""
rc = 0
loads = 0
async with Loader(args.delay, page_budget) as loader:
try:
if plan is None:
print("Строю план коридоров...", flush=True)
corridors = await build_plan(loader, scraper, base_url,
args.target_count, cache)
plan = Plan(base_url=base_url, target=args.target_count,
batch_id=args.batch_id, corridors=corridors,
created_at=started_at)
plan.save(plan_path)
print(f"План сохранён: {plan_path} ({len(corridors)} коридоров)",
flush=True)
for corridor in plan.corridors:
if corridor.status == "done":
continue
pages = corridor.planned_pages()
print(f"Коридор {corridor.label()} count={corridor.count} "
f"страниц={pages} (с {corridor.pages_done + 1})", flush=True)
for page in range(corridor.pages_done + 1, pages + 1):
key = (corridor.lo, corridor.hi)
if page == 1 and key in cache:
lots = cache.pop(key) # страница 1 уже скачана при планировании
else:
url = build_url(base_url, page, corridor.lo, corridor.hi)
html, _ = await loader.fetch(url)
_, lots = parse_page(scraper, html, url)
for lot in lots:
sink.add(lot)
total += len(lots)
corridor.pages_done = page
print(f" стр.{page}/{pages}: карточек {len(lots)}, "
f"итого {total}", flush=True)
sink.maybe_flush()
plan.save(plan_path)
if not lots:
print(" пустая страница — конец коридора", flush=True)
break
corridor.status = "done"
plan.save(plan_path)
except BudgetExhausted:
stop_reason = "потолок --measure исчерпан"
print(f"Стоп: {stop_reason}", flush=True)
except Blocked as exc:
stop_reason = f"BLOCKED/{exc.reason}: {exc.detail}"
print(f"СТОП: {stop_reason}", file=sys.stderr, flush=True)
rc = 2
finally:
loads = loader.loads
if plan is not None:
plan.save(plan_path)
truncated = [c for c in (plan.corridors if plan else []) if c.truncated]
missed = sum(c.missed for c in truncated)
notes = "; ".join(x for x in [
f"mode={mode}", f"loads={loads}", f"cards={total}",
f"skipped_non_numeric={sink.skipped_non_numeric}",
(f"truncated_corridors={len(truncated)} missed~{missed}" if truncated else ""),
stop_reason,
] if x)
try:
sink.finalize(notes)
except Exception as exc: # noqa: BLE001 — не прятать исходную причину стопа
print(f"finalize провалился: {exc}", file=sys.stderr)
rc = rc or 1
print(f"Готово. Загрузок: {loads}; карточек: {total}; отправлено: {sink.sent}; "
f"пропущено нечисловых source_id: {sink.skipped_non_numeric}; "
f"notes: {notes}", flush=True)
return rc
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
p = argparse.ArgumentParser(
prog="collect.py",
description="Ручной сбор SERP Авито (вторичка, Москва+МО) в прод-схему msk_raw.",
)
p.add_argument("--base-url", default=DEFAULT_BASE_URL,
help="базовый URL выдачи (дефолт — вкладка владельца)")
p.add_argument("--measure", type=int, default=100, metavar="N",
help="режим замера: не больше N загрузок страниц (дефолт 100)")
p.add_argument("--full", action="store_true",
help="полный проход без потолка страниц (включается только явно)")
p.add_argument("--dry-run", action="store_true",
help="ничего не слать на прод, писать CSV локально")
p.add_argument("--resume", action="store_true",
help="продолжить по сохранённому плану коридоров")
p.add_argument("--delay", type=float, default=8.0,
help="пауза между загрузками, с (±20%% джиттер, дефолт 8.0)")
p.add_argument("--batch-size", type=int, default=1000,
help="карточек в одной заливке (дефолт 1000)")
p.add_argument("--target-count", type=int, default=1500,
help="целевой размер коридора; больше — делим (дефолт 1500)")
p.add_argument("--batch-id", default=None,
help="batch_id в msk_raw.batches (дефолт msk-serp-<UTC>)")
p.add_argument("--out-dir", default=str(Path(__file__).resolve().parent / "runs"),
help="каталог плана/CSV")
p.add_argument("--ssh-host", default="selectel", help="ssh-хост прода")
p.add_argument("--container", default="tradein-postgres",
help="имя контейнера Postgres на проде")
p.add_argument("--db-user", default="tradein")
p.add_argument("--db-name", default="tradein")
args = p.parse_args(argv)
if args.batch_id is None:
args.batch_id = "msk-serp-" + datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
if not _BATCH_ID_RE.match(args.batch_id):
p.error("--batch-id: допустимы только символы [A-Za-z0-9._-]")
if args.measure < 1:
p.error("--measure должен быть >= 1")
return args
def main(argv: list[str] | None = None) -> int:
return asyncio.run(collect(parse_args(argv)))
if __name__ == "__main__":
raise SystemExit(main())