gendesign/_wt-cadval/data/sql/77_nspd_geo_jobs.sql
bot-backend 542ff1c9ca docs(tradein/domclick): два комментария описывали пул, которого нет с миграции 253
Оба места утверждали, что у Домклика один выделенный резидентный прокси и
пула нет. Это перестало быть правдой ещё в #2800 (миграция 253 сняла
резервацию узла), но текст остался — и именно на него опирался тикет #3189,
поставленный под «калибровку» ограничения, которого не существует.
Свип к тому же ходит через пул давно (serp.py:359), а бэкфилл подключён к
нему в PR #3222.

Заодно докстринг называл не тот ограничитель: свип кладёт не счётчик
провалов lease, а break по первому DomClickBlockedError (#2854) — до
ротации дело не доходит ни при каком счётчике, отсюда buckets_completed=0.

Только комментарии, поведение не меняется. Миграция 175 уже применена, а
_schema_migrations трекает по имени файла без checksum — правка текста
её не перезапустит.
2026-08-29 18:57:25 +03:00

149 lines
7.7 KiB
SQL
Raw Blame History

This file contains ambiguous Unicode characters

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

-- Resume-friendly NSPD geo backfill jobs.
--
-- Зачем:
-- Bulk-загрузка кадастровых данных по 100k+ cad-номеров занимает часы.
-- При redeploy worker'а Celery теряет in-memory state. Решение —
-- хранить очередь и текущий прогресс в БД, после редеплоя
-- worker_ready hook подхватывает незавершённые jobs.
--
-- Архитектура:
-- nspd_geo_jobs — общий журнал заданий (один job = одна bulk-операция)
-- nspd_geo_targets — список cad-номеров для обработки (1 строка = 1 cad)
-- со статусом (pending/done/failed) — это и есть resume-state.
--
-- При redeploy:
-- worker_ready signal → SELECT FROM nspd_geo_jobs WHERE status='running'
-- → re-enqueue Celery task с тем же job_id
-- → task SELECT FROM nspd_geo_targets WHERE job_id=:id AND status='pending'
-- → продолжает с той точки где остановился (idempotent UPSERT в cad_quarters_geom)
--
-- Применять через mcp postgres execute_sql или psql direct.
CREATE TABLE IF NOT EXISTS nspd_geo_jobs (
job_id BIGSERIAL PRIMARY KEY,
name TEXT, -- человекочитаемое имя ('УрФО expansion 2026-Q2')
job_kind TEXT NOT NULL, -- 'quarters' | 'parcels' | 'buildings' | 'mixed'
-- Источник целей
source_kind TEXT NOT NULL, -- 'rosreestr_pending' | 'manual_list' | 'region_bbox'
source_params jsonb, -- что-то вроде {"region_codes":[66,74,72,59], "thematic_id":2}
-- Параметры fetch (всегда через rosreestr2coord lib v5+;
-- legacy колонка `use_rosreestr2coord` удалена миграцией 78_).
rate_ms INTEGER NOT NULL DEFAULT 600,
-- Состояние
status TEXT NOT NULL DEFAULT 'queued', -- queued | running | done | failed | paused
started_at timestamptz,
finished_at timestamptz,
heartbeat_at timestamptz,
-- Счётчики
targets_total INTEGER NOT NULL DEFAULT 0,
targets_done INTEGER NOT NULL DEFAULT 0,
targets_failed INTEGER NOT NULL DEFAULT 0,
targets_skipped INTEGER NOT NULL DEFAULT 0, -- уже было в БД (idempotent skip)
requests_count INTEGER NOT NULL DEFAULT 0,
waf_blocked_count INTEGER NOT NULL DEFAULT 0,
-- Audit
error TEXT,
triggered_by TEXT NOT NULL DEFAULT 'manual',-- manual | beat | resume
created_at timestamptz NOT NULL DEFAULT NOW(),
updated_at timestamptz NOT NULL DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS nspd_geo_jobs_status_idx
ON nspd_geo_jobs(status, created_at DESC);
CREATE INDEX IF NOT EXISTS nspd_geo_jobs_running_idx
ON nspd_geo_jobs(heartbeat_at DESC) WHERE status = 'running';
COMMENT ON TABLE nspd_geo_jobs IS
'Журнал bulk geo-jobs для NSPD-fetcher. Resume-friendly: при redeploy '
'worker_ready hook подхватывает status=running с устаревшим heartbeat.';
-- ── Targets (cad-номера к обработке) ─────────────────────────────────────────
CREATE TABLE IF NOT EXISTS nspd_geo_targets (
target_id BIGSERIAL PRIMARY KEY,
job_id BIGINT NOT NULL REFERENCES nspd_geo_jobs(job_id) ON DELETE CASCADE,
cad_num TEXT NOT NULL, -- '66:41:0204016' для квартала, '...:10' для участка
thematic_id SMALLINT NOT NULL, -- 1=parcel | 2=quarter | 5=building
status TEXT NOT NULL DEFAULT 'pending', -- pending | done | failed | skipped
attempts SMALLINT NOT NULL DEFAULT 0,
-- Результат
result_features INTEGER, -- сколько features вернулось
saved_to_table TEXT, -- 'cad_quarters_geom' | 'cad_buildings' | ...
error_msg TEXT,
fetched_at timestamptz,
-- Avoid duplicate work внутри одного job
UNIQUE (job_id, cad_num, thematic_id)
);
CREATE INDEX IF NOT EXISTS nspd_geo_targets_pending_idx
ON nspd_geo_targets(job_id, status) WHERE status = 'pending';
CREATE INDEX IF NOT EXISTS nspd_geo_targets_failed_idx
ON nspd_geo_targets(job_id) WHERE status = 'failed';
COMMENT ON TABLE nspd_geo_targets IS
'Список cad-номеров для обработки в рамках job. Resume-state: '
'task просто SELECT WHERE status=pending и продолжает.';
-- ── nspd_geo_log: логи (видны через v_scrape_log_unified) ───────────────────
CREATE TABLE IF NOT EXISTS nspd_geo_log (
log_id BIGSERIAL PRIMARY KEY,
job_id BIGINT REFERENCES nspd_geo_jobs(job_id) ON DELETE CASCADE,
ts timestamptz NOT NULL DEFAULT NOW(),
level TEXT NOT NULL DEFAULT 'info',
stage TEXT,
cad_num TEXT,
message TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS nspd_geo_log_job_idx
ON nspd_geo_log(job_id, log_id DESC);
CREATE INDEX IF NOT EXISTS nspd_geo_log_ts_idx
ON nspd_geo_log(ts DESC);
-- ── Расширение unified-views на nspd_geo (UNION ALL) ────────────────────────
CREATE OR REPLACE VIEW v_scrape_runs_unified AS
SELECT 'kn' AS scraper_type, run_id, started_at, finished_at, heartbeat_at,
status, error, NULL::text AS triggered_by,
region_codes::text AS scope, requests_count,
objects_count AS items_ok, NULL::integer AS items_failed,
flats_count AS sub_items_ok,
progress_obj_index AS progress_idx, total_obj_count AS progress_total,
jsonb_build_object('developers', developer_ids, 'snapshot_date', snapshot_date,
'resumed_from', resumed_from_run_id, 'params', params) AS extra
FROM kn_scrape_runs
UNION ALL
SELECT 'nspd', run_id, started_at, finished_at, heartbeat_at, status, error,
triggered_by, region_code::text, requests_count,
quarters_ok, quarters_failed, buildings_ok,
NULL, pending_count,
jsonb_build_object('waf_429_count', waf_429_count)
FROM nspd_scrape_runs
UNION ALL
SELECT 'objective', run_id, started_at, finished_at, heartbeat_at, status, error,
triggered_by, group_name, requests_count,
reports_ok, reports_failed, rows_lots,
NULL, NULL,
jsonb_build_object('rows_corpus_room', rows_corpus_room, 'rows_history', rows_history)
FROM objective_scrape_runs
UNION ALL
SELECT 'nspd_geo', job_id, started_at, finished_at, heartbeat_at, status, error,
triggered_by, COALESCE(name, job_kind),
requests_count, targets_done, targets_failed, targets_skipped,
targets_done, targets_total,
jsonb_build_object('job_kind', job_kind, 'source_kind', source_kind,
'rate_ms', rate_ms,
'waf_blocked_count', waf_blocked_count,
'source_params', source_params)
FROM nspd_geo_jobs;
CREATE OR REPLACE VIEW v_scrape_log_unified AS
SELECT 'kn' AS scraper_type, log_id, run_id, ts, level, stage,
obj_id::text AS entity_id, message
FROM kn_scrape_log
UNION ALL
SELECT 'nspd', log_id, run_id, ts, level, stage, cad_number, message
FROM nspd_scrape_log
UNION ALL
SELECT 'nspd_geo', log_id, job_id AS run_id, ts, level, stage, cad_num, message
FROM nspd_geo_log;