Task_7: concurrent workflow

This commit is contained in:
Marco Schmid 2026-04-30 18:01:09 +02:00
parent d0ab53185c
commit 4546305c1a
7 changed files with 131 additions and 15 deletions

1
.gitignore vendored
View File

@ -3,3 +3,4 @@
GIT_INSTRUCTOR.md
README_INSTRUCTOR.md
/results/

34
TASK.md
View File

@ -1,9 +1,29 @@
# TASK 6:
# TASK 7:
* bbox für Schweiz scheint zu gross und wirft einen error ... Lösungsmöglichkeiten?
-> Wir können die Schweiz (Koordinaten) in Unterregionen aufsplitten. Macht das bitte.
-> entfernt dazu die bbox für 'davos', nehmt die 'schweiz' und splittet sie in 4, 9 oder 16 Koordinaten-Tuples auf.
Wir haben nun einmal einen (noch nicht perfekten aber funktionstüchtigen) Workflow, in welchem wir die Daten von
einer externen API (Overpass) holen, daraus POI-Objekte machen und diese als .json auf unserer Festplatte speichern.
* Speichert und loggt in welchen Koordinaten-Tuples ein Fehler auftritt (gebt am Schluss eine Zusammenfassung
dieser fehlerhaften Queries aus)
* Bildet ein neues Modul `storage.py` und baut den Code, welcher zum Speichern der POIS als .json auf der Festplatte nötig ist.
Als nächstes versuchen wir den Code zu "parallelisieren" (ThreadPool):
1. wir bilden ein Modul config.py und lagern dort unseren Konfigurationscode aus
2. Timer einbauen (Dekorator) -> ist noch unbekannt: zeige ich vor
3. Code nach vorgegebenem Schema "concurrent" machen -> findet heraus, was in unserem Beispiel der Adapter ist
- optional: Lock verwenden
```python
from concurrent.futures import ThreadPoolExecutor, as_completed
adapter = CodewarsAdapter()
usernames = [" alice ", "bob", " carol ", " dave "]
with ThreadPoolExecutor(max_workers=4) as pool:
futures = {pool.submit(adapter.get_user, n): n for n in usernames}
for future in as_completed(futures):
name = futures[future]
try:
user = future.result()
print(f"{name}: {user.honor}Honor")
except Exception as exc:
print(f"{name}: Fehler -- {exc}")
```

Binary file not shown.

58
main.py
View File

@ -3,6 +3,9 @@ from models import POI
import logging
from queries.bergbahn import BERGBAHN_QUERY
from queries.restaurant import RESTAURANT_QUERY
from storage import store_to_disk, StorageError
from pathlib import Path
# ---------------------------------------------------------------------------
# Logging konfigurieren
@ -23,34 +26,73 @@ logger = logging.getLogger(__name__)
# Konfiguration
# ---------------------------------------------------------------------------
# BBOXEN = {
# "SW": (45.8, 5.9, 46.8, 8.2),
# "SO": (45.8, 8.2, 46.8, 10.5),
# "NW": (46.8, 5.9, 47.8, 8.2),
# "NO": (46.8, 8.2, 47.8, 10.5)
# }
BBOXEN = {
"davos": (46.72, 9.70, 46.92, 10.00),
"schweiz": (45.8, 5.9, 47.8, 10.5),
1: (45.8, 5.9, 46.4667, 7.4333),
2: (45.8, 7.4333, 46.4667, 8.9667),
3: (45.8, 8.9667, 46.4667, 10.5),
4: (46.4667, 5.9, 47.1333, 7.4333),
5: (46.4667, 7.4333, 47.1333, 8.9667),
6: (46.4667, 8.9667, 47.1333, 10.5),
7: (47.1333, 5.9, 47.8, 7.4333),
8: (47.1333, 7.4333, 47.8, 8.9667),
9: (47.1333, 8.9667, 47.8, 10.5)
}
# BBOXEN = {
# 1: (45.8, 5.9, 46.3, 7.05), 2: (45.8, 7.05, 46.3, 8.2), 3: (45.8, 8.2, 46.3, 9.35), 4: (45.8, 9.35, 46.3, 10.5),
# 5: (46.3, 5.9, 46.8, 7.05), 6: (46.3, 7.05, 46.8, 8.2), 7: (46.3, 8.2, 46.8, 9.35), 8: (46.3, 9.35, 46.8, 10.5),
# 9: (46.8, 5.9, 47.3, 7.05), 10: (46.8, 7.05, 47.3, 8.2), 11: (46.8, 8.2, 47.3, 9.35), 12: (46.8, 9.35, 47.3, 10.5),
# 13: (47.3, 5.9, 47.8, 7.05), 14: (47.3, 7.05, 47.8, 8.2), 15: (47.3, 8.2, 47.8, 9.35), 16: (47.3, 9.35, 47.8, 10.5)
# }
QUERY = {"bergbahn": BERGBAHN_QUERY}
OUTPUT_DIR = Path("results")
# ---------------------------------------------------------------------------
# Hauptlogik
# ---------------------------------------------------------------------------
def main() -> None:
query_name = list(QUERY.keys())[0]
collected_pois: list[POI] = []
query_name: str = list(QUERY.keys())[0]
overpass_query: str = QUERY[query_name]
failed_query_boxes: list[int] = []
# get pois from API
for name, bbox in BBOXEN.items():
logger.info(f"Starte Abfrage für Query: {query_name}, '{name}' mit bbox={bbox}")
try:
pois: list[POI] = load_pois(overpass_query=QUERY.get(query_name,""), bbox=bbox)
pois: list[POI] = load_pois(overpass_query=overpass_query, bbox=bbox)
logger.info(f"{name}: {len(pois)} POIs gefunden")
collected_pois.extend(pois)
except OverpassApiError as exc:
logger.error(f"Fehler bei '{name}': {exc}")
failed_query_boxes.append(name)
continue
logger.info(f"\n{name}: {len(pois)} POIs gefunden")
for poi in pois:
logger.info(f" {poi.id}: ({poi.lat}, {poi.lon})")
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()

53
storage.py Normal file
View File

@ -0,0 +1,53 @@
import json
import logging
from pathlib import Path
from models import POI
from dataclasses import asdict
logger = logging.getLogger(__name__)
class StorageError(Exception):
pass
def store_to_disk(
results: list[POI],
poi_type: str = "overpass",
output_dir: Path = Path("."),
) -> Path:
"""
Speichert eine Liste von OSM-Elementen (POIs) als JSON-Datei auf der Festplatte.
Der Dateiname wird aus dem poi_type-Parameter abgeleitet:
z.B. poi_type="bergbahn" "bergbahn_results.json"
Args:
results (list[dict]): Liste von OSM-Elementen, wie sie die
Overpass API unter "elements" zurückgibt.
poi_type (str): Bezeichnung des POI-Typs. Bestimmt den
Dateinamen. Standard: "overpass"
output_dir (Path): Zielverzeichnis. Wird erstellt falls nicht
vorhanden. Standard: aktuelles Verzeichnis.
Returns:
Path: Absoluter Pfad zur gespeicherten Datei.
Raises:
StorageError: Wenn das Verzeichnis nicht erstellt oder die Datei nicht
geschrieben werden kann (z.B. fehlende Schreibrechte).
"""
output_dir.mkdir(parents=True, exist_ok=True)
output_path = output_dir / f"{poi_type}_results.json"
try:
with output_path.open("w", encoding="utf-8") as f:
json.dump([asdict(poi) for poi in results], f, indent=2, ensure_ascii=False)
except OSError as exc:
logger.error(f"Fehler beim Speichern:{exc}")
raise StorageError("Fehler beim Speichern der POIs") from exc
return output_path.resolve()