import snowflake.connector import config import traceback from datetime import datetime from typing import List from utils.json_helpers import json_load from utils.list_utils import chunks from utils.snowflake.tiktok_data.utils import TikTokImporterUtils from utils.reporting.worker_logger import WorkerLogger from utils.snowflake.tiktok_data.queries import Query from utils.snowflake.tiktok_data.models import ( TikTokCampaignsQueryResult, TikTokCampaignModel, TikTokAdSetsQueryResult, TikTokAdSetModel ) from workers.tiktok_campaigns_worker import TikTokCampaignsWorker from workers.tiktok_ad_sets_worker import TikTokAdSetsWorker 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 TikTokCampaignsFetcher: utils = TikTokImporterUtils() def fetch(self) -> List[TikTokCampaignsQueryResult]: with snowflake_connector as con: marketing_accounts = self.utils.marketing_accounts() data = con.cursor().execute(Query.fetch_campaigns_query(marketing_accounts)).fetchall() return [TikTokCampaignsQueryResult._make(result) for result in data] class TikTokAdSetsFetcher: utils = TikTokImporterUtils() def fetch(self): with snowflake_connector as con: marketing_accounts = self.utils.marketing_accounts() data = con.cursor().execute(Query.fetch_campaign_ads_query(marketing_accounts)).fetchall() return [TikTokAdSetsQueryResult._make(result) for result in data] class TikTokCampaignsImporter: campaigns_importer: TikTokCampaignsFetcher logger: WorkerLogger = WorkerLogger() utils: TikTokImporterUtils = TikTokImporterUtils() def import_campaigns(self): self.logger.info("TikTok campaigns import started", f"Started at {datetime.now()}") try: self.__import_data() except Exception as error: self.logger.error("TikTok campaigns import error", f"{error}\n{traceback.format_exc()}") def __import_data(self): self.campaigns_importer = TikTokCampaignsFetcher() self.__import_campaigns() def __import_campaigns(self): results = self.campaigns_importer.fetch() if not results: self.logger.error("TikTok campaigns import error", "No campaigns returned from datasource") for chunk in chunks(results, 500): campaigns = list(map(self.__map_campaign_model_from, chunk)) TikTokCampaignsWorker().perform_async(campaigns) self.logger.success( "TikTok campaigns import finished", f"Import finished at {datetime.now()} count: {len(results)}" ) def __map_campaign_model_from(self, item) -> TikTokCampaignModel: campaign = TikTokCampaignModel( id=item.external_id, name=item.campaign_name, start_date=item.start_date, end_date=item.end_date, genders=json_load(item.genders), ages=json_load(item.ages), budget_spend=item.budget_spend, objective=item.objective, planned_budget=item.planned_budget, marketing_account_id=item.marketing_account_id, country_codes=json_load(item.country_codes), artists=self.utils.map_artists_external_ids(item) ) return campaign class TikTokAdSetsImporter: importer: TikTokAdSetsFetcher logger: WorkerLogger = WorkerLogger() utils: TikTokImporterUtils = TikTokImporterUtils() def import_ads(self): self.logger.info("Tiktok ad sets import started", f"Started at {datetime.now()}") try: self.__import_data() except Exception as error: self.logger.error("Tiktok ad sets import error", f"{error}\n{traceback.format_exc()}") def __import_data(self): self.importer = TikTokAdSetsFetcher() self.__import_ads() def __import_ads(self): results = self.importer.fetch() if not results: self.logger.error("Tiktok ad sets import error", "No ad sets returned from datasource") for chunk in chunks(results, 500): ad_sets = list(map(self.__map_ad_set_model, chunk)) TikTokAdSetsWorker().perform_async(ad_sets) self.logger.success("Tiktok ad sets import finished", f"Import {len(results)} ads sets finished at {datetime.now()}") def __map_ad_set_model(self, item) -> TikTokAdSetModel: return TikTokAdSetModel( id=item.ad_set_id, name=item.ad_set_name, start_date=item.start_date, end_date=item.end_date, campaign_id=item.campaign_id, objective=item.objective, destination_link=item.destination_link, planned_budget=item.planned_budget, budget_spend=item.budget_spend, genders=json_load(item.genders), ages=json_load(item.ages), country_codes=json_load(item.country_codes), )