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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 17:06:03