"""Sample Python Scripts to Generate Daily Bulk Reports.""" from concurrent import futures import csv import os import datetime import logging import sys import urllib.request from time import sleep from google.cloud import bigquery from googleapiclient.discovery import build from oauth2client.service_account import ServiceAccountCredentials from ratelimit import limits from ratelimit import sleep_and_retry ##################### CONSTANT VARIABLES (DO NOT MODIFY) ##################### API_SERVICE_NAME = "youtubemusicanalytics" API_VERSION = "v1" SCOPES = ["https://www.googleapis.com/auth/yt-music-analytics"] CALLS = 3 # number of requests per RATE_LIMIT seconds (MAXIMUM is 3) RATE_LIMIT = 60 # seconds ROW_LIMIT = 100000 # limit 100k rows dump to non-temp table ########################################################################## ######## API CREDENTIAL PARAMETERS TO MODIFY (PLEASE MODIFY IF NEEDED) ####### SERVICE_ACCOUNT_CREDS = ".json" CONTENT_OWNER_ID = "" CONTENT_OWNER = "contentOwners/" + CONTENT_OWNER_ID API_KEY = "" ######## API CREDENTIAL PARAMETERS TO MODIFY (PLEASE MODIFY IF NEEDED) ####### # SERVICE_ACCOUNT_CREDS = "client_secret_767817654343-uvm148m9nl83vimdoqt0t6ud4htmmqvk.apps.googleusercontent.com.json" SERVICE_ACCOUNT_CREDS = "datalytics-etls-244219-89ea4803eaa0.json" # CONTENT_OWNER_ID = "oVAQ96fyMa57I6WozL2gUw" # orch ent CONTENT_OWNER_ID = "J8vAyKuNSYBIN_9RIdxggQ" # orch music CONTENT_OWNER = "contentOwners/" + CONTENT_OWNER_ID API_KEY = "" ####################################################################### ####################################################################### ######### QUERY PARAMETERS TO MODIFY (PLEASE MODIFY IF NEEDED) ############### # Parameters for concurrent requests MAX_ATTEMPTS = 3 MAX_THREAD_NUM = 5 # (MAXIMUM is 10) # Parameters for output destination BigQuery dataset # Set this parameter to empty if output to CSV files # BQ_DATASET_PATH_PREFIX = ".." BQ_DATASET_PATH_PREFIX = "" # Parameters for dataset slice to query YEAR = 2023 MONTH = 7 DAY = 10 TIMEZONE = "America/Los_Angeles" # parameters for enabling filtered rows summary # Type should set for "SUM", "CONSTANT", "SUMMARY_TYPE_UNSPECIFIED" FILTERED_ROW_SUMMARY = { "columns": { "plays": { "type": "SUM", }, "watch_time_sec": { "type": "SUM", }, "users": { "type": "SUM", }, "stream_date": { "type": "CONSTANT", "value": { "value": f"{YEAR:0>4}-{MONTH:0>2}-{DAY:0>2}" } }, "isrc": { "type": "CONSTANT", "value": { "value": "filtered_row_isrc" } }, } } # parameters for customer sql query CUSTOMER_SQL_QUERY = f""" SELECT DATE(DATETIME(TIMESTAMP_MICROS(hour_usec), \"{TIMEZONE}\")) AS stream_date, video_id, primary_asset_id, a.asset_id, a.isrc, gender, age_group, SUM(plays) AS plays, SUM(watch_time_sec) AS watch_time_sec, -- DO NOT REMOVE: required for pre-filtering APPROX_COUNT_DISTINCT(user_id) AS users, FROM ytma.activity LEFT JOIN UNNEST(asset_isrc_rollup) AS a GROUP BY stream_date, video_id, primary_asset_id, asset_id, isrc, gender, age_group -- NOTE: Do not add ORDER BY or LIMIT after GROUP BY, as it will be added in the script """ # parameters for pre-filtering aggregation threshold AGGREGATION_THRESHOLD = 50 # parameters for temp table names TEMP_TABLE_NAME = "" TEMP_FRS_TABLE_NAME = "" TEMP_TABLE_NAME = "temp_efed_table_20230726_01" TEMP_FRS_TABLE_NAME = "temp_efed_table_20230726_frs" # parameters for temp table query texts CREATE_TEMP_TABLE_QUERY_TEXT = f"""CREATE OR REPLACE TEMP TABLE {TEMP_TABLE_NAME} AS( SELECT GENERATE_UUID() AS uuid, * FROM ( {CUSTOMER_SQL_QUERY} HAVING users >= {AGGREGATION_THRESHOLD} ) ORDER BY uuid ) """ CREATE_TEMP_TABLE_FRS_QUERY_TEXT = f"""CREATE OR REPLACE TEMP TABLE {TEMP_FRS_TABLE_NAME} AS( SELECT GENERATE_UUID() AS uuid, * FROM ( {CUSTOMER_SQL_QUERY} HAVING users < {AGGREGATION_THRESHOLD} ) ORDER BY uuid ) """ # parameter for Query Job ID # If any of the query job ids is unset, the script will create a new one CREATE_TEMP_TABLE_QUERY_JOB_ID = "contentOwners//queryJobs/" COUNT_TEMP_TABLE_QUERY_JOB_ID = "contentOwners//queryJobs/" QUERY_JOB_ID = "contentOwners//queryJobs/" CREATE_TEMP_TABLE_FRS_QUERY_JOB_ID = "contentOwners//queryJobs/" QUERY_TEMP_TABLE_FRS_JOB_ID = "contentOwners//queryJobs/" # parameter for logfile name LOG_NAME = f"""bulk_report_{CONTENT_OWNER_ID}_{YEAR}_{MONTH}_{DAY}.log""" ######################################################################## def get_authenticated_service(): # See https://cloud.google.com/docs/authentication/production for all the ways # you can authenticate with a service account. credentials = ServiceAccountCredentials.from_json_keyfile_name( SERVICE_ACCOUNT_CREDS, scopes=SCOPES) return build( API_SERVICE_NAME, API_VERSION, static_discovery=False, # Required for restricted APIs credentials=credentials, cache_discovery=False, developerKey=API_KEY) def create_or_patch_query(service, query_job_id, query_title, query_text, query_parameter_types, filtered_row_summary): """Creates a new query job if exists, otherwise updates an existing query job. Args: service: Authenticaed service of music analytics API. query_job_id: Job id to create or patch the query. query_title: Query title. query_text: SQL format of query text. query_parameter_types: Parameter types for the query. filtered_row_summary: Filtered row summary. Returns: Name of the query job. """ try: query_job = { "query": query_text, "title": query_title, "queryParameterTypes": query_parameter_types, "filteredRowSummary": filtered_row_summary } if not query_job_id: query_job = service.contentOwners().queryJobs().create( parent=CONTENT_OWNER, body=query_job).execute() else: query_job = service.contentOwners().queryJobs().patch( name=query_job_id, body=query_job).execute() return query_job["name"] except Exception as e: logging.error("Create or patch query failed with %s.", str(e)) def generate_run_req(query_parameter_values, bq_data_path): """Generate run query request. Args: query_parameter_values: Parameter values for the query. bq_data_path: Set if data is exported to BigQuery. Returns: Run query request which will be used to run query job. """ run_req = { "spec": { # startDate and endDate are required to set # even though we are retrieving data from temp table. "timeZone": TIMEZONE, "startDate": { "year": YEAR, "month": MONTH, "day": DAY }, "endDate": { "year": YEAR, "month": MONTH, "day": DAY }, "queryParameterValues": query_parameter_values }, "destTable": bq_data_path } return run_req # Set rate limit on run query requests. @sleep_and_retry @limits(calls=CALLS, period=RATE_LIMIT) def run_query_job(service, query_job_id, query_parameter_values, bq_data_path): """Run query job and return the long-running operation. Args: service: Authenticaed service of music analytics API. query_job_id: Job id to create or patch the query. query_parameter_values: Parameter values for the query. bq_data_path: Set if data is exported to BigQuery. Returns: A long-running operation that is the result of an API call. """ run_req = generate_run_req(query_parameter_values, bq_data_path) operation = service.contentOwners().queryJobs().run( name=query_job_id, body=run_req).execute() return operation def poll_op_result(service, operation_name, sleep_time): """Polling the status of the query job until operation.done is true. Args: service: Authenticaed service of music analytics API. operation_name: Long-running operation name. sleep_time: Sleep time to set for polling operation status. Returns: A long-running operation that is the result of an API call. """ try: operation = service.operations().get(name=operation_name).execute() # Poll operation status every minute until operation.done is true. while ("done" not in operation) or (not operation["done"]): sleep(sleep_time) operation = service.operations().get(name=operation_name).execute() return operation except Exception as e: logging.error("Poll operation failed with %s", str(e)) def create_and_run_query(service, query_job_id, query_title, query_text, bq_data_path, query_parameter_types, filtered_row_summary): """Create and run a query job. Args: service: Authenticaed service of music analytics API. query_job_id: Job id to create or patch the query. query_title: Query title. query_text: SQL format of query text. bq_data_path: Set if data is exported to BigQuery. query_parameter_types: Parameter types for the query. filtered_row_summary: Filtered row summary. Returns: A long-running operation that is the result of an API call. Raises: ValueError: If the query_text is not set. RuntimeError: If the operation returned with errors. """ try: if (not query_job_id and not query_text): raise ValueError("Please set query_text parameter value.") query_job_id = create_or_patch_query(service, query_job_id, query_title, query_text, query_parameter_types, filtered_row_summary) logging.info("QueryJob ID %s", (query_job_id)) op = run_query_job(service, query_job_id, None, bq_data_path) logging.info("Operation ID %s ", (op["name"])) operation = poll_op_result(service, op["name"], 60) if (operation is None) or (("error" in operation) and (operation["error"])): raise RuntimeError("Operation %s failed with error: %s" % (str(operation["name"]), str(operation["error"]))) return operation except Exception as e: logging.error("Create and run query failed with %s.", str(e)) sys.exit(1) def query_count_temp_bq(bq_data_path): """Comment out code block below if google.cloud.bigquery package is not imported. Args: bq_data_path: Set if data is exported to BigQuery. Returns: Query result stored in BigQuery. """ client = bigquery.Client() query = f"""SELECT count FROM {bq_data_path}""" try: results = client.query(query) for row in results: count = row["count"] return count except Exception as e: logging.error("Query count temp table in BigQuery failed with %s", str(e)) sys.exit(1) def get_nth_line_from_csv(resp, n): i = 1 while i < n: resp.readline() i += 1 return resp.readline() def get_access_token(): # Get the Access Token for downloading the CSV file. return ServiceAccountCredentials.from_json_keyfile_name( SERVICE_ACCOUNT_CREDS, scopes=SCOPES).get_access_token().access_token def query_count_temp_csv(csv_url): request = urllib.request.Request(csv_url) request.add_header("Authorization", "Bearer " + get_access_token()) csv_file = urllib.request.urlopen(request) # Get first non-header line from csv and convert the result from byte to # string. return get_nth_line_from_csv(csv_file, 2).split(b"\n")[0].decode("utf-8") def send_request(offset, service, query_job_id): """Send request to the YTMA API. Args: offset: A non-negative number of rows to skip in queryJobs.run() requests. service: Authenticaed service of music analytics API. query_job_id: Job id to run the query. Returns: A long-running operation that is the result of an API call. """ for attempt in range(MAX_ATTEMPTS): log_prefix = "[OFFSET %s ATTEMPT %s]: " % (str(offset), str(attempt)) sleep(attempt * 120) try: if BQ_DATASET_PATH_PREFIX: bq_data_path = BQ_DATASET_PATH_PREFIX + str(offset) else: bq_data_path = None # OFFSET is set as type INT64 in query_parameter_types. # However, value in query_parameter_values should always pass as STRING. value = str(offset * ROW_LIMIT) query_parameter_values = {"offset": {"value": value}} run_query_op = run_query_job(service, query_job_id, query_parameter_values, bq_data_path) if run_query_op: logging.info("%s Operation ID %s", log_prefix, run_query_op["name"]) operation = service.operations().get(name=run_query_op["name"]).execute() # Poll operation status every 4 minutes until operation.done is true. # Fetch 100K rows from Temp Table usually takes around 4 minutes. operation = poll_op_result(service, run_query_op["name"], 240) if ("error" not in operation) or (not operation["error"]): logging.info("%s SUCCEED OPERATION %s", log_prefix, operation["name"]) if bq_data_path is None: logging.info("%s Output to CSV file: %s", log_prefix, operation["response"]["uri"]) return operation elif attempt == 2: logging.error("%s FAILED Operation %s:", log_prefix, operation) else: logging.warning("%s FAILED Operation %s:", log_prefix, operation) except Exception as e: if attempt == 2: logging.error("%s FAILED Operation %s: %s", log_prefix, operation, e) else: logging.warning("%s FAILED Operation %s: %s", log_prefix, operation, e) if __name__ == "__main__": logging.basicConfig( # filename=LOG_NAME, level=logging.INFO, # format="%(asctime)s | %(levelname)s | %(message)s") format="[%(asctime)s] p%(process)s %(levelname)s {%(filename)s:%(lineno)d} - %(message)s") service = get_authenticated_service() timestamp = str(int(datetime.datetime.now().timestamp())) if BQ_DATASET_PATH_PREFIX: logging.info("Output query results to BigQuery Dataset.") create_temp_bq_data_path = f"""{BQ_DATASET_PATH_PREFIX}createtemp{YEAR}{MONTH}{DAY}""" count_temp_bq_data_path = f"""{BQ_DATASET_PATH_PREFIX}counttemp{YEAR}{MONTH}{DAY}""" create_temp_frs_bq_data_path = f"""{BQ_DATASET_PATH_PREFIX}createtempfrs{YEAR}{MONTH}{DAY}""" query_temp_frs_bq_data_path = f"""{BQ_DATASET_PATH_PREFIX}querytempfrs{YEAR}{MONTH}{DAY}""" else: logging.info("Output query results to CSV files.") count_temp_bq_data_path = create_temp_bq_data_path = create_temp_frs_bq_data_path = query_temp_frs_bq_data_path = None ######################################################################### # 1. Export query results to a temp table logging.info("-- CREATE TEMP TABLE --") create_temp_query_title = f"create_temp_{timestamp}" create_temp_op = create_and_run_query(service, CREATE_TEMP_TABLE_QUERY_JOB_ID, create_temp_query_title, CREATE_TEMP_TABLE_QUERY_TEXT, create_temp_bq_data_path, None, None) logging.info("-- DONE CREATE TEMP TABLE --") ######################################################################### # 2. Get the number of rows in the temp table logging.info("-- GET TOTAL REQUEST COUNT --") count_temp_query_title = f"'count_temp_{timestamp}" count_temp_query_text = f"SELECT CEIL(COUNT(*) / 100000) AS count FROM tmp.{TEMP_TABLE_NAME}" count_temp_op = create_and_run_query(service, COUNT_TEMP_TABLE_QUERY_JOB_ID, count_temp_query_title, count_temp_query_text, count_temp_bq_data_path, None, None) if BQ_DATASET_PATH_PREFIX: total_request_count = query_count_temp_bq(count_temp_bq_data_path) else: total_request_count = query_count_temp_csv(count_temp_op["response"]["uri"]) logging.info("TOTAL REQUEST COUNT = %s", str(total_request_count)) logging.info("-- DONE TOTAL REQUEST COUNT --") ######################################################################### # 3. Send bulk queries to export temp table externally logging.info("-- GENERATE BULK REPORTS --") query_title = f"get_table_{timestamp}" query_text = f"SELECT * FROM tmp.{TEMP_TABLE_NAME} LIMIT 100000 OFFSET @offset" query_parameter_types = { "offset": { "type": { "type": "int64" }, "description": "temp table offset", "defaultValue": { "value": "0" } } } query_job_id = create_or_patch_query(service, QUERY_JOB_ID, query_title, query_text, query_parameter_types, FILTERED_ROW_SUMMARY) logging.info("QueryJob ID = %s", (query_job_id)) executor = futures.ThreadPoolExecutor(max_workers=MAX_THREAD_NUM) responses = [] for offset in range(int(total_request_count)): sleep(offset * 20) responses.append( executor.submit(send_request, offset, service, query_job_id)) futures.wait(responses) logging.info("-- DONE BULK REPORTS GENERATION --") ######################################################################### # 4. Generate filtered row summary logging.info("-- GENERATE FILTERED_ROW_SUMMARY USERS < %s --", {AGGREGATION_THRESHOLD}) create_temp_frs_query_title = f"""create_temp_frs_{timestamp}""" create_temp_frs_op = create_and_run_query(service, CREATE_TEMP_TABLE_FRS_QUERY_JOB_ID, create_temp_frs_query_title, CREATE_TEMP_TABLE_FRS_QUERY_TEXT, create_temp_frs_bq_data_path, None, None) query_frs_title = f"get_table_user_less_than_{AGGREGATION_THRESHOLD}_{timestamp}" query_frs_text = "SELECT * FROM tmp." + TEMP_FRS_TABLE_NAME op_less_than_aggregation_threshold = create_and_run_query( service, QUERY_TEMP_TABLE_FRS_JOB_ID, query_frs_title, query_frs_text, query_temp_frs_bq_data_path, None, FILTERED_ROW_SUMMARY) logging.info("-- DONE GENERATE FILTERED_ROW_SUMMARY USERS < %s --", {AGGREGATION_THRESHOLD})