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

Airflow DAG:S3KeySensor匹配结果XCom传递及多文件并行处理问题

解决Airflow DAG监控MinIO并并行触发API的问题

一、验证MinIO连接与S3KeySensor路径解析

1. 确认Airflow S3连接配置

在Airflow UI的连接管理中,检查你的MinIO连接:

  • Conn Type选择S3
  • Extra字段必须填入{"endpoint_url": "http://s3:9000"}(匹配docker-compose内的MinIO服务地址)
  • Access Key和Secret Key填写MinIO的账号密码

2. 正确配置S3KeySensor

要监控桶内新增的.jp2文件,需开启通配符匹配,示例代码:

from airflow.providers.amazon.aws.sensors.s3_key import S3KeySensor

wait_for_jp2 = S3KeySensor(
    task_id='wait_for_new_jp2',
    bucket_name='your-target-bucket',
    bucket_key='*.jp2',
    wildcard_match=True,
    aws_conn_id='your-minio-conn-id',
    poke_interval=30,  # 每30秒检查一次
)

查看Sensor日志时,若日志中显示的endpoint是http://s3:9000则配置正确;若显示AWS默认的s3.amazonaws.com,说明连接的Extra参数配置错误。

二、解决SimpleHttpOperator XCom返回None的问题

问题根源

你误用了SimpleHttpOperator来获取MinIO文件列表——获取S3/MinIO文件列表应使用S3Hook,而非SimpleHttpOperator。若确实需要用SimpleHttpOperator调用外部API并返回XCom,需满足两个条件:

  1. 确保do_xcom_push=True(默认开启,但API响应为空或不可序列化时XCom会为None)
  2. 用response_filter处理非JSON格式的响应

正确获取文件列表的方式

用PythonOperator调用S3Hook列出符合条件的文件并写入XCom:

from airflow.operators.python import PythonOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook

def fetch_jp2_files(**context):
    s3_hook = S3Hook(aws_conn_id='your-minio-conn-id')
    all_keys = s3_hook.list_keys(bucket_name='your-target-bucket')
    # 过滤出.jp2后缀的文件
    jp2_keys = [key for key in all_keys if key.endswith('.jp2')]
    # 写入XCom供后续任务调用
    context['ti'].xcom_push(key='jp2_file_list', value=jp2_keys)
    return jp2_keys

list_jp2_task = PythonOperator(
    task_id='list_jp2_files',
    python_callable=fetch_jp2_files,
    provide_context=True,
    do_xcom_push=True,
)

三、实现每个文件并行触发API调用

使用Airflow 2.2+支持的动态任务映射(Dynamic Task Mapping),通过expand方法为每个文件生成独立的API调用任务,实现并行执行:

from airflow.providers.http.operators.http import SimpleHttpOperator

# 定义API调用模板
api_task = SimpleHttpOperator(
    task_id='call_processing_api',
    http_conn_id='your-api-conn-id',  # 提前配置API的连接信息
    endpoint='/your-api-endpoint',
    method='POST',
    data='{"s3_key": "{{ params.file_key }}"}',
    headers={"Content-Type": "application/json"},
    do_xcom_push=True,
)

# 基于文件列表动态生成并行任务
parallel_api_calls = api_task.expand(
    params=[{"file_key": key} for key in list_jp2_task.output['jp2_file_list']]
)

完整DAG示例

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.amazon.aws.sensors.s3_key import S3KeySensor
from airflow.operators.python import PythonOperator
from airflow.providers.http.operators.http import SimpleHttpOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'minio_jp2_monitor',
    default_args=default_args,
    description='Monitor MinIO for new JP2 files and trigger parallel API calls',
    schedule_interval=timedelta(hours=1),
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=['minio', 'jp2', 'api'],
) as dag:

    # 1. 等待桶内出现新的.jp2文件
    wait_for_files = S3KeySensor(
        task_id='wait_for_new_jp2',
        bucket_name='your-bucket',
        bucket_key='*.jp2',
        wildcard_match=True,
        aws_conn_id='minio_conn',
        poke_interval=30,
    )

    # 2. 获取所有符合条件的.jp2文件
    def list_jp2_files(**context):
        s3_hook = S3Hook(aws_conn_id='minio_conn')
        all_keys = s3_hook.list_keys(bucket_name='your-bucket')
        jp2_keys = [key for key in all_keys if key.endswith('.jp2')]
        context['ti'].xcom_push(key='jp2_file_list', value=jp2_keys)
        return jp2_keys

    list_files_task = PythonOperator(
        task_id='list_jp2_files',
        python_callable=list_jp2_files,
        provide_context=True,
    )

    # 3. 并行触发API调用
    api_call_task = SimpleHttpOperator(
        task_id='trigger_processing_api',
        http_conn_id='api_conn',
        endpoint='/process-jp2',
        method='POST',
        data='{"s3_key": "{{ params.file_key }}"}',
        headers={"Content-Type": "application/json"},
    )

    parallel_tasks = api_call_task.expand(
        params=[{"file_key": key} for key in list_files_task.output['jp2_file_list']]
    )

    # 设置任务依赖
    wait_for_files >> list_files_task >> parallel_tasks

额外注意事项

  • 确保Airflow版本≥2.2,动态任务映射是2.2版本引入的功能
  • 若要避免重复处理文件,可在API调用成功后将文件移动到“已处理”目录,或用数据库记录已处理的S3密钥
  • 检查API连接的地址、认证信息是否配置正确

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 21:07:20