wanderlust-chatbot / scripts /run_pipeline.py
Kiriten892's picture
feat: security audit fixes, performance improvements & global data pipeline
dea44a6
Raw
History Blame Contribute Delete
11 kB
"""
Wanderlust Data Scraping Pipeline
===================================
Orchestrates scraping from OSM, TripAdvisor, Google Maps, Booking.com
and merges results into the chatbot knowledge base and backend seeder files.
Usage:
# Full pipeline (all sources, all cities)
python scripts/run_pipeline.py
# Only OSM (free, fast) for specific cities
python scripts/run_pipeline.py --sources osm --cities bangkok,tokyo,paris
# TripAdvisor + Google Maps for all cities
python scripts/run_pipeline.py --sources tripadvisor,google_maps
# Hotels only (for backend expansion)
python scripts/run_pipeline.py --sources booking
# Option B: OmkarCloud scrapers (requires cloned repos)
python scripts/run_pipeline.py --sources omkarcloud --gm-repo path/to/google-maps-scraper --ta-repo path/to/tripadvisor-scraper
# Full enrichment: OSM + OmkarCloud
python scripts/run_pipeline.py --sources osm,omkarcloud --cities all
# Dry run (scrape only, don't merge into DB files)
python scripts/run_pipeline.py --dry-run
# Ignore cache (re-scrape everything)
python scripts/run_pipeline.py --no-cache
"""
import argparse
import logging
import sys
import os
import json
import time
from datetime import datetime
# Allow running from chatbot-ml-service root
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from scripts.scrapers.config import TARGET_CITIES
from scripts.scrapers import osm_scraper, tripadvisor_scraper, google_maps_scraper, booking_scraper
from scripts.scrapers.omkarcloud_adapter import OmkarcloudAdapter
from scripts.processors import normalize_restaurants, normalize_events, normalize_destinations, normalize_hotels
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[
logging.StreamHandler(sys.stdout),
logging.FileHandler(f"scrape_{datetime.now().strftime('%Y%m%d_%H%M%S')}.log"),
],
)
logger = logging.getLogger(__name__)
ALL_SOURCES = ["osm", "tripadvisor", "google_maps", "booking", "omkarcloud"]
def parse_args():
parser = argparse.ArgumentParser(description="Wanderlust Data Scraping Pipeline")
parser.add_argument(
"--sources",
default="osm",
help=f"Comma-separated list of sources: {', '.join(ALL_SOURCES)} (default: osm)",
)
parser.add_argument(
"--cities",
default="all",
help="Comma-separated city names or 'all' (default: all)",
)
parser.add_argument(
"--dry-run",
action="store_true",
help="Scrape data but don't write to knowledge base files",
)
parser.add_argument(
"--no-cache",
action="store_true",
help="Ignore cached scrape results and re-scrape",
)
parser.add_argument(
"--limit",
type=int,
default=None,
help="Limit number of cities to scrape (for testing)",
)
parser.add_argument(
"--gm-repo",
default="omkarcloud/google-maps-scraper",
help="Path to cloned omkarcloud/google-maps-scraper repo (default: omkarcloud/google-maps-scraper)",
)
parser.add_argument(
"--ta-repo",
default="omkarcloud/tripadvisor-scraper",
help="Path to cloned omkarcloud/tripadvisor-scraper repo (default: omkarcloud/tripadvisor-scraper)",
)
return parser.parse_args()
def filter_cities(cities_arg: str) -> list:
"""Parse --cities argument and return filtered city list."""
if cities_arg.lower() == "all":
return TARGET_CITIES
names = [c.strip().lower() for c in cities_arg.split(",")]
filtered = [c for c in TARGET_CITIES if c["city"].lower() in names
or c["destination_id"] in names]
if not filtered:
logger.warning(f"No cities matched '{cities_arg}'. Running all.")
return TARGET_CITIES
return filtered
def run_osm(cities: list, use_cache: bool) -> list:
"""Run OSM scraper for all cities."""
logger.info(f"=== OSM Scraper: {len(cities)} cities ===")
return osm_scraper.scrape_all_cities(cities, cache=use_cache)
def run_omkarcloud(cities: list, use_cache: bool, gm_repo: str, ta_repo: str) -> list:
"""Run omkarcloud (Google Maps + TripAdvisor) scrapers for all cities."""
logger.info(f"=== OmkarCloud Scraper: {len(cities)} cities ===")
adapter = OmkarcloudAdapter(
gm_repo_path=gm_repo,
ta_repo_path=ta_repo,
use_cache=use_cache,
)
return adapter.scrape_all_cities(cities)
def run_tripadvisor(cities: list, use_cache: bool) -> list:
"""Run TripAdvisor scraper for all cities."""
logger.info(f"=== TripAdvisor Scraper: {len(cities)} cities ===")
return tripadvisor_scraper.scrape_all_cities(cities, cache=use_cache)
def run_google_maps(cities: list, use_cache: bool) -> list:
"""Run Google Maps scraper for all cities."""
logger.info(f"=== Google Maps Scraper: {len(cities)} cities ===")
return google_maps_scraper.scrape_all_cities(cities, cache=use_cache)
def run_booking(cities: list, use_cache: bool) -> list:
"""Run Booking.com scraper for all cities."""
logger.info(f"=== Booking.com Scraper: {len(cities)} cities ===")
return booking_scraper.scrape_all_cities(cities, cache=use_cache)
def merge_all(city_results: list, dry_run: bool) -> dict:
"""
Collect all scraped data and merge into knowledge base files.
city_results: list of per-city dicts with keys:
- restaurants (list)
- attractions (list)
- events (list)
- hotels (list, optional)
"""
all_restaurants = []
all_events = []
all_hotels_cities = []
for city_data in city_results:
restaurants = city_data.get("restaurants", [])
events = city_data.get("events", [])
hotels = city_data.get("hotels", [])
all_restaurants.extend(restaurants)
all_events.extend(events)
if hotels:
all_hotels_cities.append(city_data)
stats = {
"raw_restaurants": len(all_restaurants),
"raw_events": len(all_events),
"raw_hotels_cities": len(all_hotels_cities),
}
if dry_run:
logger.info(f"DRY RUN – would merge: {stats}")
return stats
# Merge restaurants β†’ cuisine_database.json
if all_restaurants:
logger.info(f"Merging {len(all_restaurants)} restaurants...")
r_stats = normalize_restaurants.merge_into_cuisine_db(all_restaurants)
stats["restaurants"] = r_stats
# Merge events β†’ events_calendar.json
if all_events:
logger.info(f"Merging {len(all_events)} events...")
e_stats = normalize_events.merge_into_events_db(all_events)
stats["events"] = e_stats
# Enrich destinations with attractions
logger.info("Enriching destinations with attraction data...")
d_stats = normalize_destinations.enrich_destinations(city_results)
stats["destinations"] = d_stats
# Save international hotels
if all_hotels_cities:
logger.info(f"Saving hotels from {len(all_hotels_cities)} cities...")
h_stats = normalize_hotels.save_hotels(all_hotels_cities)
stats["hotels"] = h_stats
return stats
def print_summary(stats: dict, elapsed: float):
"""Print a human-readable summary of the pipeline run."""
print("\n" + "=" * 60)
print(" WANDERLUST SCRAPING PIPELINE – SUMMARY")
print("=" * 60)
print(f" Time elapsed : {elapsed:.1f}s ({elapsed/60:.1f} min)")
print(f" Raw restaurants : {stats.get('raw_restaurants', 0)}")
print(f" Raw events : {stats.get('raw_events', 0)}")
if "restaurants" in stats:
r = stats["restaurants"]
print(f"\n cuisine_database.json:")
print(f" Added : {r.get('added', 0)}")
print(f" Skipped (dupes) : {r.get('skipped_duplicate', 0)}")
print(f" Total in DB : {r.get('total_in_db', 0)}")
if "events" in stats:
e = stats["events"]
print(f"\n events_calendar.json:")
print(f" Added : {e.get('added', 0)}")
print(f" Total in DB : {e.get('total_in_db', 0)}")
if "destinations" in stats:
d = stats["destinations"]
print(f"\n destinations.json:")
print(f" Destinations enriched : {d.get('destinations_enriched', 0)}")
print(f" Activities added : {d.get('activities_added', 0)}")
if "hotels" in stats:
h = stats["hotels"]
print(f"\n hotels_international.json:")
print(f" Added : {h.get('added', 0)}")
print(f" Total in file : {h.get('total_hotels', 0)}")
print(f" Output : {h.get('output_file', '')}")
print("=" * 60)
print("\nNext steps:")
print(" 1. Review scraped_data/ for quality check")
print(" 2. Rebuild vector store: python -c \"from app.services.vector_store import VectorStore; VectorStore().initialize()\"")
print(" 3. Restart chatbot: uvicorn app.main:app --port 8000")
print(" 4. For backend hotels: copy hotels_international.json β†’ BackEnd seeder\n")
def main():
args = parse_args()
sources = [s.strip().lower() for s in args.sources.split(",")]
cities = filter_cities(args.cities)
use_cache = not args.no_cache
if args.limit:
cities = cities[:args.limit]
logger.info(f"Pipeline start: sources={sources}, cities={len(cities)}, cache={use_cache}, dry_run={args.dry_run}")
start = time.time()
all_city_data: dict = {} # destination_id β†’ merged city data
def _merge_city(base: dict, extra: dict):
"""Merge two city data dicts by combining their lists."""
result = dict(base)
for key in ["restaurants", "attractions", "events", "hotels"]:
result.setdefault(key, [])
result[key].extend(extra.get(key, []))
return result
def _index_results(results: list):
for city_data in results:
did = city_data.get("destination_id", "")
if did in all_city_data:
all_city_data[did] = _merge_city(all_city_data[did], city_data)
else:
all_city_data[did] = city_data
# Run each source
if "osm" in sources:
_index_results(run_osm(cities, use_cache))
if "tripadvisor" in sources:
_index_results(run_tripadvisor(cities, use_cache))
if "google_maps" in sources:
_index_results(run_google_maps(cities, use_cache))
if "booking" in sources:
_index_results(run_booking(cities, use_cache))
if "omkarcloud" in sources:
_index_results(run_omkarcloud(cities, use_cache, args.gm_repo, args.ta_repo))
# Merge into DB files
merged_list = list(all_city_data.values())
stats = merge_all(merged_list, dry_run=args.dry_run)
elapsed = time.time() - start
print_summary(stats, elapsed)
if __name__ == "__main__":
main()