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

Airflow解析Dataflow输出JSON传递给BigQueryInsertJobOperator审计方案咨询

Airflow Composer 审计日志链路实现方案

核心实现逻辑

整体链路为:DataflowJavaOperator 输出JSON到GCS data目录 → PythonOperator 读取JSON提取目标字段并通过XCom传递 → BigQueryInsertJobOperator 拉取XCom值作为入参调用存储过程写入BQ

详细实现步骤

  • 第一步:使用PythonOperator完成JSON读取与字段提取
    注意Google Cloud Composer 会自动将关联GCS桶的data/目录挂载到Worker本地的/home/airflow/gcs/data/路径,可直接通过本地文件操作读取目标JSON:
import json
from airflow.operators.python import PythonOperator

def extract_audit_info(**context):
    # 替换为你的实际JSON文件路径
    json_path = "/home/airflow/gcs/data/your_dataflow_job_audit.json"
    with open(json_path, "r", encoding="utf-8") as f:
        audit_data = json.load(f)
    # 提取需要的字段,可根据实际JSON结构调整
    job_status = audit_data.get("run_status", "UNKNOWN")
    error_msg = audit_data.get("error_detail", "")
    # 返回值会自动推送至XCom,供下游任务拉取
    return {"job_status": job_status, "error_msg": error_msg}

# 定义提取任务
extract_audit_task = PythonOperator(
    task_id="extract_audit_info",
    python_callable=extract_audit_info,
    dag=dag,
)
  • 第二步:配置BigQueryInsertJobOperator拉取XCom值调用存储过程
    直接通过Airflow模板语法{{ ti.xcom_pull() }}拉取上游返回的字段,作为存储过程的入参:
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator

call_sp_task = BigQueryInsertJobOperator(
    task_id="call_audit_sp",
    project_id="your_gcp_project_id",
    configuration={
        "query": {
            "query": """
                CALL `your_project.your_dataset.your_audit_sp`(
                    @job_status,
                    @error_msg
                );
            """,
            "useLegacySql": False,
            "queryParameters": [
                {
                    "name": "job_status",
                    "parameterType": {"type": "STRING"},
                    "parameterValue": {"value": "{{ ti.xcom_pull(task_ids='extract_audit_info')['job_status'] }}"}
                },
                {
                    "name": "error_msg",
                    "parameterType": {"type": "STRING"},
                    "parameterValue": {"value": "{{ ti.xcom_pull(task_ids='extract_audit_info')['error_msg'] }}"}
                }
            ]
        }
    },
    dag=dag
)
  • 第三步:配置任务依赖
your_dataflow_task >> extract_audit_task >> call_sp_task

注意事项

  • 若JSON文件名包含动态生成的任务ID/时间戳,可在上游DataflowJavaOperator执行完成后将生成的JSON路径推送至XCom,在extract_audit_info函数中拉取该路径即可,无需硬编码
  • 可在提取字段时增加异常捕获逻辑,若JSON读取失败可直接返回失败状态,避免下游存储过程入参异常
  • 若错误信息长度超过BQ表对应字段的长度限制,可在提取阶段做截断处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 19:45:03