#!/usr/bin/env python import io import os import zipfile from typing import Literal, TypedDict import click import httpx import pandas as pd import sqlalchemy as sa from fansifter_common.adapters.aws.secretsmanager import SecretsManager from dmp.adapters.db import ReportingDB from dmp.config import settings COUNTRIES_URL = "https://download.geonames.org/export/dump/countryInfo.txt" CITIES_SIZE: Literal["500", "1000", "5000", "15000"] = "500" CITIES_URL = f"https://download.geonames.org/export/dump/cities{CITIES_SIZE}.zip" REGIONS_URL = "https://download.geonames.org/export/dump/admin1CodesASCII.txt" POSTAL_CODES_URL = "https://download.geonames.org/export/zip/{country_code2}.zip" POSTAL_CODES_COUNTRIES = [ "US", "GB", "FR", "ES", "CA", "MX", "AU", "DE", "BR", "IT", "NL", "AR", "BE", "IE", "SE", "CH", "PL", "CL", "DK", "IN", "AT", "LV", "CO", "PE", "NO", "LI", ] class Country(TypedDict): code2: str code3: str name: str class Region(TypedDict): code: str country_code2: str name: str name_native: str class City(TypedDict): id: str name: str name_native: str alternate_names: str country_code2: str region_code: str latitude: float longitude: float population: int admin2_code: str admin3_code: str class PostalCode(TypedDict): country_code2: str postal_code: str place_name: str admin1_code: str admin2_code: str admin3_code: str latitude: float longitude: float class GeoNamesHandler: def __init__(self, conn: sa.Connection) -> None: self.conn = conn def handle(self) -> None: self.load_countries() self.load_regions() self.load_cities() for country_code2 in POSTAL_CODES_COUNTRIES: self.load_postal_codes(country_code2) def load_countries(self) -> None: click.echo("Loading countries...") countries = [] response = httpx.get(COUNTRIES_URL) response.raise_for_status() for line in response.iter_lines(): if line.startswith("#") or line.strip() == "": continue parts = line.split("\t") code2 = parts[0].strip() code3 = parts[1].strip() name = parts[4].strip() countries.append( Country( code2=code2, code3=code3, name=name, ) ) df = pd.DataFrame(countries) df.to_sql( "location_country_tmp", self.conn, if_exists="replace", index=False, ) self.conn.execute( sa.text(""" MERGE INTO location_country AS country USING ( SELECT * FROM location_country_tmp ) AS tmp ON country.code2 = tmp.code2 WHEN MATCHED THEN UPDATE SET code3 = tmp.code3, name = tmp.name WHEN NOT MATCHED THEN INSERT (code2, code3, name) VALUES (tmp.code2, tmp.code3, tmp.name) """) ) self.conn.execute(sa.text("DROP TABLE location_country_tmp")) click.echo("Countries loaded.") def load_cities(self) -> None: click.echo("Loading cities...") cities = [] response = httpx.get(CITIES_URL) response.raise_for_status() with ( zipfile.ZipFile(io.BytesIO(response.content)) as zf, zf.open(f"cities{CITIES_SIZE}.txt") as zf_io, ): for line in zf_io.read().decode().splitlines(): if line.strip() == "": continue parts = line.split("\t") city_id = parts[0].strip() name = parts[1].strip() name_native = parts[2].strip() alternate_names = parts[3].strip() latitude = parts[4].strip() longitude = parts[5].strip() country_code2 = parts[8].strip() region_code = parts[10].strip() admin2_code = parts[11].strip() admin3_code = parts[12].strip() population = int(parts[14].strip() or "0") # Add parsed data to the list cities.append( City( id=city_id, name=name, name_native=name_native, alternate_names=alternate_names, country_code2=country_code2, region_code=region_code, latitude=float(latitude), longitude=float(longitude), population=population, admin2_code=admin2_code, admin3_code=admin3_code, ) ) df = pd.DataFrame(cities) df.to_sql( "location_city_tmp", self.conn, if_exists="replace", chunksize=200_000, index=False, ) self.conn.execute( sa.text(""" MERGE INTO location_city AS city USING ( SELECT * FROM location_city_tmp ) AS tmp ON city.id = tmp.id WHEN MATCHED THEN UPDATE SET name = tmp.name, name_native = tmp.name_native, alternate_names = tmp.alternate_names, country_code2 = tmp.country_code2, region_code = tmp.region_code, latitude = tmp.latitude, longitude = tmp.longitude, population = tmp.population, admin2_code = tmp.admin2_code, admin3_code = tmp.admin3_code WHEN NOT MATCHED THEN INSERT ( id, name, name_native, alternate_names, country_code2, region_code, latitude, longitude, population, admin2_code, admin3_code ) VALUES ( tmp.id, tmp.name, tmp.name_native, tmp.alternate_names, tmp.country_code2, tmp.region_code, tmp.latitude, tmp.longitude, tmp.population, tmp.admin2_code, tmp.admin3_code ) """) ) self.conn.execute(sa.text("DROP TABLE location_city_tmp")) click.echo("Cities loaded.") def load_regions(self) -> None: click.echo("Loading regions...") regions = [] response = httpx.get(REGIONS_URL) response.raise_for_status() for line in response.iter_lines(): if line.strip() == "": continue parts = line.split("\t") country_region_code = parts[0].strip() country_code2, code = country_region_code.split(".") name = parts[1].strip() name_native = parts[2].strip() regions.append( Region( code=code, country_code2=country_code2, name=name, name_native=name_native, ) ) df = pd.DataFrame(regions) df.to_sql( "location_region_tmp", self.conn, chunksize=200_000, if_exists="replace", index=False, ) self.conn.execute( sa.text(""" MERGE INTO location_region AS region USING ( SELECT * FROM location_region_tmp ) AS tmp ON region.code = tmp.code AND region.country_code2 = tmp.country_code2 WHEN MATCHED THEN UPDATE SET name = tmp.name, name_native = tmp.name_native WHEN NOT MATCHED THEN INSERT (code, country_code2, name, name_native) VALUES (tmp.code, tmp.country_code2, tmp.name, tmp.name_native) """) ) self.conn.execute(sa.text("DROP TABLE location_region_tmp")) click.echo("Regions loaded.") def load_postal_codes(self, country_code2: str) -> None: click.echo(f"Loading postal codes for {country_code2} ...") postal_codes = [] response = httpx.get(POSTAL_CODES_URL.format(country_code2=country_code2)) response.raise_for_status() with ( zipfile.ZipFile(io.BytesIO(response.content)) as zf, zf.open(f"{country_code2}.txt") as zf_io, ): for line in zf_io.read().decode().splitlines(): if line.strip() == "": continue parts = line.split("\t") postal_code = parts[1].strip() place_name = parts[2].strip() admin1_code = parts[4].strip() admin2_code = parts[6].strip() admin3_code = parts[7].strip() latitude = parts[9].strip() longitude = parts[10].strip() postal_codes.append( PostalCode( country_code2=country_code2, postal_code=postal_code, place_name=place_name, admin1_code=admin1_code, admin2_code=admin2_code, admin3_code=admin3_code, latitude=latitude, longitude=longitude, ) ) df = pd.DataFrame(postal_codes) df.to_sql( "location_postal_code_tmp", self.conn, chunksize=200_000, if_exists="replace", index=False, ) self.conn.execute( sa.text(""" MERGE INTO location_postal_code AS postal_code USING ( SELECT * FROM location_postal_code_tmp ) AS tmp ON postal_code.postal_code = tmp.postal_code AND postal_code.country_code2 = tmp.country_code2 AND postal_code.admin1_code = tmp.admin1_code AND postal_code.admin2_code = tmp.admin2_code AND postal_code.admin3_code = tmp.admin3_code AND postal_code.place_name = tmp.place_name WHEN MATCHED AND ( postal_code.latitude IS DISTINCT FROM tmp.latitude OR postal_code.longitude IS DISTINCT FROM tmp.longitude ) THEN UPDATE SET latitude = tmp.latitude, longitude = tmp.longitude WHEN NOT MATCHED THEN INSERT (postal_code, country_code2, place_name, admin1_code, admin2_code, admin3_code, latitude, longitude) VALUES (tmp.postal_code, tmp.country_code2, tmp.place_name, tmp.admin1_code, tmp.admin2_code, tmp.admin3_code, tmp.latitude, tmp.longitude) """) ) self.conn.execute(sa.text("DROP TABLE location_postal_code_tmp")) click.echo(f"Postal codes for {country_code2} loaded.") @click.command @click.option("-e", "--environment", help="Environment", required=True) @click.option("--region-name", help="AWS region", required=True, default="us-east-1") def sync_geonames_data(environment: str, region_name: str) -> None: """Sync GeoNames data.""" click.echo(f"Syncing GeoNames data for {environment}") # Initialize secrets manager secrets_manager = SecretsManager(region_name=region_name) # Get secrets snowflake_private_key = secrets_manager.get_secret( f"{environment}/ows-dmp/SNOWFLAKE_PRIVATE_KEY" ) snowflake_key_passphrase = secrets_manager.get_secret( f"{environment}/ows-dmp/SNOWFLAKE_KEY_PASSPHRASE" ) os.environ["ENVIRONMENT"] = environment os.environ["AWS_REGION_NAME"] = region_name os.environ["ORM_MODELS"] = "" os.environ["APP_RUN_FROM_CLI"] = "1" os.environ["SNOWFLAKE_PRIVATE_KEY"] = snowflake_private_key os.environ["SNOWFLAKE_KEY_PASSPHRASE"] = snowflake_key_passphrase os.environ["SNOWFLAKE_ACCOUNT"] = "delphi.us-east-1" os.environ["SNOWFLAKE_DATABASE"] = "FANSIFTER_APP_REPORTING" os.environ["SNOWFLAKE_WAREHOUSE"] = f"{environment.upper()}_ETL_WH" os.environ["SNOWFLAKE_ROLE"] = f"{environment.upper()}_OWS_DMP" os.environ["SNOWFLAKE_SCHEMA"] = environment.upper() os.environ["SNOWFLAKE_USER"] = f"{environment.upper()}_OWS_DMP" os.environ["DMP_KMS_KEY_ID"] = "" # Initialize Snowflake connection db = ReportingDB( url=settings.snowflake_url, engine_args={ "echo": settings.snowflake_echo, "poolclass": sa.pool.QueuePool, "pool_size": settings.snowflake_pool_size, "max_overflow": settings.snowflake_pool_max_overflow, "pool_recycle": settings.snowflake_pool_recycle, "pool_pre_ping": settings.snowflake_pool_pre_ping, "pool_reset_on_return": settings.snowflake_pool_reset_on_return, "connect_args": settings.snowflake_connect_args, }, ) with db, db.engine.begin() as conn: # Test connection conn.execute(sa.text("SELECT 1")) click.echo("Connected to Snowflake.") click.echo("Syncing GeoNames data...") handler = GeoNamesHandler(conn) handler.handle() click.echo("GeoNames data synced.") if __name__ == "__main__": sync_geonames_data()