如何判断BigQuery中的表是否正进行INSERT、UPDATE等写入更新操作
解决方案
以下是两种可直接落地的实现方案,你可以根据实际场景选择:
方案1:直接检测目标表是否有正在运行的写入操作
你可以通过查询BigQuery的INFORMATION_SCHEMA.JOBS_BY_PROJECT视图,判断目标表是否存在正在运行的INSERT/UPDATE/LOAD等写入类作业,如果有则跳过本次提取,查询逻辑如下:
SELECT COUNT(1) AS running_write_jobs FROM `region-<你的BQ区域>.INFORMATION_SCHEMA.JOBS_BY_PROJECT` WHERE state = 'RUNNING' AND job_type IN ('QUERY', 'LOAD') AND destination_table.project_id = '<你的项目ID>' AND destination_table.dataset_id = '<你的数据集ID>' AND destination_table.table_id = '<待检测的events_yyyymmdd表名>' AND statement_type IN ('INSERT', 'UPDATE', 'MERGE', 'BULK INSERT', 'CREATE TABLE AS SELECT')
如果查询返回的running_write_jobs大于0,说明表正处于写入状态,跳过本次提取即可。
方案2:判断表最后修改时间(更适合GA事件表场景)
你使用的events_yyyymmdd是Google Analytics导出的日表,这类表的写入操作是集中完成的,不会反复更新,你可以判断表的最后修改时间距离当前时间超过固定阈值(比如1小时),再进行提取,逻辑更简单也更稳定,查询方法如下:
SELECT TIMESTAMP_DIFF(CURRENT_TIMESTAMP(), last_modified_time, HOUR) AS hours_since_last_modified FROM `<你的项目ID>.<你的数据集ID>.__TABLES__` WHERE table_id = '<待检测的events_yyyymmdd表名>'
如果返回的hours_since_last_modified >= 1,说明表已经完成写入,可以安全提取。
现有代码修改参考
你可以在检测到新的日表之后、启动提取之前,增加上述校验逻辑,只有校验通过才执行提取,否则等待下一次轮询再判断,修改后的核心逻辑片段如下(顺便修复了原代码中变量名冲突的问题):
# Check if new daily table exists print(f"Checking if new daily table exists") logging.info(f"Checking if new daily table exists") sql = f'''\ begin create table if not exists `{tempDatasetName}` (id string, insert_date timestamp); create table if not exists `{dailyLogDatasetName}` (table_name string, insert_date timestamp); select distinct table_name from `{informationSchemaTables}` where regexp_contains(table_name, '^events_\\d{{8}}$') and table_name not in (select table_name from `{dailyLogDatasetName}`) order by 1; end ''' query_job = bqClient.query(sql) query_job.result(timeout=7200) for job in bqClient.list_jobs(parent_job=query_job.job_id): if job.statement_type == "SELECT": results = list(job.result().to_arrow().to_pydict().values())[0] # Get missing rows from the daily table if results != []: for table in results: dailyDatasetName = f"{projectName}.{databaseName}.{table}" # 新增:检查表是否已完成写入,阈值可根据实际GA导出完成时间调整 check_sql = f""" SELECT TIMESTAMP_DIFF(CURRENT_TIMESTAMP(), last_modified_time, MINUTE) AS mins_since_modified FROM `{projectName}.{databaseName}.__TABLES__` WHERE table_id = '{table}' """ check_job = bqClient.query(check_sql) check_res = check_job.result() mins_since_modified = list(check_res)[0]['mins_since_modified'] if mins_since_modified < 60: print(f"Daily table {dailyDatasetName} is still being updated, skip this round.") logging.info(f"Daily table {dailyDatasetName} is still being updated, skip this round.") continue print(f"New Daily table {dailyDatasetName} found. Commencing extraction") logging.info(f"New Daily table {dailyDatasetName} found. Commencing extraction") requested_session = types.ReadSession( table=f"projects/{projectName}/datasets/{databaseName}/tables/{table}", data_format=types.DataFormat.ARROW, read_options=types.ReadSession.TableReadOptions(selected_fields=[]) ) read_session = bqStorageClient.create_read_session( parent=f"projects/{projectName}", read_session=requested_session, max_stream_count=streamCount ) streams = read_session.streams obtainedStreamCount = len(streams) print(f"Obtained Streams. Starting {obtainedStreamCount} threads") with concurrent.futures.ThreadPoolExecutor(max_workers=obtainedStreamCount) as executor: futures = [ executor.submit( scroller, bqStorageClient, stream, table, partition, lowerFormat, storageContainerName ) for stream, partition in zip(streams, range(obtainedStreamCount)) ] for f in concurrent.futures.as_completed(futures): f.result(timeout=None) try: streams.clear() except: pass # Insert into the dailylog table the table_name processed so it will skip it on the next run print(f"Inserting {table} into {dailyLogDatasetName}") logging.info(f"Inserting {table} into {dailyLogDatasetName}") sql = f'''\ insert into {dailyLogDatasetName} values ('{table}', current_timestamp()); ''' query_job = bqClient.query(sql) query_job.result(timeout=7200)
内容的提问来源于stack exchange,提问作者Andrei Budaes
相关产品推荐
相关产品推荐

