Airflow DAG导出ETL处理后CSV到GCS存储桶无输出且运行过久问题求助
问题原因
1. DAG顶层执行逻辑错误
你将GCS文件读取、DataFrame初始化的代码写在了DAG的全局顶层,而非make_csv函数内部。Airflow的调度器会每隔几秒重新解析一次DAG文件,这就意味着每次解析都会触发一次GCS文件读取操作,不仅会导致调度卡顿、DAG运行耗时极长,更关键的是:Airflow的DAG解析进程和任务执行的Worker进程是完全隔离的,Worker执行make_csv函数时根本拿不到顶层预加载的gas_data变量,会直接触发变量未定义的报错,自然不会生成输出文件。
2. 代码缩进和语法错误
PythonOperator的定义被缩进到了make_csv函数内部,DAG根本无法识别到这个任务节点,任务执行逻辑完全不符合预期logging.info行缩进错误,to_csv的参数quotechar错误使用了HTML转义字符"而非实际的双引号,会直接触发写入异常,导致文件生成失败。
3. 权限和依赖适配问题
如果以上问题修复后仍无输出,需要确认:Airflow Worker运行时使用的GCP服务账号,是否同时拥有源存储桶的读取权限、目标存储桶的写入权限;另外pandas直接写入gs://路径依赖gcsfs的版本适配,版本不匹配也会导致写入静默失败。
4. 大文件处理性能问题
如果读取的CSV文件体积过大,单节点pandas处理会占用极高的内存和CPU,触发Worker超时、OOM崩溃,也会出现运行时间过长、无输出的问题。
修复方案
- 将GCS文件读取逻辑全部移到
make_csv函数内部,确保只有任务执行时才会触发读取操作 - 修正代码缩进、替换
quotechar='"',将Operator定义移到DAG全局作用域 - 给Airflow服务账号授予对应GCS存储桶的
roles/storage.objectAdmin权限,也可以在初始化GCSFileSystem时显式传入Airflow GCP连接对应的凭证,避免硬编码项目和密钥信息 - 大文件场景下建议替换为分块读取、或使用BigQuery做ETL处理后再导出到GCS,避免单节点性能瓶颈
以下是修正后的参考代码:
import gcsfs from airflow.operators.python import PythonOperator from airflow import models import pandas as pd import logging import csv import datetime from airflow.providers.google.common.hooks.base_google import GoogleBaseHook # 仅保留导入、DAG和Operator定义在顶层 def make_csv(): # 读取逻辑移到函数内部,通过Airflow GCP hook获取凭证 gcp_hook = GoogleBaseHook(gcp_conn_id='你的GCP连接ID') credentials = gcp_hook.get_credentials() fs = gcsfs.GCSFileSystem(project='project-goes-here', token=credentials) with fs.open('gs://源存储桶路径/file.csv') as f: gas_data = pd.read_csv(f) gas_data['Asset collection'] = 'Distribution' gas_data['Asset collection address 1'] = 'Distribution' gas_data['Asset collection address 2'] = 'Units1+2 Central City' gas_data['Asset collection address 3'] = 'ind Est' gas_data['Asset collection city'] = 'Coventry' gas_data['Asset collection postcode'] = 'CV6 5RY' gas_data['Multiplication Factor'] = '1.000' gas_data['Conversion Factor'] = '1.022640' output_path = f'gs://目标存储桶路径/Clean_zenos_data_{datetime.datetime.today().strftime("%m%d%Y%H%M%S")}.csv' # 修正quotechar参数 gas_data.to_csv(output_path, index=False, quotechar='"', sep=',', quoting=csv.QUOTE_NONNUMERIC) logging.info('Added Map, Asset collection, Asset collection address and Saved CSV') # Operator定义移到全局作用域 make_csv_function = PythonOperator( task_id='make_csv', python_callable=make_csv )
内容的提问来源于stack exchange,提问作者Mizanur Choudhury
相关产品推荐
相关产品推荐

