You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何判断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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 04:48:02