123 lines
4.8 KiB
Python
123 lines
4.8 KiB
Python
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() |