fix
This commit is contained in:
parent
25b73035a1
commit
db6711446a
6 changed files with 229 additions and 5 deletions
|
|
@ -1278,7 +1278,9 @@ def _district_cadastre_baseline(db: Session, *, district_name: str) -> dict[str,
|
||||||
WHERE cb.cost_value IS NOT NULL
|
WHERE cb.cost_value IS NOT NULL
|
||||||
AND cb.area IS NOT NULL
|
AND cb.area IS NOT NULL
|
||||||
AND cb.area >= 100
|
AND cb.area >= 100
|
||||||
AND (cb.floors IS NOT NULL AND cb.floors >= 3
|
-- floors хранится как TEXT (встречаются '1-2', '2-3') —
|
||||||
|
-- считаем только чистые числа ≥3, либо purpose-fallback.
|
||||||
|
AND ((cb.floors ~ '^[0-9]+$' AND cb.floors::int >= 3)
|
||||||
OR cb.purpose ILIKE '%многокв%')
|
OR cb.purpose ILIKE '%многокв%')
|
||||||
AND (cb.cost_value / NULLIF(cb.area, 0))
|
AND (cb.cost_value / NULLIF(cb.area, 0))
|
||||||
BETWEEN 5000 AND 500000
|
BETWEEN 5000 AND 500000
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,18 @@ fi
|
||||||
|
|
||||||
cd "$(dirname "$0")/../raw"
|
cd "$(dirname "$0")/../raw"
|
||||||
|
|
||||||
|
# Имя CSV-файла внутри zip может варьироваться (rosreestr менял схему в Q4 2024).
|
||||||
|
# Скрипт автоматически пробует обе именные конвенции.
|
||||||
|
# Если файла нет — квартал пропускается с warning (не fail). Это позволяет
|
||||||
|
# догружать 2023 / 2024Q1-Q2 / новые кварталы постепенно по мере выкачки.
|
||||||
declare -a JOBS=(
|
declare -a JOBS=(
|
||||||
|
# Формат CSV изменился в Q4 2024 (delim ; → ~, имя r_all → r-r_01-92).
|
||||||
|
# 2024Q1+Q2 уже опубликованы по новой схеме (r-r_01-92_y_YYYY_q_N), но
|
||||||
|
# delimiter мог не успеть перейти — попробуй сначала ~, при ошибках смени на ;.
|
||||||
|
# 2023 в виде квартальных CSV НЕ ПУБЛИКУЕТСЯ (архив до 2023г. содержит только
|
||||||
|
# JSON-снимки от ноября-декабря 2021 в /data-sets/Архив до 2023г. включительно/).
|
||||||
|
"2024Q1:2024-01-01:dataset_СДЕЛКИ_r-r_01-92_y_2024_q_1.csv:;"
|
||||||
|
"2024Q2:2024-04-01:dataset_СДЕЛКИ_r-r_01-92_y_2024_q_2.csv:;"
|
||||||
"2024Q3:2024-07-01:dataset_СДЕЛКИ_r_all_q_3.csv:;"
|
"2024Q3:2024-07-01:dataset_СДЕЛКИ_r_all_q_3.csv:;"
|
||||||
"2024Q4:2024-10-01:dataset_СДЕЛКИ_r-r_01-92_y_2024_q_4.csv:~"
|
"2024Q4:2024-10-01:dataset_СДЕЛКИ_r-r_01-92_y_2024_q_4.csv:~"
|
||||||
"2025Q1:2025-01-01:dataset_СДЕЛКИ_r-r_01-92_y_2025_q_1.csv:~"
|
"2025Q1:2025-01-01:dataset_СДЕЛКИ_r-r_01-92_y_2025_q_1.csv:~"
|
||||||
|
|
@ -33,8 +44,55 @@ declare -a JOBS=(
|
||||||
"2026Q1:2026-01-01:dataset_СДЕЛКИ_r-r_01-92_y_2026_q_1.csv:~"
|
"2026Q1:2026-01-01:dataset_СДЕЛКИ_r-r_01-92_y_2026_q_1.csv:~"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Возможные имена zip-ов в data/raw/ (можно класть как есть, скрипт распакует).
|
||||||
|
# Имя в data/raw/ — sdelki_YYYYqN.csv.zip ИЛИ оригинальное dataset_СДЕЛКИ_*.csv.zip.
|
||||||
|
|
||||||
|
resolve_csv() {
|
||||||
|
# Принимает comma-separated список кандидатов имён csv. Возвращает первый
|
||||||
|
# существующий (распакованный) файл. Если ничего нет — пытается распаковать
|
||||||
|
# из zip с известными именами. Возвращает пустую строку если файл не найден.
|
||||||
|
local CANDIDATES="$1"
|
||||||
|
local SRC_Q="$2"
|
||||||
|
IFS=',' read -ra NAMES <<< "$CANDIDATES"
|
||||||
|
for name in "${NAMES[@]}"; do
|
||||||
|
if [[ -f "$name" ]]; then
|
||||||
|
echo "$name"
|
||||||
|
return 0
|
||||||
|
fi
|
||||||
|
done
|
||||||
|
# Try unzip from sdelki_YYYYqN.csv.zip (lowercase quarter).
|
||||||
|
local LOWER=$(echo "$SRC_Q" | tr '[:upper:]' '[:lower:]')
|
||||||
|
local ZIP_LOCAL="sdelki_${LOWER}.csv.zip"
|
||||||
|
if [[ -f "$ZIP_LOCAL" ]]; then
|
||||||
|
unzip -o -q "$ZIP_LOCAL" >&2 || true
|
||||||
|
for name in "${NAMES[@]}"; do
|
||||||
|
if [[ -f "$name" ]]; then echo "$name"; return 0; fi
|
||||||
|
done
|
||||||
|
fi
|
||||||
|
# Try original rosreestr-named zips.
|
||||||
|
for name in "${NAMES[@]}"; do
|
||||||
|
if [[ -f "${name}.zip" ]]; then
|
||||||
|
unzip -o -q "${name}.zip" >&2 || true
|
||||||
|
[[ -f "$name" ]] && { echo "$name"; return 0; }
|
||||||
|
fi
|
||||||
|
done
|
||||||
|
echo ""
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
|
||||||
for job in "${JOBS[@]}"; do
|
for job in "${JOBS[@]}"; do
|
||||||
IFS=':' read -r SRC_Q PERIOD FILE SEP <<< "$job"
|
IFS=':' read -r SRC_Q PERIOD FILE_CANDIDATES SEP <<< "$job"
|
||||||
|
FILE=$(resolve_csv "$FILE_CANDIDATES" "$SRC_Q")
|
||||||
|
if [[ -z "$FILE" ]]; then
|
||||||
|
echo "=== $SRC_Q SKIP — нет CSV в data/raw/ (искал: $FILE_CANDIDATES; sdelki_${SRC_Q,,}.csv.zip)"
|
||||||
|
continue
|
||||||
|
fi
|
||||||
|
# Idempotency: пропускаем если уже загружен с тем же source_quarter.
|
||||||
|
EXISTING=$($PSQL -t -A -c "SELECT COUNT(*) FROM rosreestr_deals WHERE source_quarter = '$SRC_Q'" || echo 0)
|
||||||
|
if [[ "${EXISTING:-0}" -gt 0 ]]; then
|
||||||
|
echo "=== $SRC_Q SKIP — уже загружен ($EXISTING строк). Удалить: DELETE FROM rosreestr_deals WHERE source_quarter='$SRC_Q'"
|
||||||
|
continue
|
||||||
|
fi
|
||||||
echo "=== $SRC_Q ($FILE, sep='$SEP') ==="
|
echo "=== $SRC_Q ($FILE, sep='$SEP') ==="
|
||||||
$PSQL -c "TRUNCATE rosreestr_deals_staging;"
|
$PSQL -c "TRUNCATE rosreestr_deals_staging;"
|
||||||
start=$(date +%s)
|
start=$(date +%s)
|
||||||
|
|
|
||||||
146
data/sql/02b_load_via_tunnel.py
Normal file
146
data/sql/02b_load_via_tunnel.py
Normal file
|
|
@ -0,0 +1,146 @@
|
||||||
|
"""Загрузчик кварталов dataset_СДЕЛКИ через SSH-tunnel на localhost:15432.
|
||||||
|
|
||||||
|
Эквивалент 02_load_all_quarters.sh, но через psycopg вместо docker — удобнее
|
||||||
|
запускать с Windows / любой машины с уже открытым SSH-tunnel.
|
||||||
|
|
||||||
|
Idempotent: пропускает кварталы, уже залитые в rosreestr_deals.
|
||||||
|
|
||||||
|
Usage:
|
||||||
|
cd backend
|
||||||
|
uv run python ../data/sql/02b_load_via_tunnel.py 2024Q1 2024Q2
|
||||||
|
uv run python ../data/sql/02b_load_via_tunnel.py # все из JOBS
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import time
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import psycopg
|
||||||
|
|
||||||
|
PG_DSN = os.environ.get(
|
||||||
|
"TUNNEL_DSN",
|
||||||
|
"postgresql://gendesign:2J2SBPMKuS998fiwhtQqDhMI@localhost:15432/gendesign",
|
||||||
|
)
|
||||||
|
|
||||||
|
RAW = Path(__file__).resolve().parent.parent / "raw"
|
||||||
|
|
||||||
|
# (source_quarter, period_start, csv_filename, delimiter)
|
||||||
|
JOBS: list[tuple[str, str, str, str]] = [
|
||||||
|
("2024Q1", "2024-01-01", "dataset_СДЕЛКИ_r-r_01-92_y_2024_q_1.csv", ";"),
|
||||||
|
("2024Q2", "2024-04-01", "dataset_СДЕЛКИ_r-r_01-92_y_2024_q_2.csv", ";"),
|
||||||
|
# Q3+ уже залиты, но оставляем как fallback / re-run target.
|
||||||
|
("2024Q3", "2024-07-01", "dataset_СДЕЛКИ_r_all_q_3.csv", ";"),
|
||||||
|
("2024Q4", "2024-10-01", "dataset_СДЕЛКИ_r-r_01-92_y_2024_q_4.csv", "~"),
|
||||||
|
("2025Q1", "2025-01-01", "dataset_СДЕЛКИ_r-r_01-92_y_2025_q_1.csv", "~"),
|
||||||
|
("2025Q2", "2025-04-01", "dataset_СДЕЛКИ_r-r_01-92_y_2025_q_2.csv", "~"),
|
||||||
|
("2025Q3", "2025-07-01", "dataset_СДЕЛКИ_r-r_01-92_y_2025_q_3.csv", "~"),
|
||||||
|
("2025Q4", "2025-10-01", "dataset_СДЕЛКИ_r-r_01-92_y_2025_q_4.csv", "~"),
|
||||||
|
("2026Q1", "2026-01-01", "dataset_СДЕЛКИ_r-r_01-92_y_2026_q_1.csv", "~"),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
COPY_SQL = """
|
||||||
|
COPY rosreestr_deals_staging
|
||||||
|
FROM STDIN (FORMAT csv, HEADER true, DELIMITER %r, QUOTE '"', NULL '')
|
||||||
|
"""
|
||||||
|
|
||||||
|
INSERT_SQL = """
|
||||||
|
INSERT INTO rosreestr_deals (
|
||||||
|
source_quarter, period_start_date,
|
||||||
|
region_code, okato, district, city, quarter_cad_number, street,
|
||||||
|
realestate_type_code, wall_material_code, purpose_code,
|
||||||
|
year_build, floor, area,
|
||||||
|
doc_type, deal_price, currency, deal_count
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
%s, period_start_date,
|
||||||
|
region_code, NULLIF(okato, '-'), NULLIF(district, ''), NULLIF(city, ''),
|
||||||
|
quarter_cad_number, NULLIF(street, ''),
|
||||||
|
NULLIF(realestate_type_code, ''),
|
||||||
|
NULLIF(wall_material_code, ''),
|
||||||
|
NULLIF(purpose_code, ''),
|
||||||
|
CASE WHEN year_build ~ '^[12][0-9]{3}$' THEN year_build::INT ELSE NULL END,
|
||||||
|
NULLIF(floor, ''),
|
||||||
|
CASE WHEN area > 0 AND area < 1e8 THEN area ELSE NULL END,
|
||||||
|
NULLIF(doc_type, ''),
|
||||||
|
CASE WHEN deal_price > 0 AND deal_price < 1e15 THEN deal_price ELSE NULL END,
|
||||||
|
COALESCE(NULLIF(currency, ''), 'рубль'),
|
||||||
|
COALESCE(deal_count, 1)
|
||||||
|
FROM rosreestr_deals_staging
|
||||||
|
WHERE region_code IS NOT NULL
|
||||||
|
AND quarter_cad_number IS NOT NULL
|
||||||
|
AND quarter_cad_number <> ''
|
||||||
|
AND doc_type IS NOT NULL
|
||||||
|
AND doc_type <> ''
|
||||||
|
AND period_start_date IS NOT NULL
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> int:
|
||||||
|
targets = sys.argv[1:] or [j[0] for j in JOBS]
|
||||||
|
selected = [j for j in JOBS if j[0] in targets]
|
||||||
|
if not selected:
|
||||||
|
print(f"Нет matching кварталов в JOBS. Targets: {targets}")
|
||||||
|
return 1
|
||||||
|
|
||||||
|
conn = psycopg.connect(PG_DSN, connect_timeout=15)
|
||||||
|
try:
|
||||||
|
for sq, period, fname, sep in selected:
|
||||||
|
with conn.cursor() as cur:
|
||||||
|
cur.execute(
|
||||||
|
"SELECT COUNT(*) FROM rosreestr_deals WHERE source_quarter = %s",
|
||||||
|
(sq,),
|
||||||
|
)
|
||||||
|
already = cur.fetchone()[0]
|
||||||
|
if already:
|
||||||
|
print(f"=== {sq} SKIP — уже {already} строк (DELETE FROM rosreestr_deals WHERE source_quarter='{sq}' чтобы перезалить)")
|
||||||
|
continue
|
||||||
|
csv_path = RAW / fname
|
||||||
|
if not csv_path.exists():
|
||||||
|
print(f"=== {sq} SKIP — нет файла {csv_path}")
|
||||||
|
continue
|
||||||
|
print(f"=== {sq} ({fname}, sep='{sep}', {csv_path.stat().st_size/1e6:.1f} MB) ===")
|
||||||
|
t0 = time.time()
|
||||||
|
with conn.cursor() as cur:
|
||||||
|
cur.execute("TRUNCATE rosreestr_deals_staging")
|
||||||
|
copy_sql = (
|
||||||
|
f"COPY rosreestr_deals_staging FROM STDIN "
|
||||||
|
f"(FORMAT csv, HEADER true, DELIMITER '{sep}', QUOTE '\"', NULL '')"
|
||||||
|
)
|
||||||
|
with cur.copy(copy_sql) as copy:
|
||||||
|
with open(csv_path, "rb") as f:
|
||||||
|
while True:
|
||||||
|
chunk = f.read(8 * 1024 * 1024)
|
||||||
|
if not chunk:
|
||||||
|
break
|
||||||
|
copy.write(chunk)
|
||||||
|
t_copy = time.time() - t0
|
||||||
|
cur.execute(INSERT_SQL, (sq,))
|
||||||
|
inserted = cur.rowcount
|
||||||
|
conn.commit()
|
||||||
|
t_total = time.time() - t0
|
||||||
|
print(f" COPY {t_copy:.0f}s, INSERT {inserted} rows, total {t_total:.0f}s")
|
||||||
|
|
||||||
|
# Final summary
|
||||||
|
with conn.cursor() as cur:
|
||||||
|
cur.execute(
|
||||||
|
"""
|
||||||
|
SELECT source_quarter, COUNT(*) AS rows, SUM(deal_count) AS deals
|
||||||
|
FROM rosreestr_deals
|
||||||
|
GROUP BY source_quarter
|
||||||
|
ORDER BY source_quarter
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
print("\n=== Итоги по rosreestr_deals ===")
|
||||||
|
for sq, rows, deals in cur.fetchall():
|
||||||
|
print(f" {sq} rows={rows:>9} deals={deals:>10}")
|
||||||
|
finally:
|
||||||
|
conn.close()
|
||||||
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
sys.exit(main())
|
||||||
|
|
@ -18,7 +18,9 @@ CREATE TABLE IF NOT EXISTS cad_buildings (
|
||||||
building_name TEXT,
|
building_name TEXT,
|
||||||
readable_address TEXT,
|
readable_address TEXT,
|
||||||
area NUMERIC,
|
area NUMERIC,
|
||||||
floors INT,
|
-- NSPD отдаёт floors как строку — встречаются диапазоны '1-2', '2-3'.
|
||||||
|
-- Не приводим к INT, чтобы не терять данные. Запросы используют regex-CAST.
|
||||||
|
floors TEXT,
|
||||||
year_built INT,
|
year_built INT,
|
||||||
year_commisioning INT,
|
year_commisioning INT,
|
||||||
cost_value NUMERIC,
|
cost_value NUMERIC,
|
||||||
|
|
@ -26,7 +28,7 @@ CREATE TABLE IF NOT EXISTS cad_buildings (
|
||||||
status TEXT,
|
status TEXT,
|
||||||
ownership_type TEXT,
|
ownership_type TEXT,
|
||||||
cultural_heritage TEXT,
|
cultural_heritage TEXT,
|
||||||
underground_floors INT,
|
underground_floors TEXT,
|
||||||
build_record_area NUMERIC,
|
build_record_area NUMERIC,
|
||||||
build_record_type TEXT,
|
build_record_type TEXT,
|
||||||
common_data_status TEXT,
|
common_data_status TEXT,
|
||||||
|
|
|
||||||
|
|
@ -69,7 +69,9 @@ quarter_cadastre AS (
|
||||||
AND cost_value IS NOT NULL
|
AND cost_value IS NOT NULL
|
||||||
AND area IS NOT NULL
|
AND area IS NOT NULL
|
||||||
AND area >= 100
|
AND area >= 100
|
||||||
AND (floors IS NOT NULL AND floors >= 3
|
-- floors хранится как TEXT (встречаются '1-2', '2-3') —
|
||||||
|
-- считаем только чистые числа ≥3, либо purpose-fallback.
|
||||||
|
AND ((floors ~ '^[0-9]+$' AND floors::int >= 3)
|
||||||
OR purpose ILIKE '%многокв%')
|
OR purpose ILIKE '%многокв%')
|
||||||
AND (cost_value / NULLIF(area, 0)) BETWEEN 5000 AND 500000
|
AND (cost_value / NULLIF(area, 0)) BETWEEN 5000 AND 500000
|
||||||
GROUP BY quarter_cad_num
|
GROUP BY quarter_cad_num
|
||||||
|
|
|
||||||
14
data/sql/65_partitions_rosreestr_2024_h1.sql
Normal file
14
data/sql/65_partitions_rosreestr_2024_h1.sql
Normal file
|
|
@ -0,0 +1,14 @@
|
||||||
|
-- Partitions для 2024 Q1+Q2 (расширение rosreestr_deals назад во времени).
|
||||||
|
-- Применять однократно перед загрузкой 02_load_all_quarters.sh. Идемпотентно.
|
||||||
|
--
|
||||||
|
-- 2023 в виде квартальных CSV rosreestr НЕ ПУБЛИКУЕТСЯ — поэтому партиций для
|
||||||
|
-- 2023 не создаём. Архив /data-sets/Архив до 2023г. включительно/ содержит
|
||||||
|
-- только JSON-снимки от ноя-дек 2021 в формате export-json-region-NN-date-TS.zip
|
||||||
|
-- (per-region полные дампы), несовместимые с rosreestr_deals_staging.
|
||||||
|
--
|
||||||
|
-- Range FROM inclusive, TO exclusive — [start_q, start_next_q).
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS rosreestr_deals_2024q1 PARTITION OF rosreestr_deals
|
||||||
|
FOR VALUES FROM ('2024-01-01') TO ('2024-04-01');
|
||||||
|
CREATE TABLE IF NOT EXISTS rosreestr_deals_2024q2 PARTITION OF rosreestr_deals
|
||||||
|
FOR VALUES FROM ('2024-04-01') TO ('2024-07-01');
|
||||||
Loading…
Add table
Reference in a new issue