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 store_to_disk, StorageError 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) -> 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) 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) -> 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): 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) -> 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): 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}") collected_pois, failed_query_boxes = fetch_fn(BBOXEN, overpass_query) logger.warning(f"failed query boxes: {failed_query_boxes}") # store pois if collected_pois: try: saved_path = store_to_disk( results=collected_pois, poi_type=query_name, output_dir=OUTPUT_DIR, ) logger.info(f"{len(collected_pois)} POIs gespeichert in {saved_path}") except StorageError as exc: logger.error(f"Fehler beim Speichern: {exc}") else: logger.warning(f"Nichts zu speichern") if __name__ == "__main__": main()