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

如何在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:
    1. 进入Airflow webserver容器:docker exec -it <airflow-webserver-container-name> bash
    2. 执行命令生成token(替换占位符):
      airflow users create_token -u your-username -e your-email@example.com -f FirstName -l LastName
      
      也可通过Airflow UI生成:Admin > Users > 编辑目标用户 > Generate Token
  • 测试API触发:
    用curl验证手动触发DAG是否正常(替换占位符):
    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": {}}'
    
    本地Docker环境中,若MinIO和Airflow在同一自定义网络,可直接用容器名作为<airflow-host>(如airflow-webserver);否则用host.docker.internal映射宿主机端口。

步骤2:配置MinIO Webhook通知

  • 使用MinIO客户端mc配置事件通知:
    1. 配置MinIO服务器别名(替换占位符):
      mc alias set myminio http://<minio-container-ip>:9000 <minio-access-key> <minio-secret-key>
      
    2. 添加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>"
      
    3. 验证规则配置:mc event list myminio/<target-bucket>

步骤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:验证流程

  1. 上传符合规则的文件到MinIO指定存储桶
  2. 查看Airflow UI的DAG Runs页面,确认是否有新的DAG实例启动
  3. 若触发失败,可通过以下方式排查:
    • 查看MinIO通知配置:mc admin config get myminio notify_webhook
    • 查看Airflow webserver日志:docker logs <airflow-webserver-container-name>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 14:52:37