# 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, }