443 lines
18 KiB
Python
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,
|
|
}
|