基于Composer部署的Airflow如何将文件写入GCS存储桶
Composer Airflow 实现BigQuery读取处理后存GCS操作指南
1. 前置权限配置
- 进入GCP Composer环境详情页,在「环境配置」栏找到当前环境关联的服务账号
- 给该服务账号授予对应权限:
- 目标BigQuery表所在数据集的
roles/bigquery.dataViewer权限,用于读取表数据 - 目标GCS存储桶的
roles/storage.objectCreator权限,用于上传生成的文件
- 目标BigQuery表所在数据集的
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
相关产品推荐
相关产品推荐

