# ============================================================================== # Author: Simon-Pierre Boucher # File: restoka/connectors/yelp.py # Desc: Connecteur d'ENRICHISSEMENT Yelp Fusion (avis/notes) — gratuit, # 500 requêtes/jour. N'émet AUCUNE fiche restaurant : il complète les # fiches existantes avec rating, review_count, price, categories et # phone dans details.yelp (colonne d'enrichissement, jamais écrasée # par les re-crawls des sources). # # Croisement CONSERVATEUR (jamais fusionner deux établissements) : # 1. par TÉLÉPHONE (/businesses/search/phone) — signal fort ; si # plusieurs résultats, on exige le GPS le plus proche (<300 m) ; # 2. sinon par NOM+ADRESSE+GPS (/businesses/matches) — l'endpoint # d'appariement officiel de Yelp, seuil par défaut (conservateur). # # Clé requise : YELP_API_KEY dans .env (https://www.yelp.com/developers # — app gratuite). Sans clé, le connecteur lève SkipSource et le # pipeline passe son tour SANS marquer d'échec (statut « clé requise » # dans data/sources.json). # Base légale : API officielle Yelp Fusion, conditions Display # Requirements (attribution Yelp affichée avec la note sur la fiche). # ============================================================================== from __future__ import annotations import datetime import json import math import os import sys import time from ..normalize import normalize_phone from ..schema import Restaurant from .base import BaseConnector, SkipSource API = "https://api.yelp.com/v3" REFRESH_DAYS = 30 # re-vérification des notes captées MAX_REQUESTS_PER_RUN = int(os.environ.get("YELP_BUDGET", "450")) # quota 500/j MAX_PHONE_DISTANCE_M = 300 # ambiguïté téléphone : GPS requis def _haversine_m(lat1, lng1, lat2, lng2) -> float: r = 6371000.0 p1, p2 = math.radians(lat1), math.radians(lat2) dp, dl = math.radians(lat2 - lat1), math.radians(lng2 - lng1) a = math.sin(dp / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dl / 2) ** 2 return 2 * r * math.asin(math.sqrt(a)) class YelpConnector(BaseConnector): source_id = "yelp" request_delay = 0.35 # politesse API (limite 5 QPS) timeout = 20 use_detail_cache = False enrichment_only = True # n'émet aucune fiche (ingest.run) def _api(self, path: str, params: dict) -> dict: resp = self.get(f"{API}{path}", params=params, headers={"Authorization": f"Bearer {self._key}"}) return resp.json() def _payload(self, biz: dict, matched_by: str) -> dict: return { "id": biz.get("id"), "url": (biz.get("url") or "").split("?")[0], "name": biz.get("name"), "rating": biz.get("rating"), "review_count": biz.get("review_count"), "price": biz.get("price"), "categories": [c.get("title") for c in biz.get("categories") or []], "phone": biz.get("phone"), "matched_by": matched_by, "fetched_at": datetime.datetime.now(datetime.timezone.utc) .strftime("%Y-%m-%dT%H:%M:%SZ"), } def _match_by_phone(self, phone: str, lat, lng) -> dict | None: data = self._api("/businesses/search/phone", {"phone": phone}) businesses = data.get("businesses") or [] if not businesses: return None if len(businesses) == 1: return self._payload(businesses[0], "telephone") if lat is None or lng is None: return None # plusieurs candidats, pas de GPS : on passe best, best_d = None, None for b in businesses: c = b.get("coordinates") or {} if c.get("latitude") is None or c.get("longitude") is None: continue d = _haversine_m(lat, lng, c["latitude"], c["longitude"]) if best_d is None or d < best_d: best, best_d = b, d if best is not None and best_d is not None and best_d <= MAX_PHONE_DISTANCE_M: return self._payload(best, "telephone+gps") return None def _match_by_name(self, row) -> dict | None: if not (row["name"] and row["address"] and row["city"] and row["lat"] is not None and row["lng"] is not None): return None data = self._api("/businesses/matches", { "name": row["name"][:64], "address1": row["address"][:64], "city": row["city"][:64], "state": "QC", "country": "CA", "latitude": row["lat"], "longitude": row["lng"], "limit": 1, }) businesses = data.get("businesses") or [] if not businesses: return None biz = businesses[0] # /matches ne retourne pas rating/review_count : détail requis detail = self._api(f"/businesses/{biz['id']}", {}) return self._payload(detail or biz, "nom+gps") def fetch(self) -> list[Restaurant]: self._key = os.environ.get("YELP_API_KEY", "").strip() if not self._key: raise SkipSource("YELP_API_KEY manquant (.env) — connecteur prêt, " "clé gratuite : https://www.yelp.com/developers") from .. import db con = db.connect() now = time.time() stale = now - REFRESH_DAYS * 86400 budget = MAX_REQUESTS_PER_RUN enriched = skipped = 0 rows = con.execute( "SELECT uid, name, address, city, lat, lng, phone, details" " FROM restaurants WHERE active=1 AND dup_of IS NULL" " ORDER BY phone<>'' DESC, updated_at DESC").fetchall() for row in rows: if budget <= 0: break try: details = json.loads(row["details"] or "{}") except ValueError: details = {} yelp = details.get("yelp") or {} if yelp.get("fetched_at"): try: ts = datetime.datetime.strptime( yelp["fetched_at"], "%Y-%m-%dT%H:%M:%SZ") \ .replace(tzinfo=datetime.timezone.utc).timestamp() if ts > stale: continue # déjà frais (<30 j) except ValueError: pass payload = None try: phone = normalize_phone(row["phone"]) if phone: budget -= 1 payload = self._match_by_phone(phone, row["lat"], row["lng"]) if payload is None: budget -= 2 # /matches + détail payload = self._match_by_name(row) except Exception as exc: # un resto raté ne bloque pas skipped += 1 print(f"[resto-ka] yelp: {row['uid']} erreur: {exc}", file=sys.stderr) continue if payload: db.merge_details(con, row["uid"], {"yelp": payload}) con.commit() enriched += 1 con.commit() con.close() self.enriched_count = enriched self.enrich_message = (f"{enriched} resto(s) enrichis (avis/notes), " f"budget restant {max(budget, 0)} req") print(f"[resto-ka] yelp: {self.enrich_message}" + (f", {skipped} erreurs" if skipped else "")) return []