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

基于Composer部署的Airflow如何将文件写入GCS存储桶

Composer Airflow 实现BigQuery读取处理后存GCS操作指南

1. 前置权限配置

  • 进入GCP Composer环境详情页,在「环境配置」栏找到当前环境关联的服务账号
  • 给该服务账号授予对应权限:
    • 目标BigQuery表所在数据集的roles/bigquery.dataViewer权限,用于读取表数据
    • 目标GCS存储桶的roles/storage.objectCreator权限,用于上传生成的文件

2. DAG代码示例

你可以直接使用Airflow GCP Provider内置的Hook来处理权限和API调用,不需要手动配置凭证,以下是可直接修改使用的完整DAG代码:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
from airflow.providers.google.cloud.hooks.gcs import GCSHook
from datetime import datetime
import pandas as pd

# 自定义参数,替换为你自己的配置
BQ_PROJECT = "你的GCP项目ID"
BQ_TABLE = "你的数据集.你的表名"
GCS_BUCKET = "目标存储桶名称"
GCS_SAVE_PATH = "要保存到桶里的文件路径,比如output/processed_data.csv"
LOCAL_TMP_PATH = "/tmp/processed_data.csv"

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": datetime(2024, 1, 1),
    "retries": 1
}

def process_bq_data():
    # 步骤1:从BigQuery读取数据转DataFrame
    bq_hook = BigQueryHook(use_legacy_sql=False, project_id=BQ_PROJECT)
    sql = f"SELECT * FROM `{BQ_PROJECT}.{BQ_TABLE}`"
    df = bq_hook.get_pandas_df(sql=sql)

    # 步骤2:你的自定义处理逻辑,替换为你现有的pandas处理代码
    # 示例:新增一列、过滤行等
    df["process_time"] = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
    df = df[df["status"] == "valid"]

    # 步骤3:保存为本地临时文件
    df.to_csv(LOCAL_TMP_PATH, index=False, encoding="utf-8")

    # 步骤4:上传到GCS存储桶
    gcs_hook = GCSHook()
    gcs_hook.upload(
        bucket_name=GCS_BUCKET,
        object_name=GCS_SAVE_PATH,
        filename=LOCAL_TMP_PATH
    )

with DAG(
    dag_id="bq_process_to_gcs",
    default_args=default_args,
    schedule_interval=None,  # 测试用设为None,后续可以修改为需要的调度周期,比如"0 1 * * *"每天凌晨1点运行
    catchup=False,
    tags=["bigquery", "gcs", "data_process"]
) as dag:
    process_task = PythonOperator(
        task_id="process_bq_data_task",
        python_callable=process_bq_data
    )

# 你可以在此处添加更多上下游任务,DAG会按依赖顺序执行

3. 依赖配置

如果你的pandas处理逻辑用到其他第三方包,或者默认环境缺少依赖,需要在Composer环境的「PYPI包」配置页添加对应包,常用依赖包括:

  • pandas
  • pyarrow(加速BigQuery数据转DataFrame的效率)

4. 部署测试

  • 将修改好的DAG文件(后缀为.py)上传到Composer环境关联的GCS存储桶的dags/目录下
  • 等待1-2分钟后刷新Airflow UI,即可找到ID为bq_process_to_gcs的DAG,手动触发即可测试运行

常见注意事项

  • 处理超过10万行的大表时建议分批读取处理,避免worker节点内存不足
  • 不要在代码中硬编码GCP密钥,Hook会自动读取Composer默认服务账号的权限,安全性更高
  • 临时文件存储在worker的/tmp目录下,任务运行结束后会自动清理,不需要手动删除

内容的提问来源于stack exchange,提问作者Carlos Eduardo Abrantes

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 00:09:03