loko/streetup/schools/geocoding.py
2026-07-22 14:48:40 +02:00

443 lines
18 KiB
Python

# geocode_photon.py
# ------------------------------------------------------------
# Geocode a list of addresses from an Excel file using Photon.
#
# Based on the original geocode_photon.py script provided.
# Adapted for use within Django application.
# Now with internal database geocoding using fuzzy matching as first attempt,
# with fallback to Photon API.
# ------------------------------------------------------------
import asyncio
import aiohttp
import math
import time
import logging
from typing import Any, Dict, Optional, Tuple, List, Union
import pandas as pd
from .internal_geocoding import geocode_internal, geocode_internal_batch, DEFAULT_SCORE_THRESHOLD
logger = logging.getLogger(__name__)
PHOTON_URL = "https://photon.komoot.io/api/"
def chunked(seq: List[Any], size: int):
for i in range(0, len(seq), size):
yield seq[i:i+size]
def _preview(q: str, n: int = 80) -> str:
q = q.replace("\n", " ").replace("\r", " ").strip()
return (q if len(q) <= n else q[: n - 1] + "…")
def _format_found_address(properties: Dict[str, Any]) -> str:
"""Format the address from Photon response properties into a readable string."""
parts = []
# House number and street
street = properties.get("street", "")
housenumber = properties.get("housenumber", "")
if street:
if housenumber:
parts.append(f"{street} {housenumber}")
else:
parts.append(street)
elif properties.get("name"):
parts.append(properties.get("name"))
# Postcode and city
postcode = properties.get("postcode", "")
city = properties.get("city", "") or properties.get("locality", "") or properties.get("district", "")
if postcode and city:
parts.append(f"{postcode} {city}")
elif city:
parts.append(city)
elif postcode:
parts.append(postcode)
# State (if different from city)
state = properties.get("state", "")
if state and state != city:
parts.append(state)
# Country
country = properties.get("country", "")
if country:
parts.append(country)
return ", ".join(parts) if parts else properties.get("name", "Unknown")
async def fetch_photon(
session: aiohttp.ClientSession,
query: str,
limit: int = 1,
timeout_s: int = 10,
max_retries: int = 4,
) -> Tuple[Optional[float], Optional[float], str]:
"""Call Photon for a single query and return (lon, lat, found_address_or_error).
Returns:
Tuple of (longitude, latitude, found_address_or_error)
- On success: (lon, lat, formatted_address)
- On not found: (None, None, "Not found")
- On error: (None, None, error_description)
"""
params = {"q": query, "limit": str(limit)}
attempts = 0
last_http: Optional[int] = None
last_error: Optional[str] = None
t0 = time.perf_counter()
backoff = 1.0
for attempt in range(max_retries):
attempts += 1
try:
async with session.get(PHOTON_URL, params=params, timeout=timeout_s) as resp:
last_http = resp.status
# Retry on rate limit or server errors
if resp.status in (429, 500, 502, 503, 504):
await resp.read() # drain
await asyncio.sleep(backoff)
backoff = min(backoff * 2, 16)
continue
resp.raise_for_status()
data = await resp.json()
features = data.get("features") or []
if not features:
elapsed = time.perf_counter() - t0
logger.debug(f"[req] '{_preview(query)}' -> not_found (http={last_http}) "
f"in {elapsed:.2f}s, attempts={attempts}")
return (None, None, "Not found")
# Take the first feature
feature = features[0]
geom = feature.get("geometry") or {}
properties = feature.get("properties") or {}
coords = geom.get("coordinates") or []
if isinstance(coords, list) and len(coords) >= 2:
lon, lat = coords[0], coords[1]
if all(isinstance(v, (int, float)) and math.isfinite(v) for v in (lon, lat)):
found_address = _format_found_address(properties)
elapsed = time.perf_counter() - t0
logger.debug(f"[req] '{_preview(query)}' -> success lon={lon}, lat={lat} "
f"(http={last_http}) in {elapsed:.2f}s, attempts={attempts}")
return (float(lon), float(lat), found_address)
# If we reached here the response is 200 but unusable
elapsed = time.perf_counter() - t0
logger.warning(f"[req] '{_preview(query)}' -> invalid_geometry (http={last_http}) "
f"in {elapsed:.2f}s, attempts={attempts}")
return (None, None, "Invalid geometry in response")
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
last_error = e.__class__.__name__
await asyncio.sleep(backoff)
backoff = min(backoff * 2, 16)
# loop to retry
# Exhausted retries
elapsed = time.perf_counter() - t0
if last_error:
logger.warning(f"[req] '{_preview(query)}' -> error {last_error} "
f"after {elapsed:.2f}s, attempts={attempts}")
return (None, None, f"Error: {last_error}")
else:
logger.warning(f"[req] '{_preview(query)}' -> failed (http={last_http}) "
f"after {elapsed:.2f}s, attempts={attempts}")
return (None, None, f"Failed (HTTP {last_http})")
async def fetch_with_internal_fallback(
session: aiohttp.ClientSession,
query: str,
limit: int = 1,
timeout_s: int = 10,
max_retries: int = 4,
internal_threshold: float = DEFAULT_SCORE_THRESHOLD,
use_internal: bool = True,
) -> Tuple[Optional[float], Optional[float], str]:
"""
Try to geocode using internal database first, then fall back to Photon API.
Args:
session: aiohttp session for Photon API calls
query: Address string to geocode
limit: Photon limit parameter
timeout_s: Timeout for Photon API calls
max_retries: Max retries for Photon API calls
internal_threshold: Score threshold for internal geocoding (0-100)
use_internal: Whether to try internal DB first (default: True)
Returns:
Tuple of (longitude, latitude, found_address_or_source)
The found_address includes source info: "[Internal] address" or "[Photon] address"
"""
if not query or not query.strip():
return (None, None, "Empty address")
# Try internal geocoding first if enabled
if use_internal:
try:
# Run synchronous database query in a thread pool to avoid blocking the event loop
loop = asyncio.get_event_loop()
lon, lat, found_addr, score = await loop.run_in_executor(
None, geocode_internal, query, internal_threshold
)
if lon is not None and lat is not None:
# Success with internal DB
logger.debug(f"Internal DB success for '{_preview(query)}' (score: {score:.1f})")
return (lon, lat, f"[URBIS] {found_addr}")
else:
logger.debug(f"Internal DB no match for '{_preview(query)}', falling back to Photon")
except Exception as e:
logger.warning(f"Internal geocoding error for '{_preview(query)}': {e}, falling back to Photon")
# Fall back to Photon API
lon, lat, found_addr = await fetch_photon(session, query, limit, timeout_s, max_retries)
if lon is not None and lat is not None:
return (lon, lat, f"[Photon API] {found_addr}")
return (lon, lat, found_addr)
async def geocode_block(
session: aiohttp.ClientSession,
rows: List[Tuple[int, str]],
limit: int = 1,
internal_threshold: float = DEFAULT_SCORE_THRESHOLD,
use_internal: bool = True,
internal_concurrency: int = 5,
) -> Dict[int, Tuple[Optional[float], Optional[float], str]]:
"""Geocode a block (list) of (row_index, address) concurrently.
Uses a two-phase approach:
1. First, try internal DB geocoding in parallel for all addresses
2. Then, fallback to Photon API only for addresses that failed internal geocoding
Args:
session: aiohttp session for API calls
rows: List of (row_index, address) tuples
limit: Photon API limit parameter
internal_threshold: Score threshold for internal geocoding
use_internal: Whether to try internal DB first
internal_concurrency: Max concurrent internal DB queries (default 5)
Returns:
Dict mapping row_index to (lon, lat, found_address_or_error)
"""
out: Dict[int, Tuple[Optional[float], Optional[float], str]] = {}
if not rows:
return out
# Phase 1: Try internal geocoding for all addresses in parallel
need_photon: List[Tuple[int, str]] = []
if use_internal:
try:
internal_results = await geocode_internal_batch(
rows,
threshold=internal_threshold,
max_concurrent=internal_concurrency
)
for idx, addr in rows:
result = internal_results.get(idx)
if result:
lon, lat, found_addr, score = result
if lon is not None and lat is not None:
# Success with internal DB
out[idx] = (lon, lat, f"[URBIS] {found_addr}")
logger.debug(f"Internal DB success for idx={idx} (score: {score:.2f})")
else:
# Need to try Photon
need_photon.append((idx, addr))
else:
need_photon.append((idx, addr))
except Exception as e:
logger.warning(f"Internal batch geocoding error: {e}, falling back to Photon for all")
need_photon = list(rows)
else:
need_photon = list(rows)
# Phase 2: Fallback to Photon API for addresses that failed internal geocoding
if need_photon:
logger.debug(f"Falling back to Photon for {len(need_photon)} addresses")
photon_tasks = [
fetch_photon(session, addr, limit=limit)
for (idx, addr) in need_photon
]
photon_results = await asyncio.gather(*photon_tasks, return_exceptions=True)
for (idx, _), res in zip(need_photon, photon_results):
if isinstance(res, Exception):
out[idx] = (None, None, f"Error: {res.__class__.__name__}")
elif res is None:
out[idx] = (None, None, "Unknown error")
else:
lon, lat, found_addr = res
if lon is not None and lat is not None:
out[idx] = (lon, lat, f"[Photon API] {found_addr}")
else:
out[idx] = (lon, lat, found_addr)
return out
async def geocode_all(
addresses: List[str],
chunk_size: int = 5,
limit: int = 1,
internal_threshold: float = DEFAULT_SCORE_THRESHOLD,
use_internal: bool = True,
) -> List[Tuple[Optional[float], Optional[float], str]]:
"""Geocode all addresses in blocks of `chunk_size`.
First tries internal database with fuzzy matching, then falls back to Photon API.
Args:
addresses: List of address strings to geocode
chunk_size: Number of concurrent requests per block
limit: Photon API limit parameter
internal_threshold: Score threshold for internal geocoding (0-100)
use_internal: Whether to try internal DB first (default: True)
Returns:
List of (lon, lat, found_address_or_error) aligned with input order."""
connector = aiohttp.TCPConnector(limit=chunk_size) # matches concurrency per block
headers = {
"User-Agent": "streetup-geocoder/1.0 (+https://photon.komoot.io/)",
"Accept": "application/json",
}
total = len(addresses)
total_blocks = (total + max(1, chunk_size) - 1) // max(1, chunk_size)
async with aiohttp.ClientSession(connector=connector, headers=headers) as session:
results: Dict[int, Tuple[Optional[float], Optional[float], str]] = {}
# Prepare indexed rows
indexed = [(i, (a or "").strip()) for i, a in enumerate(addresses)]
for block_idx, block in enumerate(chunked(indexed, chunk_size), start=1):
t0 = time.perf_counter()
# Skip empty addresses quickly
non_empty = [(i, a) for i, a in block if a]
empty_count = len(block) - len(non_empty)
if not non_empty:
for i, _ in block:
results[i] = (None, None, "Empty address")
took = time.perf_counter() - t0
logger.info(f"[{block_idx}/{total_blocks}] Skipped empty block of {len(block)} rows (all empty). Took {took:.2f}s.")
continue
block_results = await geocode_block(
session, non_empty, limit=limit,
internal_threshold=internal_threshold,
use_internal=use_internal
)
for i, a in block:
if not a:
results[i] = (None, None, "Empty address")
else:
results[i] = block_results.get(i, (None, None, "Unknown error"))
successes = sum(1 for i, _ in non_empty if results.get(i, (None, None, ""))[0] is not None)
failures = len(non_empty) - successes
took = time.perf_counter() - t0
logger.info(f"[{block_idx}/{total_blocks}] Completed block of {len(block)} rows: {successes} success, {failures} not found/errors, {empty_count} empty. Took {took:.2f}s.")
return [results[i] for i in range(len(addresses))]
def resolve_address_series(df: pd.DataFrame, address_col: Union[int, str]) -> pd.Series:
if isinstance(address_col, int):
if address_col < 0 or address_col >= df.shape[1]:
raise IndexError(f"Address column index {address_col} out of range.")
return df.iloc[:, address_col]
else:
if address_col not in df.columns:
raise KeyError(f"Column '{address_col}' not found. Available: {list(df.columns)}")
return df[address_col]
def geocode_excel_file(
input_path: str,
output_path: str,
sheet_name: Optional[str] = None,
address_col: Union[int, str] = 0,
chunk_size: int = 3,
limit: int = 1,
output_sheet: str = "geocoded",
internal_threshold: float = DEFAULT_SCORE_THRESHOLD,
use_internal: bool = True,
) -> dict:
"""
Geocode addresses from an Excel file and write results to a new Excel file.
First tries internal database with fuzzy matching in both FR and NL,
then falls back to Photon API if score is below threshold.
Args:
input_path: Path to input Excel file
output_path: Path to output Excel file
sheet_name: Sheet name to read (None for first sheet)
address_col: Address column (index like 0 or column name)
chunk_size: Concurrent requests per block (default 3)
limit: Photon 'limit' parameter (default 1)
output_sheet: Name of the output sheet (default "geocoded")
internal_threshold: Score threshold for internal geocoding (0-100, default 80)
use_internal: Whether to try internal DB first (default True)
Returns:
dict: Statistics about the geocoding operation
"""
# Load Excel
if sheet_name:
df = pd.read_excel(input_path, sheet_name=sheet_name)
else:
df = pd.read_excel(input_path)
# Resolve address column
try:
# Try to parse address-col as int; if fails use as name
try:
addr_col: Union[int, str] = int(address_col)
except (ValueError, TypeError):
addr_col = str(address_col)
s = resolve_address_series(df, addr_col)
except (IndexError, KeyError) as e:
raise ValueError(f"Error: {e}")
addresses = s.astype(str).fillna("").tolist()
logger.info(f"Starting geocoding of {len(addresses)} rows in blocks of {chunk_size}...")
t0_all = time.perf_counter()
# Run async geocoding in blocks
results = asyncio.run(geocode_all(
addresses, chunk_size=chunk_size, limit=limit,
internal_threshold=internal_threshold,
use_internal=use_internal
))
lons = [r[0] for r in results]
lats = [r[1] for r in results]
found_addresses = [r[2] for r in results]
took_all = time.perf_counter() - t0_all
total = len(addresses)
total_ok = sum(1 for r in results if r[0] is not None and r[1] is not None)
total_empty = sum(1 for a in addresses if not a.strip())
total_fail = total - total_ok - total_empty
logger.info(f"Finished geocoding: {total_ok} success, {total_fail} not found/errors, {total_empty} empty. Took {took_all:.2f}s.")
# Attach results
df_out = df.copy()
df_out["longitude"] = lons
df_out["latitude"] = lats
df_out["found_address"] = found_addresses
# Write to Excel
with pd.ExcelWriter(output_path, engine="openpyxl") as writer:
df_out.to_excel(writer, index=False, sheet_name=output_sheet)
logger.info(f"Done. Wrote results to sheet '{output_sheet}' in {output_path}.")
return {
'total': total,
'success': total_ok,
'failed': total_fail,
'empty': total_empty,
'duration': took_all,
}