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,需满足两个条件:
- 确保
do_xcom_push=True(默认开启,但API响应为空或不可序列化时XCom会为None) - 用
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
相关产品推荐
相关产品推荐

