feat(db): backfill complexes nulls + add job_settings (migrations 80, 81)
80_complexes_backfill_nulls.sql: - cad_quarter: +722 rows via ST_Contains spatial join on cad_quarters_geom - developer_id: +247 rows via fuzzy name match on domrf_developers - obj_class: +2 rows via objective_lots.class mode - cad_building_num TEXT dropped (100% NULL; relation lives in cad_buildings.complex_id) - v_complex_full rebuilt: cad_building_nums TEXT[] replaces dead column - v_complex_buildings created: complex_id → cad_building_nums[], buildings_count 81_job_settings.sql: - New table job_settings (job_type PK, enabled, queue_name, cron_schedule, rate_ms, max_retries, max_concurrency, extra_config JSONB) - Seeded 4 rows: scrape_kn (queue=scrape_kn, '15 4 * * mon'), nspd_geo (queue=geo, rate_ms=600), objective_etl, objective_sync (cron from prod objective_sync_config) - objective_sync_config marked DEPRECATED via COMMENT — kept for backward compat until backend refactor Both already applied on prod via postgres MCP. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
de21fac0a8
commit
cbe1253461
2 changed files with 366 additions and 0 deletions
223
data/sql/80_complexes_backfill_nulls.sql
Normal file
223
data/sql/80_complexes_backfill_nulls.sql
Normal file
|
|
@ -0,0 +1,223 @@
|
||||||
|
-- 80_complexes_backfill_nulls.sql
|
||||||
|
-- Quick-win backfill для public.complexes через spatial joins + lookups.
|
||||||
|
-- Закрывает NULL'ы без новых scrapers.
|
||||||
|
--
|
||||||
|
-- Что делает:
|
||||||
|
-- A. cad_quarter — spatial join ST_Contains(cad_quarters_geom.geom, complexes.geom)
|
||||||
|
-- B. cad_building_num — колонка 100% NULL, никто не пишет в неё.
|
||||||
|
-- Выбрана Опция A: колонка дропается, relation живёт через
|
||||||
|
-- cad_buildings.complex_id (FK 1:N, уже существует).
|
||||||
|
-- Создаётся VIEW v_complex_buildings.
|
||||||
|
-- v_complex_full пересоздаётся без cad_building_num + с
|
||||||
|
-- cad_building_nums TEXT[] из агрегации по cad_buildings.
|
||||||
|
-- C. district_name — spatial join ST_Contains(ekb_districts.geom, complexes.geom)
|
||||||
|
-- (только EKB-покрытие, ~9 районов)
|
||||||
|
-- D. developer_id — fuzzy match lower(trim(developer_name)) с domrf_developers.
|
||||||
|
-- Ambiguous (>1 match по имени) — пропускаем, берём только
|
||||||
|
-- уникальные совпадения.
|
||||||
|
-- E. obj_class — источник данных domrf_kn_objects.obj_class = NULL для
|
||||||
|
-- всех 4548 строк; objective_lots.class даёт только 2 complexes.
|
||||||
|
-- Шаг выполняется (закрывает эти 2), но большого эффекта нет.
|
||||||
|
-- F. wall_type — domrf_kn_objects.wall_type = NULL для всех строк;
|
||||||
|
-- objective_lots не содержит wall_type.
|
||||||
|
-- Источника данных нет — шаг пропускается.
|
||||||
|
--
|
||||||
|
-- Dependencies:
|
||||||
|
-- - public.complexes (таблица)
|
||||||
|
-- - public.cad_quarters_geom (GIST индекс cad_quarters_geom_gist)
|
||||||
|
-- - public.ekb_districts (GIST индекс ekb_districts_geom_idx)
|
||||||
|
-- - public.domrf_developers (btree на developer_id PK)
|
||||||
|
-- - public.objective_lots (FK complex_id, поле class)
|
||||||
|
-- - public.cad_buildings (FK complex_id)
|
||||||
|
-- - public.v_complex_full (пересоздаётся — DROP CASCADE + CREATE)
|
||||||
|
--
|
||||||
|
-- Apply order: после 78_drop_use_rosreestr2coord.sql
|
||||||
|
--
|
||||||
|
-- Идемпотентность:
|
||||||
|
-- - UPDATE ... WHERE column IS NULL — повторный запуск не перезапишет NOT NULL.
|
||||||
|
-- - DROP TABLE IF EXISTS / CREATE OR REPLACE VIEW — безопасны.
|
||||||
|
-- - DROP COLUMN IF EXISTS — безопасен.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- A. cad_quarter — spatial join
|
||||||
|
-- ============================================================
|
||||||
|
-- Spatial indexes используются (EXPLAIN: Index Scan on cad_quarters_geom).
|
||||||
|
-- LIMIT 1 защищает от boundary случаев когда точка лежит на двух кварталах.
|
||||||
|
UPDATE complexes c
|
||||||
|
SET cad_quarter = (
|
||||||
|
SELECT q.cad_number
|
||||||
|
FROM cad_quarters_geom q
|
||||||
|
WHERE ST_Contains(q.geom, c.geom)
|
||||||
|
LIMIT 1
|
||||||
|
)
|
||||||
|
WHERE c.cad_quarter IS NULL AND c.geom IS NOT NULL;
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- B. cad_building_num — Опция A: дропаем колонку, relation через cad_buildings.complex_id
|
||||||
|
-- ============================================================
|
||||||
|
|
||||||
|
-- Шаг B.1: пересоздать v_complex_full без cad_building_num.
|
||||||
|
-- Вместо неё добавляем cad_building_nums TEXT[] (агрегация из cad_buildings).
|
||||||
|
-- DROP CASCADE нужен потому что PostgreSQL не позволяет CREATE OR REPLACE VIEW
|
||||||
|
-- если изменяется список/порядок колонок (cad_building_num убираем, cad_building_nums добавляем).
|
||||||
|
DROP VIEW IF EXISTS v_complex_full CASCADE;
|
||||||
|
|
||||||
|
CREATE OR REPLACE VIEW v_complex_full AS
|
||||||
|
SELECT
|
||||||
|
c.id AS complex_id,
|
||||||
|
c.canonical_name,
|
||||||
|
c.developer_name,
|
||||||
|
c.region_id,
|
||||||
|
c.district_name,
|
||||||
|
c.obj_class,
|
||||||
|
c.address,
|
||||||
|
c.latitude,
|
||||||
|
c.longitude,
|
||||||
|
(c.geom IS NOT NULL) AS has_geom,
|
||||||
|
c.flat_count,
|
||||||
|
c.square_living,
|
||||||
|
c.site_status,
|
||||||
|
c.cad_quarter,
|
||||||
|
-- cad_building_nums заменяет cad_building_num: агрегат из cad_buildings.cad_num
|
||||||
|
(
|
||||||
|
SELECT array_agg(b.cad_num ORDER BY b.cad_num)
|
||||||
|
FROM cad_buildings b
|
||||||
|
WHERE b.complex_id = c.id
|
||||||
|
) AS cad_building_nums,
|
||||||
|
(
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM cad_buildings b
|
||||||
|
WHERE b.complex_id = c.id
|
||||||
|
) AS cad_buildings_n,
|
||||||
|
(
|
||||||
|
SELECT array_agg(DISTINCT cs.source ORDER BY cs.source)
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
) AS sources,
|
||||||
|
(
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
) AS source_links_n,
|
||||||
|
(
|
||||||
|
SELECT COUNT(DISTINCT cs.source)
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
) AS source_kinds_n,
|
||||||
|
(
|
||||||
|
SELECT array_agg(cs.source_external_id ORDER BY cs.source_external_id)
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
AND cs.source = 'domrf_kn'
|
||||||
|
) AS domrf_obj_ids,
|
||||||
|
(
|
||||||
|
SELECT cs.source_id
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
AND cs.source = 'objective'
|
||||||
|
LIMIT 1
|
||||||
|
) AS objective_project_name,
|
||||||
|
(
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM domrf_kn_flats f
|
||||||
|
WHERE f.obj_id IN (
|
||||||
|
SELECT cs.source_external_id
|
||||||
|
FROM complex_sources cs
|
||||||
|
WHERE cs.complex_id = c.id
|
||||||
|
AND cs.source = 'domrf_kn'
|
||||||
|
AND cs.source_external_id IS NOT NULL
|
||||||
|
)
|
||||||
|
) AS kn_flats_n,
|
||||||
|
(
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM objective_lots ol
|
||||||
|
WHERE ol.complex_id = c.id
|
||||||
|
) AS objective_lots_n,
|
||||||
|
c.primary_source,
|
||||||
|
c.created_at,
|
||||||
|
c.updated_at
|
||||||
|
FROM complexes c;
|
||||||
|
|
||||||
|
-- Шаг B.2: дропаем колонку cad_building_num (100% NULL, индекс дропается автоматически)
|
||||||
|
ALTER TABLE complexes DROP COLUMN IF EXISTS cad_building_num;
|
||||||
|
|
||||||
|
-- Шаг B.3: VIEW v_complex_buildings — удобный агрегат complex → buildings
|
||||||
|
CREATE OR REPLACE VIEW v_complex_buildings AS
|
||||||
|
SELECT
|
||||||
|
c.id AS complex_id,
|
||||||
|
c.canonical_name,
|
||||||
|
array_agg(DISTINCT b.cad_num ORDER BY b.cad_num) AS cad_building_nums,
|
||||||
|
COUNT(b.cad_num) AS buildings_count
|
||||||
|
FROM complexes c
|
||||||
|
LEFT JOIN cad_buildings b ON b.complex_id = c.id
|
||||||
|
GROUP BY c.id, c.canonical_name;
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- C. district_name — spatial join (EKB coverage only)
|
||||||
|
-- ============================================================
|
||||||
|
-- EXPLAIN: Index Scan on ekb_districts_geom_idx — GIST index используется.
|
||||||
|
-- Закрывает только complexes внутри 9 районов EKB (geom IS NOT NULL).
|
||||||
|
-- Complexes без geom или за пределами EKB остаются NULL.
|
||||||
|
UPDATE complexes c
|
||||||
|
SET district_name = (
|
||||||
|
SELECT d.district_name
|
||||||
|
FROM ekb_districts d
|
||||||
|
WHERE d.geom IS NOT NULL
|
||||||
|
AND ST_Contains(d.geom, c.geom)
|
||||||
|
LIMIT 1
|
||||||
|
)
|
||||||
|
WHERE c.district_name IS NULL AND c.geom IS NOT NULL;
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- D. developer_id — fuzzy name match с domrf_developers
|
||||||
|
-- ============================================================
|
||||||
|
-- Берём только случаи с ровно одним совпадением по нормализованному имени
|
||||||
|
-- (lower + trim). Ambiguous (29 complexes с ≥2 candidates) — пропускаем.
|
||||||
|
UPDATE complexes c
|
||||||
|
SET developer_id = (
|
||||||
|
SELECT d.developer_id
|
||||||
|
FROM domrf_developers d
|
||||||
|
WHERE lower(trim(d.developer_name)) = lower(trim(c.developer_name))
|
||||||
|
LIMIT 1
|
||||||
|
)
|
||||||
|
WHERE c.developer_id IS NULL
|
||||||
|
AND c.developer_name IS NOT NULL
|
||||||
|
AND (
|
||||||
|
-- только если ровно одна запись в domrf_developers с таким именем
|
||||||
|
SELECT COUNT(DISTINCT d2.developer_id)
|
||||||
|
FROM domrf_developers d2
|
||||||
|
WHERE lower(trim(d2.developer_name)) = lower(trim(c.developer_name))
|
||||||
|
) = 1;
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- E. obj_class — via objective_lots.class (2 complexes покрывается)
|
||||||
|
-- ============================================================
|
||||||
|
-- domrf_kn_objects.obj_class = NULL для всех строк (источник не заполнен).
|
||||||
|
-- objective_lots.class — единственный источник; покрывает ~2 complexes с NULL obj_class.
|
||||||
|
-- Берём mode (наиболее частый класс) по complex_id.
|
||||||
|
UPDATE complexes c
|
||||||
|
SET obj_class = (
|
||||||
|
SELECT ol.class
|
||||||
|
FROM objective_lots ol
|
||||||
|
WHERE ol.complex_id = c.id
|
||||||
|
AND ol.class IS NOT NULL
|
||||||
|
GROUP BY ol.class
|
||||||
|
ORDER BY COUNT(*) DESC
|
||||||
|
LIMIT 1
|
||||||
|
)
|
||||||
|
WHERE c.obj_class IS NULL
|
||||||
|
AND EXISTS (
|
||||||
|
SELECT 1 FROM objective_lots ol2
|
||||||
|
WHERE ol2.complex_id = c.id AND ol2.class IS NOT NULL
|
||||||
|
);
|
||||||
|
|
||||||
|
-- ============================================================
|
||||||
|
-- F. wall_type — источника данных нет
|
||||||
|
-- ============================================================
|
||||||
|
-- domrf_kn_objects.wall_type = NULL для всех 4548 строк.
|
||||||
|
-- objective_lots не содержит wall_type.
|
||||||
|
-- Шаг пропускается — нет источника для backfill.
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
143
data/sql/81_job_settings.sql
Normal file
143
data/sql/81_job_settings.sql
Normal file
|
|
@ -0,0 +1,143 @@
|
||||||
|
-- Централизованная таблица настроек для всех Celery-job'ов.
|
||||||
|
--
|
||||||
|
-- Контекст:
|
||||||
|
-- До этой миграции конфиг был разбросан по трём местам:
|
||||||
|
-- • scrape_kn — env-переменные SCRAPE_KN_CRON / SCRAPE_KN_DEFAULT_REGIONS
|
||||||
|
-- • nspd_geo — per-row rate_ms в nspd_geo_jobs (нет global default)
|
||||||
|
-- • objective_etl — hardcode в settings + env OBJECTIVE_ANTON_SQLITE_PATH
|
||||||
|
-- • objective_sync — уже в objective_sync_config (оставляем как deprecated fallback)
|
||||||
|
--
|
||||||
|
-- Цель: единая таблица job_settings — один UI /admin/jobs/settings для всех job'ов.
|
||||||
|
--
|
||||||
|
-- Зависимости:
|
||||||
|
-- Нет FK на другие таблицы. objective_sync_config НЕ дропается — она помечается
|
||||||
|
-- DEPRECATED и остаётся как fallback для уже задеплоенного backend-кода.
|
||||||
|
--
|
||||||
|
-- Порядок применения:
|
||||||
|
-- Применяется один раз на чистой или работающей БД.
|
||||||
|
-- Идемпотентно (CREATE IF NOT EXISTS + INSERT ON CONFLICT DO NOTHING).
|
||||||
|
-- Backend-код должен быть задеплоен ПОСЛЕ этой миграции.
|
||||||
|
--
|
||||||
|
-- Проверка после применения:
|
||||||
|
-- SELECT * FROM job_settings; -- должна вернуть 4 строки
|
||||||
|
-- SELECT * FROM objective_sync_config; -- должна по-прежнему работать
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
-- ── Создание таблицы ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS job_settings (
|
||||||
|
job_type TEXT PRIMARY KEY,
|
||||||
|
-- 'scrape_kn' | 'nspd_geo' | 'objective_etl' | 'objective_sync'
|
||||||
|
enabled BOOLEAN NOT NULL DEFAULT true,
|
||||||
|
queue_name TEXT NOT NULL DEFAULT 'celery',
|
||||||
|
-- celery (default) | scrape_kn | geo
|
||||||
|
cron_schedule TEXT,
|
||||||
|
-- Стандартный crontab: 'minute hour dom month dow'. NULL = только ручной запуск.
|
||||||
|
rate_ms INTEGER,
|
||||||
|
-- Пауза между API-запросами (мс). NULL/0 = не throttled.
|
||||||
|
max_retries SMALLINT NOT NULL DEFAULT 2,
|
||||||
|
max_concurrency SMALLINT NOT NULL DEFAULT 1,
|
||||||
|
-- Макс. параллельных task'ов этого типа.
|
||||||
|
extra_config JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||||
|
-- Job-specific: {"regions":[66], "groups_csv":"...", "sqlite_path":"..."}.
|
||||||
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||||
|
updated_by TEXT,
|
||||||
|
description TEXT
|
||||||
|
-- Человеко-читаемое описание для UI.
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE job_settings IS
|
||||||
|
'Centralized settings for all Celery jobs. See /admin/jobs/settings UI. '
|
||||||
|
'Replaces scattered env vars and per-table configs.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN job_settings.queue_name IS
|
||||||
|
'Celery queue routing. Allows concurrency between job types within one worker container.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN job_settings.cron_schedule IS
|
||||||
|
'Standard crontab string (5-field: minute hour dom month dow). '
|
||||||
|
'NULL = manual trigger only. Beat re-reads on restart.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN job_settings.rate_ms IS
|
||||||
|
'Delay between API calls in milliseconds (used by NSPD geo). NULL/0 for non-throttled jobs.';
|
||||||
|
|
||||||
|
COMMENT ON COLUMN job_settings.extra_config IS
|
||||||
|
'Job-specific JSON config: e.g. {"regions":[66,77], "groups_csv":"Свердловская область,...", '
|
||||||
|
'"sqlite_path":"/data/anton-sqlite/analysis.db"}.';
|
||||||
|
|
||||||
|
-- ── Seed данные ───────────────────────────────────────────────────────────────
|
||||||
|
-- ON CONFLICT DO NOTHING — безопасен при повторном запуске.
|
||||||
|
-- Defaults взяты из:
|
||||||
|
-- scrape_kn → backend/app/core/config.py (scrape_kn_cron / scrape_kn_default_regions)
|
||||||
|
-- nspd_geo → nspd_geo_jobs.rate_ms most common value (600) на проде
|
||||||
|
-- objective_etl → backend/app/core/config.py (objective_anton_sqlite_path)
|
||||||
|
-- objective_sync → objective_sync_config.cron_schedule (prod value '0 6 * * tue')
|
||||||
|
|
||||||
|
INSERT INTO job_settings
|
||||||
|
(job_type, enabled, queue_name, cron_schedule, rate_ms, max_retries,
|
||||||
|
max_concurrency, extra_config, description)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
'scrape_kn',
|
||||||
|
true,
|
||||||
|
'scrape_kn',
|
||||||
|
'15 4 * * mon',
|
||||||
|
NULL,
|
||||||
|
2,
|
||||||
|
1,
|
||||||
|
'{"default_regions": [66]}'::jsonb,
|
||||||
|
'Скрейпинг КН (ЦИАН): загрузка объявлений по регионам. '
|
||||||
|
'Расписание и регионы ранее задавались через env SCRAPE_KN_CRON / SCRAPE_KN_DEFAULT_REGIONS.'
|
||||||
|
),
|
||||||
|
(
|
||||||
|
'nspd_geo',
|
||||||
|
true,
|
||||||
|
'geo',
|
||||||
|
NULL,
|
||||||
|
600,
|
||||||
|
2,
|
||||||
|
1,
|
||||||
|
'{}'::jsonb,
|
||||||
|
'Геокодирование объектов через НСПД/rosreestr2coord. '
|
||||||
|
'Запускается вручную по job-конфигурациям nspd_geo_jobs. rate_ms — глобальный дефолт.'
|
||||||
|
),
|
||||||
|
(
|
||||||
|
'objective_etl',
|
||||||
|
true,
|
||||||
|
'celery',
|
||||||
|
NULL,
|
||||||
|
NULL,
|
||||||
|
2,
|
||||||
|
1,
|
||||||
|
'{"sqlite_path": "/data/anton-sqlite/analysis.db"}'::jsonb,
|
||||||
|
'ETL: SQLite Антоновского → PG (objective_* таблицы). '
|
||||||
|
'Запускается вручную. sqlite_path ранее задавался через env OBJECTIVE_ANTON_SQLITE_PATH.'
|
||||||
|
),
|
||||||
|
(
|
||||||
|
'objective_sync',
|
||||||
|
true,
|
||||||
|
'celery',
|
||||||
|
'0 6 * * tue',
|
||||||
|
NULL,
|
||||||
|
2,
|
||||||
|
1,
|
||||||
|
'{
|
||||||
|
"groups_csv": "Свердловская область,Челябинск,Тюмень,Пермь",
|
||||||
|
"use_ddu": true,
|
||||||
|
"use_dkp": true,
|
||||||
|
"period_months_back": 1,
|
||||||
|
"inter_group_delay_s": 30,
|
||||||
|
"rate_ms": 3000
|
||||||
|
}'::jsonb,
|
||||||
|
'Синхронизация DOM.РФ Objective (отчёты). '
|
||||||
|
'Детальный конфиг дублирован в objective_sync_config (deprecated legacy fallback).'
|
||||||
|
)
|
||||||
|
ON CONFLICT (job_type) DO NOTHING;
|
||||||
|
|
||||||
|
-- ── Пометить objective_sync_config как deprecated ─────────────────────────────
|
||||||
|
-- Таблица НЕ дропается: backend продолжает читать её как fallback до рефакторинга.
|
||||||
|
COMMENT ON TABLE objective_sync_config IS
|
||||||
|
'DEPRECATED: use job_settings (job_type=''objective_sync'') for new config reads. '
|
||||||
|
'Kept for backward compat until backend refactor. Do NOT drop until refactor is deployed.';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
Loading…
Add table
Reference in a new issue