import snowflake.connector import config import traceback from datetime import datetime from typing import List from utils.list_utils import chunks from workers.google_campaigns_worker import GoogleCampaignsWorker from workers.google_ad_sets_worker import GoogleAdSetsWorker from utils.reporting.worker_logger import WorkerLogger from utils.snowflake.google_data.models import ( GRASCampaignsQueryResult, GRASAdSetsQueryResult, GoogleCampaignModel, GoogleAdSetModel, ) from utils.json_helpers import json_load from utils.snowflake.google_data.utils import GoogleImporterUtils from utils.snowflake.google_data.queries import Query snowflake_connector = snowflake.connector.connect( user=config.SNOWFLAKE_USER, password=config.SNOWFLAKE_PASS, account=config.SNOWFLAKE_ACCOUNT, database=config.ADS_SNOWFLAKE_DB, schema=config.ADS_SNOWFLAKE_SCHEMA, warehouse=config.SNOWFLAKE_WAREHOUSE, ) class GoogleCampaignsImporter: utils = GoogleImporterUtils() def fetch(self) -> List[GRASCampaignsQueryResult]: with snowflake_connector as con: marketing_accounts = self.utils.marketing_accounts() data = con.cursor().execute(Query.fetch_campaigns_query(marketing_accounts)).fetchall() return [GRASCampaignsQueryResult._make(result) for result in data] class GoogleAdSetsFetcher: utils = GoogleImporterUtils() def fetch(self): with snowflake_connector as con: marketing_accounts = self.utils.marketing_accounts() data = con.cursor().execute(Query.fetch_ad_sets_query(marketing_accounts)).fetchall() return [GRASAdSetsQueryResult._make(result) for result in data] class GoogleCampaignsAdsImporter: campaigns_importer: GoogleCampaignsImporter logger: WorkerLogger = WorkerLogger() utils: GoogleImporterUtils = GoogleImporterUtils() def import_campaigns(self): self.logger.info("Google campaigns import started", f"Started at {datetime.now()}") try: self.__import_data() except Exception as error: self.logger.error("Google campaigns import error", f"{error}\n{traceback.format_exc()}") def __import_data(self): self.campaigns_importer = GoogleCampaignsImporter() self.__import_campaigns() def __import_campaigns(self): results = self.campaigns_importer.fetch() if not results: self.logger.error("Google campaigns import error", "No campaigns returned from GRAS") for chunk in chunks(results, 500): campaigns = list(map(self.__map_campaign_model_from, chunk)) GoogleCampaignsWorker().perform_async(campaigns) self.logger.success( "Google campaigns import finished", f"Import finished at {datetime.now()} count: {len(results)}" ) def __map_campaign_model_from(self, item) -> GoogleCampaignModel: campaign = GoogleCampaignModel( id=item.campaign_id, name=item.campaign_name, bidding_type=item.bidding_type, start_date=item.start_date, end_date=item.end_date, genders=json_load(item.genders), ages=json_load(item.ages), types=self.utils.map_types(item.types), links=self.utils.map_links(item.urls), territories=json_load(item.territories), spend=item.spend, budget=item.budget, account=json_load(item.accounts)[0], linkfire_links=self.utils.map_linkfire_links(item.linkfire_links) ) campaign.artists = self.utils.map_artists_external_ids(item) return campaign class GoogleAdSetsImporter: importer: GoogleAdSetsFetcher logger: WorkerLogger = WorkerLogger() utils: GoogleImporterUtils = GoogleImporterUtils() def import_ads(self): self.logger.info("Google ad sets import started", f"Started at {datetime.now()}") try: self.__import_data() except Exception as error: self.logger.error("Google ad sets import error", f"{error}\n{traceback.format_exc()}") def __import_data(self): self.importer = GoogleAdSetsFetcher() self.__import_ads() def __import_ads(self): results = self.importer.fetch() if not results: self.logger.error("Google ad sets import error", "No ad sets returned from GRAS") for chunk in chunks(results, 500): ad_sets = list(map(self.__map_ad_set_model, chunk)) GoogleAdSetsWorker().perform_async(ad_sets) self.logger.success("Google ad sets import finished", f"Import {len(results)} ads sets finished at {datetime.now()}") def __map_ad_set_model(self, item) -> GoogleAdSetModel: return GoogleAdSetModel( id=item.ad_group_id, name=item.ad_group_name, campaign_id=str(item.campaign_id), start_date=item.start_date, end_date=item.end_date, bidding_type=item.bidding_type, ages=json_load(item.ages), genders=json_load(item.genders), types=self.utils.map_types(item.types), links=self.utils.map_links(item.urls), territories=json_load(item.territories), spend=item.spend, )