"""Bounded Overture GeoParquet reads that share one S3 client and one STAC index.
``overturemaps.record_batch_reader`` downloads the release's STAC index and opens a
new S3 client on every call, and it reads every column of the matching row groups.
The compiler issues six such reads per scene, so this module keeps the client for
the life of the process, keeps the index of each release in memory and on disk
(releases are immutable), and projects each read to the columns the providers
decode. The row filter is the same bounding-box predicate. The index is the only
way objects are selected; when it cannot be fetched after a few retries, the query
fails and the calling provider reports that.
"""
from __future__ import annotations
import os
import threading
import time
from pathlib import Path
from urllib.request import urlopen
import pyarrow as pa
import pyarrow.compute as pc
import pyarrow.dataset as ds
import pyarrow.fs as fs
import pyarrow.parquet as pq
from overturemaps.core import get_latest_release
STAC_COLLECTIONS_URL = "https://stac.overturemaps.org/{release}/collections.parquet"
S3_REGION = "us-west-2"
# The index server throttles bursts of requests; these pauses precede each retry.
STAC_RETRY_DELAYS_S = (0.5, 1.0, 2.0)
BUILDING_COLUMNS = (
"id",
"geometry",
"is_underground",
"height",
"min_height",
"num_floors",
"min_floor",
"sources",
"names",
)
BUILDING_PART_COLUMNS = (*BUILDING_COLUMNS, "building_id")
SURFACE_COLUMNS = ("id", "geometry", "subtype", "class", "surface", "sources")
SEGMENT_COLUMNS = ("id", "geometry", "subtype", "class", "road_surface", "width_rules", "sources")
_state_lock = threading.Lock()
_stac_tables: dict[str, pa.Table] = {}
_filesystems: dict[tuple[float | None, float | None], fs.S3FileSystem] = {}
[docs]
def latest_release() -> str:
"""The current Overture release."""
for delay in STAC_RETRY_DELAYS_S:
try:
return get_latest_release()
except Exception: # overturemaps raises a bare Exception for any catalog error
time.sleep(delay)
return get_latest_release()
[docs]
def index_root() -> Path:
"""Directory holding one cached STAC index per release."""
from openworld_radio_twin.config import get_settings
settings = get_settings()
configured = settings.overture_index_root
if configured is not None:
return Path(configured).expanduser()
return settings.cache_root / "overture-index"
[docs]
def bbox_filter(bbox: tuple[float, float, float, float]) -> pc.Expression:
xmin, ymin, xmax, ymax = bbox
return (
(pc.field("bbox", "xmin") < xmax)
& (pc.field("bbox", "xmax") > xmin)
& (pc.field("bbox", "ymin") < ymax)
& (pc.field("bbox", "ymax") > ymin)
)
def _download_index(release: str) -> bytes:
url = STAC_COLLECTIONS_URL.format(release=release)
for delay in STAC_RETRY_DELAYS_S:
try:
with urlopen(url) as response:
return response.read()
except OSError: # HTTPError 429 from the throttling front end is an OSError
time.sleep(delay)
with urlopen(url) as response:
return response.read()
def _stac_table(release: str) -> pa.Table:
# Concurrent theme reads share one load: the lock covers the download as well.
with _state_lock:
table = _stac_tables.get(release)
if table is None:
path = index_root() / f"collections-{release}.parquet"
if not path.is_file():
payload = _download_index(release)
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}")
temporary.write_bytes(payload)
temporary.replace(path)
table = pq.read_table(path)
_stac_tables[release] = table
return table
def _filesystem(connect_timeout: float | None, request_timeout: float | None) -> fs.S3FileSystem:
key = (connect_timeout, request_timeout)
with _state_lock:
if key not in _filesystems:
_filesystems[key] = fs.S3FileSystem(
anonymous=True,
region=S3_REGION,
connect_timeout=connect_timeout,
request_timeout=request_timeout,
)
return _filesystems[key]
[docs]
def intersecting_files(
overture_type: str,
bbox: tuple[float, float, float, float],
release: str,
) -> list[str]:
"""Parquet objects of ``overture_type`` whose extent crosses ``bbox``.
Indexes up to the 2026-08 releases name each object's type in ``collection``; from
release 2026-09-23.0 that column is empty and the type appears only in the object path
(``.../theme=buildings/type=building/part-...``), so both are accepted.
"""
selection = (pc.field("type") == "Feature") & bbox_filter(bbox)
rows = _stac_table(release).filter(selection).select(["collection", "assets"]).to_pylist()
in_path = f"/type={overture_type}/"
files = []
for row in rows:
href = row["assets"]["aws"]["alternate"]["s3"]["href"]
collection = row["collection"]
if collection == overture_type or (collection is None and in_path in href):
files.append(href[len("s3://") :])
return files
[docs]
def read_batches(
overture_type: str,
bbox: tuple[float, float, float, float],
release: str,
columns: tuple[str, ...],
*,
connect_timeout: float | None = None,
request_timeout: float | None = None,
) -> list[pa.RecordBatch]:
"""The non-empty record batches inside ``bbox``, restricted to ``columns`` that exist."""
files = intersecting_files(overture_type, bbox, release)
if not files:
return []
dataset = ds.dataset(files, filesystem=_filesystem(connect_timeout, request_timeout))
available = set(dataset.schema.names)
batches = dataset.to_batches(
columns=[name for name in columns if name in available],
filter=bbox_filter(bbox),
use_threads=True,
batch_readahead=16,
fragment_readahead=4,
)
return [batch for batch in batches if batch.num_rows > 0]