如何在MinIO上传指定文件时触发Apache Airflow DAG?
实现MinIO上传文件触发Airflow DAG(Docker环境)
核心思路
利用MinIO的Webhook通知功能,当指定存储桶有文件上传时,调用Airflow的REST API触发目标DAG。由于两者都运行在Docker容器中,需确保容器间网络互通,同时处理好API认证和请求格式匹配。
步骤1:配置Airflow API权限与触发能力
- 确保使用Airflow 2.x及以上版本(默认启用REST API)
- 生成Airflow API Token:
- 进入Airflow webserver容器:
docker exec -it <airflow-webserver-container-name> bash - 执行命令生成token(替换占位符):
也可通过Airflow UI生成:airflow users create_token -u your-username -e your-email@example.com -f FirstName -l LastNameAdmin > Users > 编辑目标用户 > Generate Token
- 进入Airflow webserver容器:
- 测试API触发:
用curl验证手动触发DAG是否正常(替换占位符):
本地Docker环境中,若MinIO和Airflow在同一自定义网络,可直接用容器名作为curl -X POST "http://<airflow-host>:8080/api/v1/dags/<target-dag-id>/dagRuns" \ -H "Authorization: Bearer <your-airflow-token>" \ -H "Content-Type: application/json" \ -d '{"conf": {}}'<airflow-host>(如airflow-webserver);否则用host.docker.internal映射宿主机端口。
步骤2:配置MinIO Webhook通知
- 使用MinIO客户端
mc配置事件通知:- 配置MinIO服务器别名(替换占位符):
mc alias set myminio http://<minio-container-ip>:9000 <minio-access-key> <minio-secret-key> - 添加Webhook事件规则,监听指定桶的文件上传事件:
mc event add myminio/<target-bucket> arn:minio:sqs::<default>:webhook \ --event "s3:ObjectCreated:*" \ --prefix "指定文件前缀(可选,如data/)" \ --suffix ".csv(可选,指定文件后缀)" \ --endpoint "http://<airflow-host>:8080/api/v1/dags/<target-dag-id>/dagRuns" \ --auth-type "bearer" \ --auth-token "<your-airflow-token>" - 验证规则配置:
mc event list myminio/<target-bucket>
- 配置MinIO服务器别名(替换占位符):
步骤3:(可选)适配MinIO事件格式到Airflow API
MinIO发送的Webhook payload包含事件详情,但Airflow API仅接受{"conf": {}}格式的请求。若需要将MinIO的文件上传信息(如桶名、文件名)传递给DAG,可搭建中转服务转换格式:
- 编写简单的Flask中转服务:
from flask import Flask, request import requests app = Flask(__name__) AIRFLOW_API = "http://airflow-webserver:8080/api/v1/dags/<target-dag-id>/dagRuns" AIRFLOW_TOKEN = "<your-airflow-token>" @app.route('/trigger', methods=['POST']) def trigger_dag(): minio_data = request.json # 提取事件关键信息 event_details = minio_data['Records'][0]['s3'] dag_conf = { "bucket": event_details['bucket']['name'], "file_key": event_details['object']['key'] } # 调用Airflow API headers = { "Authorization": f"Bearer {AIRFLOW_TOKEN}", "Content-Type": "application/json" } requests.post(AIRFLOW_API, json={"conf": dag_conf}, headers=headers) return "OK", 200 if __name__ == '__main__': app.run(host='0.0.0.0', port=5000) - 将该服务打包为Docker镜像并运行,然后把MinIO的Webhook endpoint改为
http://<flask-service-container-name>:5000/trigger
步骤4:验证流程
- 上传符合规则的文件到MinIO指定存储桶
- 查看Airflow UI的
DAG Runs页面,确认是否有新的DAG实例启动 - 若触发失败,可通过以下方式排查:
- 查看MinIO通知配置:
mc admin config get myminio notify_webhook - 查看Airflow webserver日志:
docker logs <airflow-webserver-container-name>
- 查看MinIO通知配置:
内容的提问来源于stack exchange,提问作者nelscodes
相关产品推荐
相关产品推荐

