基于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
相关产品推荐
相关产品推荐

