overpass/main.py

123 lines
4.8 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import logging
from concurrent.futures import ThreadPoolExecutor, as_completed
from threading import Lock
from models import POI, FetchMode
from overpass import load_pois, OverpassApiError
from storage import StorageError, JsonStorage
from utils import timer
from config import BBOXEN, QUERY, OUTPUT_DIR, FETCH_MODE, MAX_WORKERS
# ---------------------------------------------------------------------------
# Logging konfigurieren
# ---------------------------------------------------------------------------
# Erinnerung: Log-Levels -> DEBUG, INFO, WARNING, ERROR, CRITICAL
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
datefmt="%H:%M:%S",
)
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# 3 verschiedene Funktionen (zum Vergleichen der Zeit)
# ---------------------------------------------------------------------------
@timer
def fetch_serial(bboxen: dict, overpass_query: str, poi_type: str) -> tuple[list[POI], list[int]]:
"""Serielle Variante eine Box nach der anderen."""
collected, failed = [], []
for name, bbox in bboxen.items():
try:
pois = load_pois(overpass_query=overpass_query, bbox=bbox, poi_type=poi_type)
collected.extend(pois)
logger.info(f"{name}: {len(pois)} POIs gefunden")
except OverpassApiError as exc:
logger.error(f"Fehler bei Box '{name}': {exc}")
failed.append(name)
return collected, failed
@timer
def fetch_concurrent(bboxen: dict, overpass_query: str, poi_type: str) -> tuple[list[POI], list[int]]:
"""Concurrent Variante alle Boxen parallel."""
collected, failed = [], []
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
futures = {
pool.submit(load_pois, overpass_query, bbox, poi_type): name for name, bbox in bboxen.items()
}
for future in as_completed(futures):
name = futures[future]
try:
pois = future.result()
logger.info(f"Box {name}: {len(pois)} POIs gefunden")
collected.extend(pois)
except OverpassApiError as exc:
logger.error(f"Fehler bei Box '{name}': {exc}")
failed.append(name)
return collected, failed
@timer
def fetch_concurrent_locked(bboxen: dict, overpass_query: str, poi_type: str) -> tuple[list[POI], list[int]]:
"""Concurrent Variante alle Boxen parallel."""
lock = Lock()
collected, failed = [], []
with ThreadPoolExecutor(max_workers=MAX_WORKERS) as pool:
futures = {
pool.submit(load_pois, overpass_query, bbox, poi_type): name for name, bbox in bboxen.items()
}
for future in as_completed(futures):
name = futures[future]
try:
pois = future.result()
logger.info(f"Box {name}: {len(pois)} POIs gefunden")
with lock:
collected.extend(pois)
except OverpassApiError as exc:
logger.error(f"Fehler bei Box '{name}': {exc}")
failed.append(name)
return collected, failed
FETCH_STRATEGIES = {
FetchMode.SERIAL: fetch_serial,
FetchMode.CONCURRENT: fetch_concurrent,
FetchMode.LOCKED: fetch_concurrent_locked,
}
# ---------------------------------------------------------------------------
# Hauptlogik
# ---------------------------------------------------------------------------
def main() -> None:
query_name: str = list(QUERY.keys())[0]
overpass_query: str = QUERY[query_name]
fetch_fn = FETCH_STRATEGIES[FETCH_MODE]
logger.info(f"Starte im Modus: {FETCH_MODE}")
# NEU: 'poi_type' als Attribut in der POI Klasse speichern -> bedingt ein gewisses 'prop-drilling' an
# verschiedenen Stellen im Code, macht aber die geplante .store()-Methode dann später einfacher, einheitlicher und
# übersichtlicher.
# FOLGE: bei jedem Durchlauf werden die bestehenden pois überschrieben (und nicht upgedated), was wir aber an
# unserer JsonStorage nicht ändern (kein Nachbauen von Datenbank-Logik...)
collected_pois, failed_query_boxes = fetch_fn(bboxen=BBOXEN, overpass_query=overpass_query, poi_type=query_name)
logger.warning(f"failed query boxes: {failed_query_boxes}")
storage = JsonStorage(output_dir=OUTPUT_DIR)
# storage = PostgresStorage(connection_string="postgresql://...")
if collected_pois:
try:
location = storage.store(collected_pois)
logger.info(f"{len(collected_pois)} POIs gespeichert: {location}")
except StorageError as exc:
logger.error(f"Fehler beim Speichern: {exc}")
else:
logger.warning("Nichts zu speichern")
if __name__ == "__main__":
main()