Compare commits
No commits in common. "cfa24c3bcf0c2cf8f3e6cbc3ca199a77013f14ac" and "6553e5830f53ddfcc159368325ad5b4e73c68b50" have entirely different histories.
cfa24c3bcf
...
6553e5830f
6 changed files with 26 additions and 326 deletions
|
|
@ -1056,13 +1056,12 @@ async def scrape_cian_backfill_history(
|
||||||
Iterates Cian listings with missing offer_price_history, fetches Cian detail
|
Iterates Cian listings with missing offer_price_history, fetches Cian detail
|
||||||
pages, persists enrichment (price changes, views, agent, ceiling_height, etc.).
|
pages, persists enrichment (price changes, views, agent, ceiling_height, etc.).
|
||||||
|
|
||||||
Houses block (houses_price_dynamics): fetches newbuilding pages for houses with
|
Houses block (houses_price_dynamics) is currently a no-op — requires
|
||||||
cian_zhk_url set and no existing price dynamics rows. Requires migration
|
a cian_zhk_url column on houses (migration pending, see TODO in
|
||||||
071_houses_cian_zhk_url.sql applied and cian_zhk_url populated (via SERP scrape
|
app/tasks/cian_history_backfill.py).
|
||||||
or scripts/backfill_cian_zhk_url.py).
|
|
||||||
|
|
||||||
Rate-limit: scraper_settings 'cian' delay (default 5s) between requests.
|
Rate-limit: scraper_settings 'cian' delay (default 5s) between requests.
|
||||||
Idempotent: skips listings/houses that already have history rows.
|
Idempotent: skips listings that already have offer_price_history rows.
|
||||||
"""
|
"""
|
||||||
from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history
|
from app.tasks.cian_history_backfill import CianBackfillResult, backfill_cian_history
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -33,8 +33,6 @@ class NewbuildingEnrichment:
|
||||||
|
|
||||||
# Core house metadata (from state.newbuilding)
|
# Core house metadata (from state.newbuilding)
|
||||||
cian_internal_house_id: int | None = None
|
cian_internal_house_id: int | None = None
|
||||||
# ZHK page URL — canonical slug form stored after fetch / redirect resolution
|
|
||||||
cian_zhk_url: str | None = None
|
|
||||||
name: str | None = None
|
name: str | None = None
|
||||||
address: str | None = None
|
address: str | None = None
|
||||||
year_built: int | None = None
|
year_built: int | None = None
|
||||||
|
|
@ -123,7 +121,6 @@ async def fetch_newbuilding(
|
||||||
|
|
||||||
result = NewbuildingEnrichment(
|
result = NewbuildingEnrichment(
|
||||||
cian_internal_house_id=nb.get("id"),
|
cian_internal_house_id=nb.get("id"),
|
||||||
cian_zhk_url=zhk_url,
|
|
||||||
name=nb.get("name") or nb.get("displayName"),
|
name=nb.get("name") or nb.get("displayName"),
|
||||||
address=_extract_address(nb),
|
address=_extract_address(nb),
|
||||||
year_built=nb.get("yearBuilt") or nb.get("buildYear"),
|
year_built=nb.get("yearBuilt") or nb.get("buildYear"),
|
||||||
|
|
@ -391,7 +388,7 @@ def _extract_nested_offers(offers_state: dict[str, Any]) -> list[dict[str, Any]]
|
||||||
# ---- save helpers ----
|
# ---- save helpers ----
|
||||||
|
|
||||||
|
|
||||||
def save_newbuilding_enrichment(
|
async def save_newbuilding_enrichment(
|
||||||
db: Any,
|
db: Any,
|
||||||
house_id: int,
|
house_id: int,
|
||||||
enrichment: NewbuildingEnrichment,
|
enrichment: NewbuildingEnrichment,
|
||||||
|
|
@ -400,7 +397,7 @@ def save_newbuilding_enrichment(
|
||||||
|
|
||||||
Steps:
|
Steps:
|
||||||
1. UPSERT management_companies (if present) → get mc_id
|
1. UPSERT management_companies (if present) → get mc_id
|
||||||
2. UPDATE houses with Cian metadata (incl. cian_zhk_url if present)
|
2. UPDATE houses with Cian metadata
|
||||||
3. INSERT INTO houses_price_dynamics (chart points, ON CONFLICT DO UPDATE)
|
3. INSERT INTO houses_price_dynamics (chart points, ON CONFLICT DO UPDATE)
|
||||||
4. INSERT INTO house_reliability_checks (overall + details)
|
4. INSERT INTO house_reliability_checks (overall + details)
|
||||||
"""
|
"""
|
||||||
|
|
@ -449,7 +446,6 @@ def save_newbuilding_enrichment(
|
||||||
cian_internal_house_id = COALESCE(
|
cian_internal_house_id = COALESCE(
|
||||||
CAST(:cihi AS bigint), cian_internal_house_id
|
CAST(:cihi AS bigint), cian_internal_house_id
|
||||||
),
|
),
|
||||||
cian_zhk_url = COALESCE(:zhk_url, cian_zhk_url),
|
|
||||||
year_built = COALESCE(:yb, year_built),
|
year_built = COALESCE(:yb, year_built),
|
||||||
management_company_id = COALESCE(CAST(:mcid AS bigint), management_company_id),
|
management_company_id = COALESCE(CAST(:mcid AS bigint), management_company_id),
|
||||||
transport_accessibility_rate = COALESCE(
|
transport_accessibility_rate = COALESCE(
|
||||||
|
|
@ -464,7 +460,6 @@ def save_newbuilding_enrichment(
|
||||||
{
|
{
|
||||||
"hid": house_id,
|
"hid": house_id,
|
||||||
"cihi": enrichment.cian_internal_house_id,
|
"cihi": enrichment.cian_internal_house_id,
|
||||||
"zhk_url": enrichment.cian_zhk_url,
|
|
||||||
"yb": enrichment.year_built,
|
"yb": enrichment.year_built,
|
||||||
"mcid": mc_id,
|
"mcid": mc_id,
|
||||||
"tar": enrichment.transport_accessibility_rate,
|
"tar": enrichment.transport_accessibility_rate,
|
||||||
|
|
@ -551,68 +546,3 @@ def save_newbuilding_enrichment(
|
||||||
len(enrichment.reliability_checks),
|
len(enrichment.reliability_checks),
|
||||||
mc_id,
|
mc_id,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
async def resolve_cian_zhk_url(
|
|
||||||
cian_internal_house_id: int,
|
|
||||||
*,
|
|
||||||
session: AsyncSession | None = None,
|
|
||||||
) -> str | None:
|
|
||||||
"""Resolve canonical ZHK slug URL from a Cian internal house ID.
|
|
||||||
|
|
||||||
Cian redirects https://cian.ru/zhk/<id>/ → canonical slug URL
|
|
||||||
(e.g. https://zhk-park-ekb-i.cian.ru/). We follow the redirect and
|
|
||||||
return the final URL.
|
|
||||||
|
|
||||||
Used by backfill scripts (scripts/backfill_cian_zhk_url.py) for houses
|
|
||||||
that have cian_internal_house_id set but cian_zhk_url IS NULL.
|
|
||||||
|
|
||||||
Args:
|
|
||||||
cian_internal_house_id: Cian's numeric ЖК identifier.
|
|
||||||
session: optional shared curl_cffi AsyncSession (caller owns lifecycle).
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
Canonical ZHK URL string, or None if request failed / redirect not followed.
|
|
||||||
|
|
||||||
Note:
|
|
||||||
This function makes a real HTTP request — do NOT call it without rate limiting.
|
|
||||||
Caller must enforce per-request sleep matching scraper_settings 'cian' delay.
|
|
||||||
"""
|
|
||||||
close_session = False
|
|
||||||
if session is None:
|
|
||||||
session = AsyncSession(
|
|
||||||
impersonate="chrome120",
|
|
||||||
timeout=15.0,
|
|
||||||
headers={
|
|
||||||
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
|
|
||||||
"Accept-Language": "ru-RU,ru;q=0.9,en;q=0.8",
|
|
||||||
},
|
|
||||||
)
|
|
||||||
close_session = True
|
|
||||||
|
|
||||||
fallback_url = f"https://cian.ru/zhk/{cian_internal_house_id}/"
|
|
||||||
try:
|
|
||||||
resp = await session.get(fallback_url, allow_redirects=True)
|
|
||||||
final_url = str(resp.url)
|
|
||||||
if final_url and final_url != fallback_url:
|
|
||||||
logger.debug(
|
|
||||||
"resolve_cian_zhk_url id=%s → %s", cian_internal_house_id, final_url
|
|
||||||
)
|
|
||||||
return final_url
|
|
||||||
# Redirect not followed or same URL — return fallback as canonical
|
|
||||||
if resp.status_code == 200:
|
|
||||||
return fallback_url
|
|
||||||
logger.warning(
|
|
||||||
"resolve_cian_zhk_url id=%s: HTTP %d, no usable URL",
|
|
||||||
cian_internal_house_id,
|
|
||||||
resp.status_code,
|
|
||||||
)
|
|
||||||
return None
|
|
||||||
except Exception as exc:
|
|
||||||
logger.warning(
|
|
||||||
"resolve_cian_zhk_url id=%s failed: %s", cian_internal_house_id, exc
|
|
||||||
)
|
|
||||||
return None
|
|
||||||
finally:
|
|
||||||
if close_session:
|
|
||||||
await session.close()
|
|
||||||
|
|
|
||||||
|
|
@ -9,11 +9,11 @@
|
||||||
AND NOT EXISTS (SELECT 1 FROM offer_price_history WHERE listing_id=...)
|
AND NOT EXISTS (SELECT 1 FROM offer_price_history WHERE listing_id=...)
|
||||||
LIMIT batch_size
|
LIMIT batch_size
|
||||||
→ fetch_detail + save_detail_enrichment per listing
|
→ fetch_detail + save_detail_enrichment per listing
|
||||||
2. Houses block: SELECT houses WHERE cian_zhk_url IS NOT NULL
|
2. Houses block — TODO (PR #525): houses не имеют cian_url / cian_zhk_url колонки.
|
||||||
AND NOT EXISTS (SELECT 1 FROM houses_price_dynamics WHERE house_id=...)
|
Колонка cian_internal_house_id (020_houses_alter_cian.sql) хранит numeric Cian ID,
|
||||||
LIMIT batch_size
|
но не ZHK URL, необходимый для fetch_newbuilding().
|
||||||
→ fetch_newbuilding(cian_zhk_url) + save_newbuilding_enrichment per house
|
Backfill ЖК-динамики требует отдельной миграции, добавляющей cian_zhk_url к houses.
|
||||||
Requires migration 071_houses_cian_zhk_url.sql (cian_zhk_url column).
|
До тех пор do_houses=True — no-op (houses_total=0).
|
||||||
|
|
||||||
Rate limit: scraper_settings.get_scraper_delay('cian') between requests.
|
Rate limit: scraper_settings.get_scraper_delay('cian') between requests.
|
||||||
"""
|
"""
|
||||||
|
|
@ -71,8 +71,8 @@ async def backfill_cian_history(
|
||||||
batch_size: max entities to process per call (listings + houses + valuations counted
|
batch_size: max entities to process per call (listings + houses + valuations counted
|
||||||
separately).
|
separately).
|
||||||
do_listings: process listings batch (offer_price_history backfill).
|
do_listings: process listings batch (offer_price_history backfill).
|
||||||
do_houses: process houses batch (houses_price_dynamics backfill via
|
do_houses: reserved for future houses_price_dynamics backfill (currently no-op —
|
||||||
fetch_newbuilding). Requires migration 071_houses_cian_zhk_url.sql applied.
|
see module docstring TODO).
|
||||||
do_valuations: process Cian Valuation Calculator batch (external_valuations backfill).
|
do_valuations: process Cian Valuation Calculator batch (external_valuations backfill).
|
||||||
Default False — opt-in because each call hits Cian auth-gated API.
|
Default False — opt-in because each call hits Cian auth-gated API.
|
||||||
dry_run: skip all fetch+save; only count and log pending rows.
|
dry_run: skip all fetch+save; only count and log pending rows.
|
||||||
|
|
@ -149,85 +149,18 @@ async def backfill_cian_history(
|
||||||
await asyncio.sleep(delay)
|
await asyncio.sleep(delay)
|
||||||
|
|
||||||
# ── 2. Houses: missing houses_price_dynamics ──────────────────────────────
|
# ── 2. Houses: missing houses_price_dynamics ──────────────────────────────
|
||||||
|
# TODO (PR #525): houses table has cian_internal_house_id (bigint) but no
|
||||||
|
# cian_zhk_url column. fetch_newbuilding() requires a ZHK URL
|
||||||
|
# (e.g. https://zhk-...-i.cian.ru/). Add migration `063_houses_add_cian_url.sql`
|
||||||
|
# with `ALTER TABLE houses ADD COLUMN IF NOT EXISTS cian_zhk_url text` and populate
|
||||||
|
# it during the existing SERP scrape (CianScraper → save_listings → house upsert).
|
||||||
|
# Until that column exists, do_houses is accepted but silently no-ops.
|
||||||
if do_houses:
|
if do_houses:
|
||||||
from app.services.scrapers.cian_newbuilding import (
|
logger.info(
|
||||||
fetch_newbuilding,
|
"cian_backfill: houses block skipped — cian_zhk_url column not yet added "
|
||||||
save_newbuilding_enrichment,
|
"(see TODO in cian_history_backfill.py)"
|
||||||
)
|
)
|
||||||
|
|
||||||
house_rows = (
|
|
||||||
db.execute(
|
|
||||||
text("""
|
|
||||||
SELECT h.id, h.cian_zhk_url
|
|
||||||
FROM houses h
|
|
||||||
WHERE h.cian_zhk_url IS NOT NULL
|
|
||||||
AND NOT EXISTS (
|
|
||||||
SELECT 1 FROM houses_price_dynamics hpd
|
|
||||||
WHERE hpd.house_id = h.id
|
|
||||||
)
|
|
||||||
LIMIT :lim
|
|
||||||
"""),
|
|
||||||
{"lim": batch_size},
|
|
||||||
)
|
|
||||||
.mappings()
|
|
||||||
.all()
|
|
||||||
)
|
|
||||||
result.houses_total = len(house_rows)
|
|
||||||
|
|
||||||
if dry_run:
|
|
||||||
logger.info("dry_run: would process %d cian houses", result.houses_total)
|
|
||||||
else:
|
|
||||||
for hrow in house_rows:
|
|
||||||
house_id: int = hrow["id"]
|
|
||||||
zhk_url: str = hrow["cian_zhk_url"]
|
|
||||||
result.houses_processed += 1
|
|
||||||
|
|
||||||
enrichment = None
|
|
||||||
try:
|
|
||||||
enrichment = await fetch_newbuilding(zhk_url)
|
|
||||||
except Exception as exc:
|
|
||||||
logger.warning(
|
|
||||||
"cian_newbuilding fetch failed for house_id=%s url=%s: %s",
|
|
||||||
house_id,
|
|
||||||
zhk_url,
|
|
||||||
exc,
|
|
||||||
)
|
|
||||||
result.houses_failed_fetch += 1
|
|
||||||
await asyncio.sleep(delay)
|
|
||||||
continue
|
|
||||||
|
|
||||||
if enrichment is None:
|
|
||||||
logger.warning(
|
|
||||||
"cian_newbuilding fetch returned None for house_id=%s url=%s",
|
|
||||||
house_id,
|
|
||||||
zhk_url,
|
|
||||||
)
|
|
||||||
result.houses_failed_fetch += 1
|
|
||||||
await asyncio.sleep(delay)
|
|
||||||
continue
|
|
||||||
|
|
||||||
try:
|
|
||||||
save_newbuilding_enrichment(db, house_id, enrichment)
|
|
||||||
result.houses_succeeded += 1
|
|
||||||
except Exception as exc:
|
|
||||||
logger.warning(
|
|
||||||
"cian_newbuilding save failed for house_id=%s: %s", house_id, exc
|
|
||||||
)
|
|
||||||
result.houses_failed_save += 1
|
|
||||||
# Roll back to clean session state so next house can proceed.
|
|
||||||
# save_newbuilding_enrichment commits on success; on failure the
|
|
||||||
# transaction is left open/dirty — rollback to avoid session poison.
|
|
||||||
try:
|
|
||||||
db.rollback()
|
|
||||||
except Exception as rb_exc:
|
|
||||||
logger.warning(
|
|
||||||
"cian_newbuilding rollback failed for house_id=%s: %s",
|
|
||||||
house_id,
|
|
||||||
rb_exc,
|
|
||||||
)
|
|
||||||
|
|
||||||
await asyncio.sleep(delay)
|
|
||||||
|
|
||||||
# ── 3. Cian listings без external_valuations (price prediction backfill) ──
|
# ── 3. Cian listings без external_valuations (price prediction backfill) ──
|
||||||
if do_valuations:
|
if do_valuations:
|
||||||
rows = db.execute(
|
rows = db.execute(
|
||||||
|
|
|
||||||
|
|
@ -1,38 +0,0 @@
|
||||||
-- 071_houses_cian_zhk_url.sql
|
|
||||||
-- Purpose: Add cian_zhk_url TEXT column to houses — blocker for Cian newbuilding history backfill.
|
|
||||||
--
|
|
||||||
-- Context (#568):
|
|
||||||
-- houses already has cian_internal_house_id (bigint, 020_houses_alter_cian.sql) but
|
|
||||||
-- fetch_newbuilding() requires a slug-based ZHK URL (e.g. https://zhk-...-i.cian.ru/).
|
|
||||||
-- The numeric ID alone is not sufficient because the ZHK URL encodes a human-readable
|
|
||||||
-- slug Cian uses for routing — it cannot be reliably reconstructed from ID alone.
|
|
||||||
--
|
|
||||||
-- Backfill strategy (post-deploy, see scripts/backfill_cian_zhk_url.py):
|
|
||||||
-- For houses where cian_internal_house_id IS NOT NULL but cian_zhk_url IS NULL,
|
|
||||||
-- the canonical fallback URL https://cian.ru/zhk/<id>/ redirects to the slug URL.
|
|
||||||
-- The backfill script resolves final redirect URLs and stores the canonical slug form.
|
|
||||||
-- Do NOT run from inside the migration — requires live HTTP access + rate limiting.
|
|
||||||
--
|
|
||||||
-- Idempotency: IF NOT EXISTS on both ALTER and CREATE INDEX.
|
|
||||||
-- Dependencies: 009_houses.sql, 020_houses_alter_cian.sql
|
|
||||||
-- Apply after: 070_houses_dadata_enrichment.sql
|
|
||||||
|
|
||||||
BEGIN;
|
|
||||||
|
|
||||||
ALTER TABLE houses
|
|
||||||
ADD COLUMN IF NOT EXISTS cian_zhk_url text;
|
|
||||||
|
|
||||||
-- Partial index: used by backfill queries to find houses that already have a URL
|
|
||||||
-- and by fetch_newbuilding() callers joining on URL.
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_houses_cian_zhk_url
|
|
||||||
ON houses (cian_zhk_url)
|
|
||||||
WHERE cian_zhk_url IS NOT NULL;
|
|
||||||
|
|
||||||
COMMENT ON COLUMN houses.cian_zhk_url IS
|
|
||||||
'Canonical Cian ЖК page URL (e.g. https://zhk-slug-city-i.cian.ru/). '
|
|
||||||
'Set during SERP scrape when Cian returns the slug, or via backfill '
|
|
||||||
'(scripts/backfill_cian_zhk_url.py) resolving redirect from '
|
|
||||||
'https://cian.ru/zhk/<cian_internal_house_id>/. '
|
|
||||||
'Required by fetch_newbuilding() for houses_price_dynamics backfill (#568).';
|
|
||||||
|
|
||||||
COMMIT;
|
|
||||||
|
|
@ -1,122 +0,0 @@
|
||||||
"""Backfill houses.cian_zhk_url for rows that have cian_internal_house_id but no URL.
|
|
||||||
|
|
||||||
Requires migration 071_houses_cian_zhk_url.sql to be applied first.
|
|
||||||
|
|
||||||
What it does:
|
|
||||||
SELECT houses WHERE cian_internal_house_id IS NOT NULL AND cian_zhk_url IS NULL
|
|
||||||
For each row: resolve https://cian.ru/zhk/<id>/ → follow redirect → store final URL.
|
|
||||||
|
|
||||||
Cian redirects the canonical /zhk/<id>/ path to the slug URL
|
|
||||||
(e.g. https://zhk-park-ekb-i.cian.ru/). We capture the final redirect destination.
|
|
||||||
|
|
||||||
Usage:
|
|
||||||
cd tradein-mvp/backend
|
|
||||||
DATABASE_URL=postgresql://... uv run python ../scripts/backfill_cian_zhk_url.py
|
|
||||||
|
|
||||||
Env vars:
|
|
||||||
DATABASE_URL required — psycopg v3 DSN
|
|
||||||
DELAY_SEC default 5.0 — seconds between requests (Cian rate limit)
|
|
||||||
BATCH_LIMIT default 500 — max houses to process per run (safety cap)
|
|
||||||
|
|
||||||
Idempotent — skips houses where cian_zhk_url IS NOT NULL. Safe to re-run.
|
|
||||||
"""
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import asyncio
|
|
||||||
import logging
|
|
||||||
import os
|
|
||||||
import sys
|
|
||||||
import time
|
|
||||||
|
|
||||||
import psycopg
|
|
||||||
|
|
||||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "backend"))
|
|
||||||
|
|
||||||
from app.services.scrapers.cian_newbuilding import resolve_cian_zhk_url
|
|
||||||
|
|
||||||
logging.basicConfig(
|
|
||||||
format="%(asctime)s [%(levelname)s] %(message)s",
|
|
||||||
level=logging.INFO,
|
|
||||||
)
|
|
||||||
logger = logging.getLogger("backfill_cian_zhk_url")
|
|
||||||
|
|
||||||
DATABASE_URL = os.environ.get("DATABASE_URL", "")
|
|
||||||
DELAY_SEC = float(os.environ.get("DELAY_SEC", "5.0"))
|
|
||||||
BATCH_LIMIT = int(os.environ.get("BATCH_LIMIT", "500"))
|
|
||||||
|
|
||||||
|
|
||||||
async def main() -> None:
|
|
||||||
if not DATABASE_URL:
|
|
||||||
logger.error("DATABASE_URL env var not set")
|
|
||||||
sys.exit(1)
|
|
||||||
|
|
||||||
conn = await psycopg.AsyncConnection.connect(DATABASE_URL)
|
|
||||||
|
|
||||||
async with conn:
|
|
||||||
rows = await (
|
|
||||||
await conn.execute(
|
|
||||||
"""
|
|
||||||
SELECT id, cian_internal_house_id
|
|
||||||
FROM houses
|
|
||||||
WHERE cian_internal_house_id IS NOT NULL
|
|
||||||
AND cian_zhk_url IS NULL
|
|
||||||
ORDER BY id
|
|
||||||
LIMIT %(lim)s
|
|
||||||
""",
|
|
||||||
{"lim": BATCH_LIMIT},
|
|
||||||
)
|
|
||||||
).fetchall()
|
|
||||||
|
|
||||||
total = len(rows)
|
|
||||||
logger.info(
|
|
||||||
"Found %d houses with cian_internal_house_id and no cian_zhk_url (limit=%d)",
|
|
||||||
total,
|
|
||||||
BATCH_LIMIT,
|
|
||||||
)
|
|
||||||
if total == 0:
|
|
||||||
logger.info("Nothing to backfill.")
|
|
||||||
return
|
|
||||||
|
|
||||||
succeeded = 0
|
|
||||||
failed = 0
|
|
||||||
t0 = time.time()
|
|
||||||
|
|
||||||
conn = await psycopg.AsyncConnection.connect(DATABASE_URL)
|
|
||||||
async with conn:
|
|
||||||
for i, (house_id, cian_id) in enumerate(rows, start=1):
|
|
||||||
url = await resolve_cian_zhk_url(cian_id)
|
|
||||||
if url:
|
|
||||||
await conn.execute(
|
|
||||||
"""
|
|
||||||
UPDATE houses
|
|
||||||
SET cian_zhk_url = %(url)s
|
|
||||||
WHERE id = %(hid)s
|
|
||||||
AND cian_zhk_url IS NULL
|
|
||||||
""",
|
|
||||||
{"url": url, "hid": house_id},
|
|
||||||
)
|
|
||||||
await conn.commit()
|
|
||||||
succeeded += 1
|
|
||||||
logger.info("[%d/%d] house_id=%s id=%s → %s", i, total, house_id, cian_id, url)
|
|
||||||
else:
|
|
||||||
failed += 1
|
|
||||||
logger.warning(
|
|
||||||
"[%d/%d] house_id=%s id=%s — could not resolve URL",
|
|
||||||
i, total, house_id, cian_id,
|
|
||||||
)
|
|
||||||
|
|
||||||
if i < total:
|
|
||||||
await asyncio.sleep(DELAY_SEC)
|
|
||||||
|
|
||||||
elapsed = time.time() - t0
|
|
||||||
logger.info(
|
|
||||||
"Backfill complete: %d succeeded, %d failed, %.1fs total (%.1fs/req avg)",
|
|
||||||
succeeded,
|
|
||||||
failed,
|
|
||||||
elapsed,
|
|
||||||
elapsed / max(succeeded + failed, 1),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
asyncio.run(main())
|
|
||||||
|
|
@ -6,7 +6,7 @@ BATCH="${BATCH_SIZE:-100}"
|
||||||
MAX_ITERS="${MAX_ITERS:-1000}"
|
MAX_ITERS="${MAX_ITERS:-1000}"
|
||||||
DO_LISTINGS="${DO_LISTINGS:-true}"
|
DO_LISTINGS="${DO_LISTINGS:-true}"
|
||||||
DO_VALUATIONS="${DO_VALUATIONS:-true}"
|
DO_VALUATIONS="${DO_VALUATIONS:-true}"
|
||||||
DO_HOUSES="${DO_HOUSES:-true}" # enabled after migration 071_houses_cian_zhk_url.sql (#568)
|
DO_HOUSES="${DO_HOUSES:-false}" # houses stubbed until cian_zhk_url migration
|
||||||
|
|
||||||
usage() {
|
usage() {
|
||||||
cat <<EOF
|
cat <<EOF
|
||||||
|
|
@ -21,7 +21,7 @@ Env vars:
|
||||||
MAX_ITERS default: 1000 (safety stop)
|
MAX_ITERS default: 1000 (safety stop)
|
||||||
DO_LISTINGS default: true
|
DO_LISTINGS default: true
|
||||||
DO_VALUATIONS default: true
|
DO_VALUATIONS default: true
|
||||||
DO_HOUSES default: true (requires migration 071_houses_cian_zhk_url.sql on prod)
|
DO_HOUSES default: false (no cian_zhk_url col yet)
|
||||||
|
|
||||||
Examples:
|
Examples:
|
||||||
$0 --dry-run # count what would be processed
|
$0 --dry-run # count what would be processed
|
||||||
|
|
@ -46,15 +46,13 @@ while (( i < MAX_ITERS )); do
|
||||||
resp=$(curl -fsS -X POST "$url")
|
resp=$(curl -fsS -X POST "$url")
|
||||||
lt=$(echo "$resp" | jq -r '.listings_total // 0')
|
lt=$(echo "$resp" | jq -r '.listings_total // 0')
|
||||||
vt=$(echo "$resp" | jq -r '.valuations_total // 0')
|
vt=$(echo "$resp" | jq -r '.valuations_total // 0')
|
||||||
ht=$(echo "$resp" | jq -r '.houses_total // 0')
|
|
||||||
ls=$(echo "$resp" | jq -r '.listings_succeeded // 0')
|
ls=$(echo "$resp" | jq -r '.listings_succeeded // 0')
|
||||||
vs=$(echo "$resp" | jq -r '.valuations_succeeded // 0')
|
vs=$(echo "$resp" | jq -r '.valuations_succeeded // 0')
|
||||||
hs=$(echo "$resp" | jq -r '.houses_succeeded // 0')
|
|
||||||
dur=$(echo "$resp" | jq -r '.duration_sec // 0')
|
dur=$(echo "$resp" | jq -r '.duration_sec // 0')
|
||||||
cum_listings=$((cum_listings + ls))
|
cum_listings=$((cum_listings + ls))
|
||||||
cum_valuations=$((cum_valuations + vs))
|
cum_valuations=$((cum_valuations + vs))
|
||||||
echo "[$i] listings=$lt (ok=$ls) valuations=$vt (ok=$vs) houses=$ht (ok=$hs) ${dur}s cum: hist=$cum_listings val=$cum_valuations"
|
echo "[$i] listings=$lt (ok=$ls) valuations=$vt (ok=$vs) ${dur}s cum: hist=$cum_listings val=$cum_valuations"
|
||||||
if (( lt == 0 && vt == 0 && ht == 0 )); then
|
if (( lt == 0 && vt == 0 )); then
|
||||||
echo "Drained — stopping."
|
echo "Drained — stopping."
|
||||||
break
|
break
|
||||||
fi
|
fi
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue