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

Airflow DAG部署后无执行:GCS无格式文件转CSV故障求助

问题排查与解决方案

一、基础执行触发问题

你的DAG设置了schedule_interval=None,这意味着它不会自动调度执行,必须手动触发:

  • 登录Airflow UI,找到dag_cartolas_gcs_to_csv DAG
  • 点击右侧的"触发DAG"按钮(播放图标)手动启动一次运行

二、DAG加载与权限验证

  1. 确认DAG是否被Airflow识别

    • 检查Airflow UI的"DAGs"页面,确认dag_cartolas_gcs_to_csv是否在列表中,状态是否为"已启用"(开关为绿色)
    • 若未出现,排查:
      • DAG文件是否放在Airflow指定的dags目录下
      • 文件权限是否允许Airflow进程读取
      • 查看Airflow scheduler日志,排查是否存在语法错误或加载失败信息
  2. GCS权限验证

    • 确保Airflow使用的服务账号拥有cartolas-ahorros-test存储桶的storage.objects.list、storage.objects.get和storage.objects.create权限
    • 可在Airflow服务器上执行以下命令测试权限(替换服务账号密钥路径):
      gcloud auth activate-service-account --key-file=/path/to/your-key.json
      gsutil ls gs://cartolas-ahorros-test/CARTOLAS.*
      

三、代码逻辑优化

1. 过滤已处理的CSV文件

当前代码会列出所有以CARTOLAS.开头的对象,包括已生成的.csv文件,需添加过滤规则:

list_objects = GCSListObjectsOperator(
    task_id='list_objects',
    bucket='cartolas-ahorros-test',
    prefix='CARTOLAS.',
    exclude_pattern='*.csv',  # 排除已处理的CSV文件
    dag=dag,
)

2. 支持多文件处理

当前代码仅处理列表中的第一个对象,若需批量处理所有符合条件的文件,修改download_task:

def download_task(**kwargs):
    bucket_name = 'cartolas-ahorros-test'
    ti = kwargs['ti']
    objects = ti.xcom_pull(task_ids='list_objects')
    if not objects:
        raise ValueError("未找到前缀为'CARTOLAS.'的对象")
    local_paths = []
    for object_name in objects:
        local_path = '/tmp/' + object_name
        download_file_from_gcs(bucket_name, object_name, local_path)
        local_paths.append(local_path)
    return local_paths

同时需修改后续process_task和upload_task,适配列表类型的输入参数。

3. 添加日志输出

在Python函数中添加日志,便于排查执行细节:

import logging

logger = logging.getLogger(__name__)

def procesar_archivo(input_path, output_path):
    logger.info(f"开始处理文件: {input_path}")
    # 原文件处理逻辑...
    logger.info(f"文件处理完成,输出路径: {output_path}")

def download_task(**kwargs):
    logger.info("开始下载GCS文件")
    # 原下载逻辑...
    logger.info(f"下载完成,本地路径列表: {local_paths}")

4. 清理临时文件

处理完成后删除本地临时文件,避免占用磁盘空间:

def upload_task(**kwargs):
    ti = kwargs['ti']
    output_paths = ti.xcom_pull(task_ids='process_file')
    bucket_name = 'cartolas-ahorros-test'
    for output_path in output_paths:
        object_name = os.path.basename(output_path)
        upload_file_to_gcs(bucket_name, object_name, output_path)
        # 删除临时文件
        os.remove(output_path)
        os.remove(output_path.replace('.csv', ''))
        logger.info(f"已清理临时文件: {output_path}")

四、日志排查技巧

  • 手动触发后若无日志,检查Airflow任务日志路径(默认$AIRFLOW_HOME/logs/dag_id/task_id/execution_date)
  • 查看scheduler日志,确认DAG是否被正确调度
  • 查看worker日志,确认任务是否被分配执行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 11:35:55