# -*- coding: utf-8 -*- """Persistencia incremental en Neon/Postgres para el scraper de Idealista. Modelo analítico v2: - Estado actual por anuncio en idealista_listings. - Snapshot por corrida en idealista_listing_snapshots. - Control de corridas en idealista_runs. - Detección de desplazamiento como desaparición rápida después de ser observado. """ from __future__ import annotations import hashlib import json import os from typing import Any import pandas as pd import psycopg from psycopg.rows import dict_row from psycopg.types.json import Jsonb DATABASE_URL_ENV_NAMES = ("DATABASE_URL_TASAS", "DATABASE_URL", "NEON_DATABASE_URL") DEFAULT_PROBABLY_RENTED_THRESHOLD = int(os.getenv("PROBABLY_RENTED_MISSING_RUNS", "3")) DEFAULT_NOT_SEEN_THRESHOLD = int(os.getenv("NOT_SEEN_RECENTLY_MISSING_RUNS", "2")) # Categorías recomendadas para una frecuencia de monitoreo cada tercer día. # Se clasifican por corridas visibles, no por días exactos, porque la # desaparición sólo se observa entre cortes de scraping. RENTAL_VELOCITY_BASIS = "visible_runs_3_day_cadence" RENTAL_VELOCITY_LABELS = { "very_fast": "Muy rápido", "fast": "Rápido", "normal": "Normal", "slow": "Lento", "very_slow": "Muy lento", "unknown": "Sin clasificar", } WATCH_FIELDS = [ "url", "title", "address_text", "location_full", "price_eur", "price_text", "price_period", "tipologia", "tipologia_text", "area_m2", "area_text", "floor_info", "listed_when", "estimated_published_at", "tag", "agency_name", "agency_url", "image_main_url", "image_main_webp", "image_count", "online_booking", "has_map_button", ] DB_COLUMNS = [ "listing_id", "district_slug", "source_input", "page_hint", "position_in_page", "global_position", "url", "title", "address_text", "location_full", "price_eur", "price_text", "price_period", "tipologia", "tipologia_text", "area_m2", "area_text", "floor_info", "listed_when", "estimated_published_at", "tag", "agency_name", "agency_url", "image_main_url", "image_main_webp", "image_count", "online_booking", "has_map_button", ] def get_database_url() -> str | None: for name in DATABASE_URL_ENV_NAMES: value = os.getenv(name) if value: return value return None def require_database_url() -> str: value = get_database_url() if not value: names = ", ".join(DATABASE_URL_ENV_NAMES) raise RuntimeError(f"No se encontró cadena de conexión Neon. Configura una de estas variables: {names}") return value def connect() -> psycopg.Connection: return psycopg.connect(require_database_url(), row_factory=dict_row) def ensure_schema(conn: psycopg.Connection) -> None: with conn.cursor() as cur: cur.execute("CREATE EXTENSION IF NOT EXISTS pgcrypto;") cur.execute( """ CREATE TABLE IF NOT EXISTS idealista_runs ( run_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), started_at TIMESTAMPTZ NOT NULL DEFAULT now(), finished_at TIMESTAMPTZ, source_filename TEXT, entries_count INTEGER DEFAULT 0, scraped_count INTEGER DEFAULT 0, inserted_count INTEGER DEFAULT 0, updated_count INTEGER DEFAULT 0, unchanged_count INTEGER DEFAULT 0, reactivated_count INTEGER DEFAULT 0, missing_updated_count INTEGER DEFAULT 0, probably_rented_count INTEGER DEFAULT 0, snapshot_count INTEGER DEFAULT 0, status TEXT NOT NULL DEFAULT 'running', error_message TEXT, districts_queried JSONB DEFAULT '[]'::jsonb, metadata JSONB DEFAULT '{}'::jsonb ); """ ) cur.execute( """ CREATE TABLE IF NOT EXISTS idealista_listings ( listing_key TEXT PRIMARY KEY, listing_id TEXT, district_slug TEXT, source_input TEXT, page_hint INTEGER, position_in_page INTEGER, global_position INTEGER, url TEXT, title TEXT, address_text TEXT, location_full TEXT, price_eur NUMERIC, price_text TEXT, price_period TEXT, tipologia INTEGER, tipologia_text TEXT, area_m2 NUMERIC, area_text TEXT, floor_info TEXT, listed_when TEXT, estimated_published_at DATE, tag TEXT, agency_name TEXT, agency_url TEXT, image_main_url TEXT, image_main_webp TEXT, image_count INTEGER, online_booking BOOLEAN, has_map_button BOOLEAN, first_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), last_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), last_run_id UUID REFERENCES idealista_runs(run_id) ON DELETE SET NULL, is_active BOOLEAN NOT NULL DEFAULT TRUE, status TEXT NOT NULL DEFAULT 'active', missing_runs INTEGER NOT NULL DEFAULT 0, visible_runs INTEGER NOT NULL DEFAULT 1, days_to_displacement INTEGER, displacement_detection_lag_days INTEGER, rental_velocity_category TEXT NOT NULL DEFAULT 'unknown', rental_velocity_basis TEXT NOT NULL DEFAULT 'visible_runs_3_day_cadence', rental_velocity_classified_at TIMESTAMPTZ, deactivated_at TIMESTAMPTZ, reactivated_at TIMESTAMPTZ, content_hash TEXT, payload JSONB DEFAULT '{}'::jsonb, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); """ ) # Migraciones seguras si existe una versión anterior. alterations = [ "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS source_filename TEXT", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS entries_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS scraped_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS inserted_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS updated_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS unchanged_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS snapshot_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS metadata JSONB DEFAULT '{}'::jsonb", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS reactivated_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS missing_updated_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS probably_rented_count INTEGER DEFAULT 0", "ALTER TABLE idealista_runs ADD COLUMN IF NOT EXISTS districts_queried JSONB DEFAULT '[]'::jsonb", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS source_input TEXT", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS position_in_page INTEGER", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS global_position INTEGER", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS estimated_published_at DATE", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS is_active BOOLEAN NOT NULL DEFAULT TRUE", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS status TEXT NOT NULL DEFAULT 'active'", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS missing_runs INTEGER NOT NULL DEFAULT 0", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS visible_runs INTEGER NOT NULL DEFAULT 1", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS days_to_displacement INTEGER", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS displacement_detection_lag_days INTEGER", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS rental_velocity_category TEXT NOT NULL DEFAULT 'unknown'", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS rental_velocity_basis TEXT NOT NULL DEFAULT 'visible_runs_3_day_cadence'", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS rental_velocity_classified_at TIMESTAMPTZ", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS deactivated_at TIMESTAMPTZ", "ALTER TABLE idealista_listings ADD COLUMN IF NOT EXISTS reactivated_at TIMESTAMPTZ", ] for sql in alterations: cur.execute(sql) cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_listing_id ON idealista_listings(listing_id);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_district ON idealista_listings(district_slug);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_price ON idealista_listings(price_eur);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_last_seen ON idealista_listings(last_seen_at);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_status ON idealista_listings(status);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_missing ON idealista_listings(missing_runs);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_velocity ON idealista_listings(rental_velocity_category);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_listings_visible_runs ON idealista_listings(visible_runs);") cur.execute( """ CREATE TABLE IF NOT EXISTS idealista_listing_snapshots ( snapshot_id BIGSERIAL PRIMARY KEY, run_id UUID REFERENCES idealista_runs(run_id) ON DELETE CASCADE, listing_key TEXT NOT NULL, listing_id TEXT, district_slug TEXT, price_eur NUMERIC, area_m2 NUMERIC, title TEXT, url TEXT, page_hint INTEGER, position_in_page INTEGER, global_position INTEGER, listed_when TEXT, estimated_published_at DATE, content_hash TEXT, observed_at TIMESTAMPTZ NOT NULL DEFAULT now(), payload JSONB DEFAULT '{}'::jsonb, UNIQUE(run_id, listing_key) ); """ ) snapshot_alterations = [ "ALTER TABLE idealista_listing_snapshots ADD COLUMN IF NOT EXISTS page_hint INTEGER", "ALTER TABLE idealista_listing_snapshots ADD COLUMN IF NOT EXISTS position_in_page INTEGER", "ALTER TABLE idealista_listing_snapshots ADD COLUMN IF NOT EXISTS global_position INTEGER", "ALTER TABLE idealista_listing_snapshots ADD COLUMN IF NOT EXISTS listed_when TEXT", "ALTER TABLE idealista_listing_snapshots ADD COLUMN IF NOT EXISTS estimated_published_at DATE", ] for sql in snapshot_alterations: cur.execute(sql) cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_snapshots_listing_key ON idealista_listing_snapshots(listing_key);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_snapshots_run ON idealista_listing_snapshots(run_id);") cur.execute("CREATE INDEX IF NOT EXISTS idx_idealista_snapshots_observed ON idealista_listing_snapshots(observed_at);") conn.commit() def _none_if_nan(value: Any) -> Any: if value is None: return None try: if pd.isna(value): return None except Exception: pass if hasattr(value, "item"): try: return value.item() except Exception: pass return value def normalize_record(row: dict[str, Any]) -> dict[str, Any]: out = {k: _none_if_nan(v) for k, v in row.items()} for k in ["page_hint", "position_in_page", "global_position", "tipologia", "image_count"]: if out.get(k) is not None: out[k] = int(out[k]) for k in ["price_eur", "area_m2"]: if out.get(k) is not None: out[k] = float(out[k]) return out def listing_key_for(row: dict[str, Any]) -> str | None: listing_id = row.get("listing_id") url = row.get("url") if listing_id: return str(listing_id) if url: return str(url) return None def content_hash(row: dict[str, Any]) -> str: payload = {field: row.get(field) for field in WATCH_FIELDS} raw = json.dumps(payload, ensure_ascii=False, sort_keys=True, default=str) return hashlib.sha256(raw.encode("utf-8")).hexdigest() def rental_velocity_category_from_visible_runs(visible_runs: int | None) -> str: """Clasifica velocidad de arrendamiento inferido para corrida cada tercer día. Criterio: - 1 corrida visible: very_fast - 2 a 3 corridas visibles: fast - 4 a 6 corridas visibles: normal - 7 a 10 corridas visibles: slow - 11+ corridas visibles: very_slow """ if visible_runs is None: return "unknown" try: n = int(visible_runs) except Exception: return "unknown" if n <= 0: return "unknown" if n == 1: return "very_fast" if 2 <= n <= 3: return "fast" if 4 <= n <= 6: return "normal" if 7 <= n <= 10: return "slow" return "very_slow" def create_run( conn: psycopg.Connection, source_filename: str, entries_count: int, districts_queried: list[str] | None = None, metadata: dict | None = None, ) -> str: ensure_schema(conn) with conn.cursor() as cur: cur.execute( """ INSERT INTO idealista_runs (source_filename, entries_count, districts_queried, metadata) VALUES (%s, %s, %s, %s) RETURNING run_id; """, (source_filename, entries_count, Jsonb(districts_queried or []), Jsonb(metadata or {})), ) row = cur.fetchone() conn.commit() return str(row["run_id"]) def finish_run(conn: psycopg.Connection, run_id: str, stats: dict[str, int], status: str = "success", error_message: str | None = None) -> None: with conn.cursor() as cur: cur.execute( """ UPDATE idealista_runs SET finished_at = now(), scraped_count = %s, inserted_count = %s, updated_count = %s, unchanged_count = %s, reactivated_count = %s, missing_updated_count = %s, probably_rented_count = %s, snapshot_count = %s, status = %s, error_message = %s WHERE run_id = %s; """, ( stats.get("scraped", 0), stats.get("inserted", 0), stats.get("updated", 0), stats.get("unchanged", 0), stats.get("reactivated", 0), stats.get("missing_updated", 0), stats.get("probably_rented", 0), stats.get("snapshots", 0), status, error_message, run_id, ), ) conn.commit() def _insert_snapshot(cur: psycopg.Cursor, row: dict[str, Any], key: str, run_id: str, h: str) -> bool: cur.execute( """ INSERT INTO idealista_listing_snapshots ( run_id, listing_key, listing_id, district_slug, price_eur, area_m2, title, url, page_hint, position_in_page, global_position, listed_when, estimated_published_at, content_hash, payload ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) ON CONFLICT (run_id, listing_key) DO NOTHING; """, ( run_id, key, row.get("listing_id"), row.get("district_slug"), row.get("price_eur"), row.get("area_m2"), row.get("title"), row.get("url"), row.get("page_hint"), row.get("position_in_page"), row.get("global_position"), row.get("listed_when"), row.get("estimated_published_at"), h, Jsonb(row), ), ) return cur.rowcount > 0 def upsert_listings( conn: psycopg.Connection, df: pd.DataFrame, run_id: str, districts_queried: list[str] | None = None, probably_rented_threshold: int = DEFAULT_PROBABLY_RENTED_THRESHOLD, not_seen_threshold: int = DEFAULT_NOT_SEEN_THRESHOLD, ) -> dict[str, int]: ensure_schema(conn) stats = { "scraped": int(len(df)), "inserted": 0, "updated": 0, "unchanged": 0, "reactivated": 0, "missing_updated": 0, "probably_rented": 0, "snapshots": 0, } observed_keys: set[str] = set() with conn.cursor() as cur: for _, raw_row in df.iterrows(): row = normalize_record(raw_row.to_dict()) key = listing_key_for(row) if not key: continue # Evita sumar visible_runs más de una vez en la misma corrida si el mismo # anuncio aparece duplicado por solapamiento de distritos o paginación. if key in observed_keys: continue observed_keys.add(key) h = content_hash(row) cur.execute("SELECT content_hash, status, is_active FROM idealista_listings WHERE listing_key = %s", (key,)) existing = cur.fetchone() was_reactivated = bool(existing and existing.get("status") in {"temporarily_missing", "not_seen_recently", "probably_rented", "possibly_removed"}) if existing is None: cur.execute( """ INSERT INTO idealista_listings ( listing_key, listing_id, district_slug, source_input, page_hint, position_in_page, global_position, url, title, address_text, location_full, price_eur, price_text, price_period, tipologia, tipologia_text, area_m2, area_text, floor_info, listed_when, estimated_published_at, tag, agency_name, agency_url, image_main_url, image_main_webp, image_count, online_booking, has_map_button, last_run_id, is_active, status, missing_runs, visible_runs, days_to_displacement, displacement_detection_lag_days, rental_velocity_category, rental_velocity_basis, rental_velocity_classified_at, content_hash, payload ) VALUES ( %(listing_key)s, %(listing_id)s, %(district_slug)s, %(source_input)s, %(page_hint)s, %(position_in_page)s, %(global_position)s, %(url)s, %(title)s, %(address_text)s, %(location_full)s, %(price_eur)s, %(price_text)s, %(price_period)s, %(tipologia)s, %(tipologia_text)s, %(area_m2)s, %(area_text)s, %(floor_info)s, %(listed_when)s, %(estimated_published_at)s, %(tag)s, %(agency_name)s, %(agency_url)s, %(image_main_url)s, %(image_main_webp)s, %(image_count)s, %(online_booking)s, %(has_map_button)s, %(last_run_id)s, TRUE, 'active', 0, 1, NULL, NULL, 'unknown', 'visible_runs_3_day_cadence', NULL, %(content_hash)s, %(payload)s ); """, {**{c: row.get(c) for c in DB_COLUMNS}, "listing_key": key, "last_run_id": run_id, "content_hash": h, "payload": Jsonb(row)}, ) stats["inserted"] += 1 else: changed = existing.get("content_hash") != h set_status = "active" cur.execute( """ UPDATE idealista_listings SET listing_id = %(listing_id)s, district_slug = %(district_slug)s, source_input = %(source_input)s, page_hint = %(page_hint)s, position_in_page = %(position_in_page)s, global_position = %(global_position)s, url = %(url)s, title = %(title)s, address_text = %(address_text)s, location_full = %(location_full)s, price_eur = %(price_eur)s, price_text = %(price_text)s, price_period = %(price_period)s, tipologia = %(tipologia)s, tipologia_text = %(tipologia_text)s, area_m2 = %(area_m2)s, area_text = %(area_text)s, floor_info = %(floor_info)s, listed_when = %(listed_when)s, estimated_published_at = %(estimated_published_at)s, tag = %(tag)s, agency_name = %(agency_name)s, agency_url = %(agency_url)s, image_main_url = %(image_main_url)s, image_main_webp = %(image_main_webp)s, image_count = %(image_count)s, online_booking = %(online_booking)s, has_map_button = %(has_map_button)s, last_seen_at = now(), last_run_id = %(last_run_id)s, is_active = TRUE, status = %(status)s, missing_runs = 0, visible_runs = COALESCE(visible_runs, 0) + 1, days_to_displacement = NULL, displacement_detection_lag_days = NULL, rental_velocity_category = 'unknown', rental_velocity_basis = 'visible_runs_3_day_cadence', rental_velocity_classified_at = NULL, deactivated_at = NULL, reactivated_at = CASE WHEN %(was_reactivated)s THEN now() ELSE reactivated_at END, content_hash = %(content_hash)s, payload = %(payload)s, updated_at = now() WHERE listing_key = %(listing_key)s; """, { **{c: row.get(c) for c in DB_COLUMNS}, "listing_key": key, "last_run_id": run_id, "status": set_status, "was_reactivated": was_reactivated, "content_hash": h, "payload": Jsonb(row), }, ) if was_reactivated: stats["reactivated"] += 1 if changed or was_reactivated: stats["updated"] += 1 else: stats["unchanged"] += 1 if _insert_snapshot(cur, row, key, run_id, h): stats["snapshots"] += 1 if districts_queried: stats.update(_mark_missing(cur, observed_keys, districts_queried, probably_rented_threshold, not_seen_threshold)) conn.commit() return stats def _mark_missing( cur: psycopg.Cursor, observed_keys: set[str], districts_queried: list[str], probably_rented_threshold: int, not_seen_threshold: int, ) -> dict[str, int]: stats = {"missing_updated": 0, "probably_rented": 0} district_list = sorted(set([d for d in districts_queried if d])) if not district_list: return stats # PostgreSQL usa <> ALL(array) para excluir observados. Si observed_keys está vacío, # usamos una lista imposible para evitar SQL dinámico peligroso. observed_list = list(observed_keys) or ["__NO_OBSERVED_KEYS__"] cur.execute( """ WITH candidates AS ( SELECT listing_key, missing_runs + 1 AS next_missing_runs, COALESCE(visible_runs, 0) AS visible_runs, status FROM idealista_listings WHERE district_slug = ANY(%s) AND listing_key <> ALL(%s) AND status <> 'probably_rented' ), classified AS ( SELECT listing_key, next_missing_runs, visible_runs, CASE WHEN next_missing_runs >= %s THEN 'probably_rented' WHEN next_missing_runs >= %s THEN 'not_seen_recently' ELSE 'temporarily_missing' END AS next_status, CASE WHEN visible_runs <= 0 THEN 'unknown' WHEN visible_runs = 1 THEN 'very_fast' WHEN visible_runs BETWEEN 2 AND 3 THEN 'fast' WHEN visible_runs BETWEEN 4 AND 6 THEN 'normal' WHEN visible_runs BETWEEN 7 AND 10 THEN 'slow' ELSE 'very_slow' END AS next_velocity_category FROM candidates ), updated AS ( UPDATE idealista_listings l SET missing_runs = c.next_missing_runs, is_active = FALSE, status = c.next_status, deactivated_at = COALESCE(l.deactivated_at, now()), days_to_displacement = CASE WHEN c.next_status = 'probably_rented' THEN GREATEST(0, EXTRACT(DAY FROM (l.last_seen_at - l.first_seen_at))::int) ELSE l.days_to_displacement END, displacement_detection_lag_days = CASE WHEN c.next_status = 'probably_rented' THEN GREATEST(0, EXTRACT(DAY FROM (now() - l.last_seen_at))::int) ELSE l.displacement_detection_lag_days END, rental_velocity_category = CASE WHEN c.next_status = 'probably_rented' THEN c.next_velocity_category ELSE l.rental_velocity_category END, rental_velocity_basis = CASE WHEN c.next_status = 'probably_rented' THEN 'visible_runs_3_day_cadence' ELSE l.rental_velocity_basis END, rental_velocity_classified_at = CASE WHEN c.next_status = 'probably_rented' THEN COALESCE(l.rental_velocity_classified_at, now()) ELSE l.rental_velocity_classified_at END, updated_at = now() FROM classified c WHERE l.listing_key = c.listing_key RETURNING l.status ) SELECT COUNT(*)::int AS missing_updated, COUNT(*) FILTER (WHERE status = 'probably_rented')::int AS probably_rented FROM updated; """, (district_list, observed_list, probably_rented_threshold, not_seen_threshold), ) row = cur.fetchone() or {} stats["missing_updated"] = int(row.get("missing_updated") or 0) stats["probably_rented"] = int(row.get("probably_rented") or 0) return stats def _remove_timezone_if_datetime(value): """Convierte datetime con timezone a datetime sin timezone para exportar a Excel.""" if hasattr(value, "tzinfo") and value.tzinfo is not None: try: return value.replace(tzinfo=None) except Exception: return value return value def read_dataframe(conn: psycopg.Connection, sql: str, params: tuple | None = None) -> pd.DataFrame: """Lee una consulta SQL en DataFrame usando cursor psycopg3 y limpia timezones para Excel.""" with conn.cursor() as cur: cur.execute(sql, params or ()) rows = cur.fetchall() columns = [desc.name for desc in cur.description] if cur.description else [] df = pd.DataFrame(rows, columns=columns) # Excel no soporta datetimes con timezone. # Convertimos cualquier columna datetime tz-aware a datetime sin timezone. for col in df.columns: if pd.api.types.is_datetime64tz_dtype(df[col]): df[col] = df[col].dt.tz_localize(None) elif df[col].dtype == "object": df[col] = df[col].apply(_remove_timezone_if_datetime) return df def write_sheet(writer, conn, sql: str, sheet_name: str): df = read_dataframe(conn, sql) if df.empty: df = pd.DataFrame({"mensaje": ["Sin datos disponibles"]}) df.to_excel(writer, sheet_name=sheet_name, index=False) def export_incremental_excel(conn: psycopg.Connection, output_xlsx: str) -> str: """Exporta workbook analítico completo desde Neon.""" base_sql = """ SELECT listing_key, listing_id, district_slug, status, is_active, missing_runs, visible_runs, rental_velocity_category, CASE rental_velocity_category WHEN 'very_fast' THEN 'Muy rápido' WHEN 'fast' THEN 'Rápido' WHEN 'normal' THEN 'Normal' WHEN 'slow' THEN 'Lento' WHEN 'very_slow' THEN 'Muy lento' ELSE 'Sin clasificar' END AS rental_velocity_label, rental_velocity_basis, days_to_displacement, displacement_detection_lag_days, rental_velocity_classified_at, first_seen_at, last_seen_at, deactivated_at, reactivated_at, EXTRACT(DAY FROM (COALESCE(deactivated_at, last_seen_at, now()) - first_seen_at))::int AS days_visible, estimated_published_at, CASE WHEN estimated_published_at IS NOT NULL THEN (CURRENT_DATE - estimated_published_at)::int END AS estimated_days_on_portal, price_eur, area_m2, CASE WHEN area_m2 IS NOT NULL AND area_m2 <> 0 THEN ROUND((price_eur / area_m2)::numeric, 2) END AS eur_m2, page_hint, position_in_page, global_position, title, url, address_text, location_full, price_text, price_period, tipologia, tipologia_text, agency_name, agency_url, listed_when, tag, source_input, last_run_id, updated_at FROM idealista_listings ORDER BY district_slug, price_eur NULLS LAST; """ mayor_desplazamiento_sql = """ SELECT listing_key, listing_id, district_slug, status, rental_velocity_category, CASE rental_velocity_category WHEN 'very_fast' THEN 'Muy rápido' WHEN 'fast' THEN 'Rápido' WHEN 'normal' THEN 'Normal' WHEN 'slow' THEN 'Lento' WHEN 'very_slow' THEN 'Muy lento' ELSE 'Sin clasificar' END AS rental_velocity_label, visible_runs, missing_runs, first_seen_at, last_seen_at, deactivated_at, days_to_displacement, displacement_detection_lag_days, price_eur, area_m2, title, url, address_text, location_full, listed_when, estimated_published_at FROM idealista_listings WHERE status = 'probably_rented' ORDER BY visible_runs ASC NULLS LAST, days_to_displacement ASC NULLS LAST, displacement_detection_lag_days ASC NULLS LAST, deactivated_at DESC NULLS LAST LIMIT 500; """ mayor_antiguedad_sql = """ SELECT listing_key, listing_id, district_slug, status, first_seen_at, last_seen_at, EXTRACT(DAY FROM (now() - first_seen_at))::int AS days_in_base, price_eur, area_m2, title, url, address_text, location_full FROM idealista_listings ORDER BY first_seen_at ASC LIMIT 500; """ menor_antiguedad_sql = """ SELECT listing_key, listing_id, district_slug, status, first_seen_at, last_seen_at, EXTRACT(DAY FROM (now() - first_seen_at))::int AS days_in_base, price_eur, area_m2, title, url, address_text, location_full FROM idealista_listings ORDER BY first_seen_at DESC LIMIT 500; """ cambios_precio_sql = """ WITH ordered AS ( SELECT listing_key, district_slug, (ARRAY_AGG(price_eur ORDER BY observed_at ASC))[1] AS first_price_eur, (ARRAY_AGG(price_eur ORDER BY observed_at DESC))[1] AS last_price_eur, MIN(price_eur) AS min_price_eur, MAX(price_eur) AS max_price_eur, MIN(observed_at) AS first_observed_at, MAX(observed_at) AS last_observed_at, COUNT(*) AS observations FROM idealista_listing_snapshots WHERE price_eur IS NOT NULL GROUP BY listing_key, district_slug ) SELECT o.*, (last_price_eur - first_price_eur) AS change_abs_eur, ROUND(((last_price_eur - first_price_eur) / NULLIF(first_price_eur, 0)) * 100, 2) AS change_pct FROM ordered o WHERE observations > 1 AND first_price_eur IS DISTINCT FROM last_price_eur ORDER BY change_abs_eur ASC NULLS LAST; """ categorias_sql = """ SELECT rental_velocity_category, CASE rental_velocity_category WHEN 'very_fast' THEN 'Muy rápido' WHEN 'fast' THEN 'Rápido' WHEN 'normal' THEN 'Normal' WHEN 'slow' THEN 'Lento' WHEN 'very_slow' THEN 'Muy lento' ELSE 'Sin clasificar' END AS rental_velocity_label, COUNT(*)::int AS listings_count, ROUND(AVG(visible_runs)::numeric, 2) AS avg_visible_runs, ROUND(AVG(days_to_displacement)::numeric, 2) AS avg_days_to_displacement, ROUND(AVG(price_eur)::numeric, 2) AS avg_price_eur FROM idealista_listings WHERE status = 'probably_rented' GROUP BY rental_velocity_category ORDER BY MIN(CASE rental_velocity_category WHEN 'very_fast' THEN 1 WHEN 'fast' THEN 2 WHEN 'normal' THEN 3 WHEN 'slow' THEN 4 WHEN 'very_slow' THEN 5 ELSE 6 END); """ runs_sql = """ SELECT run_id, started_at, finished_at, source_filename, entries_count, scraped_count, inserted_count, updated_count, unchanged_count, reactivated_count, missing_updated_count, probably_rented_count, snapshot_count, status, districts_queried, error_message FROM idealista_runs ORDER BY run_id DESC LIMIT 200; """ with pd.ExcelWriter(output_xlsx, engine="openpyxl") as writer: write_sheet(writer, conn, base_sql, "Base incremental") write_sheet(writer, conn, mayor_desplazamiento_sql, "Mayor desplazamiento") write_sheet(writer, conn, categorias_sql, "Categorias velocidad") write_sheet(writer, conn, mayor_antiguedad_sql, "Mayor antiguedad") write_sheet(writer, conn, menor_antiguedad_sql, "Menor antiguedad") write_sheet(writer, conn, cambios_precio_sql, "Cambios precio") write_sheet(writer, conn, runs_sql, "Corridas")