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

基于BigQuery的CDC Merge存储过程执行审计设计及错误信息捕获方案咨询

我来帮你设计一套适配BigQuery Merge CDC流程的审计方案,同时搞定错误信息的捕获逻辑,结合你用Airflow加载数据的场景,一步步拆解:

一、先搞定审计表的结构设计

首先需要两张核心表:一张记录工作流执行的整体审计信息,另一张专门存错误详情。

1. 工作流执行审计表(比如audit.workflow_execution_logs)

这张表用来记录每次存储过程执行的关键元数据:

CREATE OR REPLACE TABLE audit.workflow_execution_logs (
  execution_id STRING NOT NULL PRIMARY KEY,
  procedure_name STRING NOT NULL,
  source_table STRING NOT NULL, -- 落地表名称
  target_table STRING NOT NULL, -- Vault表名称
  start_time TIMESTAMP NOT NULL,
  end_time TIMESTAMP,
  execution_duration_sec INT64,
  execution_status STRING NOT NULL DEFAULT 'RUNNING', -- 可选值:RUNNING/SUCCESS/FAILED
  row_count_processed INT64, -- Merge操作影响的行数
  created_by STRING NOT NULL,
  created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP()
);

2. 错误日志表(比如audit.error_details_logs)

专门存储执行失败时的错误详情,方便排查:

CREATE OR REPLACE TABLE audit.error_details_logs (
  error_id STRING NOT NULL PRIMARY KEY DEFAULT GENERATE_UUID(),
  execution_id STRING NOT NULL,
  error_message STRING NOT NULL,
  error_stack_trace STRING,
  error_timestamp TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP(),
  affected_table STRING NOT NULL,
  FOREIGN KEY(execution_id) REFERENCES audit.workflow_execution_logs(execution_id)
);
二、存储过程中的审计与错误捕获实现

核心思路是:在存储过程执行前后写入/更新审计表,用BigQuery的异常处理块捕获错误并写入错误日志。

完整存储过程示例

CREATE OR REPLACE PROCEDURE cdc.merge_landing_to_vault(
  IN source_table STRING,
  IN target_table STRING,
  IN execution_id STRING -- 可从Airflow传入dag_run_id,实现端到端追溯
)
BEGIN
  -- 1. 初始化变量
  DECLARE start_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP();
  DECLARE end_time TIMESTAMP;
  DECLARE duration_sec INT64;
  DECLARE status STRING DEFAULT 'SUCCESS';
  DECLARE row_count INT64;
  DECLARE current_user STRING DEFAULT SESSION_USER();

  -- 2. 写入初始审计记录(标记为运行中)
  INSERT INTO audit.workflow_execution_logs(
    execution_id, procedure_name, source_table, target_table, 
    start_time, execution_status, created_by
  )
  VALUES(
    execution_id, 'merge_landing_to_vault', source_table, target_table,
    start_time, 'RUNNING', current_user
  );

  BEGIN
    -- 3. 核心Merge CDC逻辑(这里替换成你的实际Merge语句)
    EXECUTE IMMEDIATE FORMAT("""
      MERGE `%s` AS target
      USING `%s` AS source
      ON target.id = source.id -- 替换成你的匹配键
      WHEN MATCHED THEN
        UPDATE SET -- 替换成你的更新逻辑
          target.last_updated = CURRENT_TIMESTAMP(),
          target.field1 = source.field1,
          target.field2 = source.field2
      WHEN NOT MATCHED THEN
        INSERT (id, field1, field2, created_at, last_updated) -- 替换成你的字段
        VALUES (source.id, source.field1, source.field2, CURRENT_TIMESTAMP(), CURRENT_TIMESTAMP())
    """, target_table, source_table);

    -- 获取Merge影响的行数
    SET row_count = @@ROWCOUNT;
  EXCEPTION
    -- 4. 捕获错误,更新状态并写入错误日志
    WHEN OTHERS THEN
      SET status = 'FAILED';
      INSERT INTO audit.error_details_logs(
        execution_id, error_message, error_stack_trace, affected_table
      )
      VALUES(
        execution_id, ERROR_MESSAGE(), ERROR_STACK_TRACE(), target_table
      );
      -- 重新抛出错误,让Airflow感知到失败(可选,根据你的需求)
      RAISE;
  END;

  -- 5. 更新审计记录的最终状态、结束时间和时长
  SET end_time = CURRENT_TIMESTAMP();
  SET duration_sec = TIMESTAMP_DIFF(end_time, start_time, SECOND);
  
  UPDATE audit.workflow_execution_logs
  SET
    end_time = end_time,
    execution_duration_sec = duration_sec,
    execution_status = status,
    row_count_processed = row_count
  WHERE execution_id = execution_id;
END;
三、Airflow层的补充审计关联

因为你用Airflow加载GCS到落地表,可以把Airflow的dag_run_id作为execution_id传入存储过程,这样就能把数据加载和CDC合并的审计信息关联起来,实现端到端的追溯:

# Airflow任务示例:调用BigQuery存储过程
from airflow.providers.google.cloud.operators.bigquery import BigQueryExecuteQueryOperator

merge_task = BigQueryExecuteQueryOperator(
    task_id="run_merge_cdc_procedure",
    sql="""
        CALL cdc.merge_landing_to_vault(
          'project.dataset.landing_table',
          'project.dataset.vault_table',
          '{{ dag_run.id }}' -- 传入Airflow的dag_run_id作为execution_id
        );
    """,
    use_legacy_sql=False,
    gcp_conn_id="google_cloud_default"
)
四、关键细节说明
  • 错误捕获函数:BigQuery提供了ERROR_MESSAGE()(获取错误描述)、ERROR_STACK_TRACE()(获取调用栈)、ERROR_CODE()(获取错误码)等函数,你可以根据需要加到错误日志表中。
  • 执行时长计算:用TIMESTAMP_DIFF(end_time, start_time, SECOND)精确到秒,如果需要更细粒度可以用MILLISECOND。
  • 幂等性:execution_id用Airflow的dag_run_id或者GENERATE_UUID()保证唯一性,避免重复记录。
  • 权限控制:确保执行存储过程的账号有审计表和错误日志表的读写权限,以及落地表和Vault表的对应权限。

内容的提问来源于stack exchange,提问作者recyclinguy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 19:57:53