Files
fn_registry/python/functions/pipelines/compute_centers_reachability_pipeline.py
T
egutierrez faac610745 feat: extraccion masiva footprint_aurgi (41 funcs + 4 types + stack Docker geo)
Extrae al registry funciones del proyecto interno footprint_aurgi:
- core (6): slugify_ascii, normalize_for_join, cp_provincia_es, infer_provincia_from_cp, safe_read_csv_fallback, csv_to_parquet_duckdb
- geo puras (7): haversine_km, point_in_ring, point_in_polygon, point_in_polygons_bbox, polygon_bbox, extent_with_padding, distance_bucket
- geo I/O (4): load_geojson_polygons, load_boundary_gdf, add_basemap_osm, add_basemap_with_timeout
- valhalla client (4): valhalla_route, valhalla_isochrone, valhalla_isochrones_async, valhalla_matrix_1_to_n
- datascience stats (7): trimmed_mean, geometric_mean, detect_distribution_type, best_central_tendency, summary_stats, kde_density_levels, alpha_shape_concave_hull
- datascience fuzzy (3): fuzzy_merge_adaptive (rapidfuzz), words_to_dataset, remove_words_from_column
- datascience viz (2): plot_kde_2d, plot_heatmap_log
- infra (4): compress_pdf_ghostscript, render_table_page_pdfpages, add_header_logo, osm2pgsql_ingest
- pipelines (4): setup_geo_stack_docker, compute_centers_reachability, generate_isochrones_by_zone, count_points_per_zone
- types geo (4): LonLat, BBox, IsochroneRequest, Centro

Incluye:
- apps/footprint_geo_stack/ (PostGIS + Martin + Valhalla via docker-compose)
- 131/132 tests pasan (1 skip esperado: osm2pgsql en PATH)
- Issue tracker dev/issues/0052-footprint-aurgi-extraction.md
- Atribucion uniforme: source_repo internal:footprint_aurgi, source_license internal-aurgi
- Build con 9 agentes en paralelo (8 wave 1 + 1 wave 2 pipelines)

Tambien commitea trabajo previo no commiteado: aggregate_extraction_results, chunk_with_overlap, clean_pdf_text, merge_entity_aliases, extract_graph_gliner2, extract_relations_mrebel, extract_triples_spacy_es, gliner2/mrebel/marianmt/rebel/spacy_es load_model, parse_rebel_output, translate_es_to_en, issue 0050/0051.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-04 23:35:22 +02:00

83 lines
2.8 KiB
Python

"""Pipeline: calcula matriz de tiempo/distancia y isócronas para centros de servicio."""
from __future__ import annotations
import asyncio
import sys
import os
_FUNCTIONS_DIR = os.path.join(os.path.dirname(__file__), "..")
if _FUNCTIONS_DIR not in sys.path:
sys.path.insert(0, _FUNCTIONS_DIR)
from geo.valhalla_matrix_1_to_n import valhalla_matrix_1_to_n
from geo.valhalla_isochrones_async import valhalla_isochrones_async
async def compute_centers_reachability_pipeline(
origins: list[tuple[float, float]],
centers: list[tuple[float, float]],
isochrone_minutes: int = 15,
base_url: str = "http://localhost:8002",
concurrency: int = 6,
) -> dict:
"""Calcula la accesibilidad de centros de servicio desde orígenes clientes.
Compone valhalla_matrix_1_to_n para obtener tiempos/distancias de todos
los pares (origin, center) y valhalla_isochrones_async para generar una
isócrona por cada centro.
Args:
origins: Lista de (lat, lon) de los clientes o puntos de origen.
centers: Lista de (lat, lon) de los centros de servicio.
isochrone_minutes: Minutos de isócrona a calcular para cada centro.
base_url: URL base del servidor Valhalla.
concurrency: Número máximo de requests async simultáneos para isócronas.
Returns:
Dict con claves:
"matrix" (list[dict]): Un dict por par (origin_i, center_j) con
{i, j, meters, seconds, error}. Orden: todos los centros
para origin[0], luego origin[1], etc.
"isochrones" (list[dict|None]): Una isócrona GeoJSON por cada center,
en el mismo orden. None si Valhalla falló para ese centro.
"""
n_origins = len(origins)
n_centers = len(centers)
# Build all pairs: (origin_idx, center_idx)
pairs = [(i, j) for i in range(n_origins) for j in range(n_centers)]
# 1. Matrix: origins as sources, centers as destinations
raw_matrix = valhalla_matrix_1_to_n(
origins=origins,
destinations=centers,
pairs=pairs,
base_url=base_url,
concurrency=concurrency,
)
matrix = [
{
"i": pairs[k][0],
"j": pairs[k][1],
"meters": raw_matrix[k]["meters"],
"seconds": raw_matrix[k]["seconds"],
"error": raw_matrix[k]["error"],
}
for k in range(len(pairs))
]
# 2. Isochrones: one per center
iso_requests = [
{"lat": lat, "lon": lon, "minutes": isochrone_minutes, "id": str(idx)}
for idx, (lat, lon) in enumerate(centers)
]
isochrones = await valhalla_isochrones_async(
requests=iso_requests,
base_url=base_url,
concurrency=concurrency,
)
return {"matrix": matrix, "isochrones": list(isochrones)}