from airflow.models import Variable import requests import json from airflow.contrib.operators.ssh_operator import SSHOperator from airflow.hooks.base_hook import BaseHook from airflow.contrib.operators.bigquery_operator import BigQueryOperator from airflow.operators.bash_operator import BashOperator def are_downloads_available(date_to_download): ready_to_proceed_for_dwlds = False #Pass sony account ID account_id = Variable.get('apple_account_to_check_dwlds') vendor = get_vendor(account_id) report = get_am_report() payload = get_payload(account_id, vendor, report, date_to_download) headers = {"content-type":"application/x-www-form-urlencoded"} #print('payload == ', payload) r = requests.post(Variable.get('apple_api_url'), headers=headers, data=payload, stream=True) print("Apple response status == ", r.status_code) if r.status_code == 200: ready_to_proceed_for_dwlds = True return ready_to_proceed_for_dwlds def get_vendor(account_id): vendors = json.loads(Variable.get('vendors')) vendor_number = None for vendor in vendors: if vendor['accountNumber'] == account_id: vendor_number = vendor['vendorNumber'] break return vendor_number def get_am_report(): reports = json.loads(Variable.get('reports')) am_report = None for report in reports: if report['name'] == 'amStreams': am_report = report break return am_report def get_payload(account_id, vendor, report, date_to_download): apple_access_token = BaseHook.get_connection('apple_access_token').password payload = 'jsonRequest={"userid":"' + Variable.get('apple_user_id') + '""' payload = payload + ',"accesstoken":"' + apple_access_token payload = payload + '","version":"1.0","mode":"Robot.XML","account":"' + account_id payload = payload + '","queryInput":"[p=Reporter.properties, Sales.getReport, ' + vendor payload = payload + ',' + report['name'] +','+ report['type'] + ',Daily,' + date_to_download.replace('-', '') + ',' + report['version'] + ']"}' return payload def get_partitioned_accounts() : accounts = json.loads(Variable.get('accounts')) list = [] count = 1 accounts_str = '' for account in accounts: accounts_str = accounts_str + account['accountId'] + ' ' count = count + 1 if count % 4 == 0 : list.append(accounts_str) accounts_str = '' return list def get_streams_bq_operator(dag, account_str, date_to_download): filter = '%apple_amStreams_%' + account_str + '_' + date_to_download + '%' + '.tsv' #amStreams bq_load_streams = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_' + account_str + '_streams', sql=''' INSERT `dsp-automation.apple.apple_amStreams` ( Datestamp, Ingest_Datestamp, Ingest_Timestamp, Apple_Identifier, Storefront_Name, Anonymized_Person_ID, Membership_Type, Membership_Mode, Membership_Partner, Postal_Code, Device_Type, Operating_System, UTC_Offset, Action_Type, End_Reason_Type, Offline, Source_of_Stream, Container_Type, Container_Sub_Type, Container_ID, Container_Name, Stream_Timestamp, Stream_Start_Position, Stream_Duration, report_date, report_licensor ) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amStreams WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) return bq_load_streams def get_content_bq_operator(dag, account_str, date_to_download): filter = '%apple_amContent_%' + account_str + '_' + date_to_download + '%' + '.tsv' #amContent bq_load_content = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_' + account_str + '_content', sql=''' INSERT `dsp-automation.apple.apple_amContent` ( Apple_Identifier, ISRC, Title, Artist, Artist_ID, Item_Type, Media_Type, Media_Duration, Vendor_Identifier, Label_Studio_Network, Grid, report_date, report_licensor) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amContent WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) return bq_load_content def set_move_task(dag, date, upstream_task, downstream_task, filters): for filter in filters: move_files = BashOperator( task_id=filter[0], bash_command='gsutil -q -m mv gs://' + Variable.get("gcs_bucket") + '/' + Variable.get("apple_folder_in_bucket") + '/incoming/' + filter[1] + \ ' gs://' + Variable.get("gcs_bucket") + '/' + Variable.get("apple_folder_in_bucket") + '/processed/' + date + ' || true', dag=dag ) move_files.set_upstream(upstream_task) move_files.set_downstream(downstream_task) # This function is for creating the download tasks and it spreads out the SSH tasks based on the country def create_download_tasks(date_to_download, dag, delegate_to_dwlds, stop_vms_task, wait_task): accounts = get_partitioned_accounts() # Download the extracted files and spread out the downloads across 4 VMs count = 0 ssh_conn_name = '' for account in accounts: if count == 0: #ssh_conn_name = 'remote_vm_conn' ssh_conn_name = 'apple-vm-1' else: #ssh_conn_name = 'remote_vm_conn' + str(count) ssh_conn_name = 'apple-vm-1' download_task = SSHOperator( ssh_conn_id=ssh_conn_name, task_id= account.strip().replace(' ', '_') + '_' + '_extracted_dwnlds', command='java -cp collect-analytics-0.0.1-SNAPSHOT-jar-with-dependencies.jar ' + \ 'com.sony.analytics.apple.AppleDownloadExecutor -downloadByAccounts -accountIDs ' + \ account + '-downloadDate ' + date_to_download, dag=dag) download_task.set_upstream(delegate_to_dwlds) for account_str in account.strip().split(): streams_task = get_streams_bq_operator(dag, account_str, date_to_download) content_task = get_content_bq_operator(dag, account_str, date_to_download) download_task.set_downstream(streams_task) download_task.set_downstream(content_task) filter_str = 'apple_amStreams_*' + account_str + '_' + date_to_download + '*' + '.tsv' filter = ('Move_stream_files_' + account_str, filter_str) filters = [filter] set_move_task(dag, date_to_download, streams_task, wait_task, filters) filter_str = 'apple_amContent_*' + account_str + '_' + date_to_download + '*' + '.tsv' filter = ('Move_content_files_' + account_str, filter_str) filters = [filter] set_move_task(dag, date_to_download, content_task, wait_task, filters) #streams_task.set_downstream(wait_task) #content_task.set_downstream(wait_task) if count == 3: count = 0 else: count = count + 1 count = 0 for account in accounts: if count == 0: #ssh_conn_name = 'remote_vm_conn' ssh_conn_name = 'apple-vm-1' else: #ssh_conn_name = 'remote_vm_conn' + str(count) ssh_conn_name = 'apple-vm-1' download_task = SSHOperator( ssh_conn_id=ssh_conn_name, task_id=account.strip().replace(' ', '_') + '_' + '_gzip_dwnlds', command='java -cp collect-analytics-0.0.1-SNAPSHOT-jar-with-dependencies.jar ' + \ 'com.sony.analytics.apple.AppleDownloadExecutor -downloadByAccounts -accountIDs ' + \ account + '-downloadDate ' + date_to_download + ' -downloadGzipsOnly', dag=dag) download_task.set_upstream(wait_task) download_task.set_downstream(stop_vms_task) if count == 3: count = 0 else: count = count + 1 filter = '%apple_amEvent_%' + date_to_download + '%' + '.tsv' #amEvent bq_load_event = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_event', sql=''' INSERT `dsp-automation.apple.apple_amEvent` ( DateStamp, Apple_Identifier, Action_Type, End_Reason_Type, Time_Bucket, UTC_Offset, Membership_Type, Membership_Mode, Membership_Partner, Offline, Storefront_Name, Streams, Ingest_Datestamp, report_date, report_licensor ) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amEvent WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) filter = '%apple_amNonRoyaltyStreams_%' + date_to_download + '%' + '.tsv' #amNonRoyaltyStreams bq_load_non_royalty_streams = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_non_royalty_streams', sql=''' INSERT `dsp-automation.apple.apple_amNonRoyaltyStreams` ( Datestamp, Ingest_Datestamp, Ingest_Timestamp, Apple_Identifier, Storefront_Name, Anonymized_Person_ID, Membership_Type, Membership_Mode, Membership_Partner, Postal_Code, Device_Type, Operating_System, UTC_Offset, Action_Type, End_Reason_Type, Offline, Source_of_Stream, Container_Type, Container_Sub_Type, Container_ID, Container_Name, Stream_Timestamp, Stream_Start_Position, Stream_Duration, report_date, report_licensor ) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amNonRoyaltyStreams WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) filter = '%apple_amContentDemographics_%' + date_to_download + '%' + '.tsv' #amContentDemographics bq_load_content_demographics = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_content_demographics', sql=''' INSERT `dsp-automation.apple.apple_amContentDemographics` ( Datestamp, Ingest_Datestamp, Apple_Identifier, Storefront_Name, Membership_Type, Membership_Mode, Action_Type, Gender, Age_Band, Listeners, Streams, report_date, report_licensor ) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amContentDemographics WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) filter = '%apple_amArtistDemographics_%' + date_to_download + '%' + '.tsv' #amArtistDemographics bq_load_am_artist_demographics = BigQueryOperator( task_id='apple_bqexternal_to_bqinternal_am_artist_demographics', sql=''' INSERT `dsp-automation.apple.apple_amArtistDemographics` ( Datestamp, Ingest_Datestamp, Artist_ID, Storefront_Name, Membership_Type, Membership_Mode, Action_Type, Gender, Age_Band, Listeners, Streams, report_date, report_licensor ) SELECT * EXCEPT(filename), CAST(REGEXP_EXTRACT(filename, r'_([0-9]*-[0-9]*-[0-9]*)_') AS date) AS report_date, REGEXP_EXTRACT(filename, r'_([0-9]*).tsv') AS report_licensor FROM ( SELECT *, _FILE_NAME AS filename FROM apple.external_amArtistDemographics WHERE _FILE_NAME like "''' + filter + '''" ) ''', write_disposition='WRITE_APPEND', use_legacy_sql=False, dag=dag ) bq_load_event.set_upstream(wait_task) filter_str = 'apple_amEvent_*' + date_to_download + '*' + '.tsv' filter = ('Move_event_files', filter_str) filters = [filter] set_move_task(dag, date_to_download, bq_load_event, stop_vms_task, filters) bq_load_non_royalty_streams.set_upstream(wait_task) filter_str = 'apple_amNonRoyaltyStreams_*' + date_to_download + '*' + '.tsv' filter = ('Move_NonRoyaltyStreams_files', filter_str) filters = [filter] set_move_task(dag, date_to_download, bq_load_non_royalty_streams, stop_vms_task, filters) bq_load_content_demographics.set_upstream(wait_task) filter_str = 'apple_amContentDemographics_*' + date_to_download + '*' + '.tsv' filter = ('Move_amContentDemographics_files', filter_str) filters = [filter] set_move_task(dag, date_to_download, bq_load_content_demographics, stop_vms_task, filters) bq_load_am_artist_demographics.set_upstream(wait_task) filter_str = 'apple_amArtistDemographics_*' + date_to_download + '*' + '.tsv' filter = ('Move_amArtistDemographics_files', filter_str) filters = [filter] set_move_task(dag, date_to_download, bq_load_am_artist_demographics, stop_vms_task, filters)