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

如何实现Airflow中Kafka监听与DAG自动化的流式数据管道?

Airflow Kafka监听器无法触发下游DAG的解决方案

需求

  • 监听Kafka主题,消息数量与到达时间不确定
  • 监听器持续运行,消息到达时自动启动数据转换+MongoDB写入任务,处理完成后结束任务

现有代码

DAG 1 (监听Kafka的DAG)

from airflow.decorators import dag, task
from airflow.operators.python_operator import PythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.utils.dates import days_ago
from datetime import datetime, timedelta
from kafka import KafkaConsumer

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

@dag(default_args=default_args, description='Streaming pipeline', start_date=days_ago(1), schedule_interval=None, catchup=False)
def job_offer_matching():
    
    def trigger_downstream_dag(message_from_topic):
        downstream_dag_id = 'downstream_kafka_messages'
        print(f"encountered {message_from_topic} - executing external dag!")
        
        trigger_dag_run_op = TriggerDagRunOperator(
            task_id="trigger_downstream_kafka_messages",
            trigger_dag_id=downstream_dag_id,
            conf={"message": message_from_topic},
            wait_for_completion=True
        )
        return trigger_dag_run_op
    
    @task
    def listen_kafka_topic():
        # Set up the Kafka consumer
        consumer = KafkaConsumer(
            'test-topic',  # Topic name
            bootstrap_servers='hide:9092',  # Kafka broker address
            security_protocol='SASL_SSL',
            sasl_mechanism='PLAIN',
            sasl_plain_username='hide',
            sasl_plain_password='hide',
            group_id='my-group',  # Consumer group ID
            auto_offset_reset='earliest',  # Start from the earliest offset
            enable_auto_commit=True,  # Automatically commit offsets
            value_deserializer=lambda x: x.decode('utf-8')  # Decode message value
        )
        
        # Retrieve the data sources from the collection
        for message in consumer:
            message_from_topic = message.value
            print(message_from_topic)
            trigger_downstream_dag(message_from_topic)
    
    listen_kafka_topic()

job_offer_matching()

DAG 2 (下游处理DAG)

from airflow.decorators import dag, task
from datetime import datetime, timedelta
from airflow.utils.dates import days_ago
from airflow.operators.python_operator import PythonOperator

default_args = {
    'owner': 'me',
}

@dag(default_args=default_args, description='Streaming pipeline Downstream', start_date=days_ago(1), schedule_interval=None, catchup=False)
def downstream_kafka_messages():
    
    @task
    def message(dag_run=None):
        print(dag_run.conf.get("message"))

    message()

downstream_kafka_messages()

问题原因

你当前的代码无法触发下游DAG,核心问题在于:

  1. TriggerDagRunOperator的使用方式错误:该Operator是在DAG定义阶段声明的任务节点,而非在任务运行时动态创建执行。在listen_kafka_topic任务内部创建并返回Operator,不会实际触发下游DAG,因为Airflow的任务运行阶段仅执行Python逻辑,不处理DAG结构定义。
  2. 长时任务的潜在风险:listen_kafka_topic是持续运行的循环任务,Airflow默认的任务超时机制可能中断该任务,且持续占用Worker资源不符合Airflow的批处理调度模型。

修复方案

方案1:使用Airflow API触发下游DAG

直接在监听任务中调用Airflow的本地API触发下游DAG,替代动态创建Operator的方式:

修改DAG 1的listen_kafka_topic任务:

@task(execution_timeout=timedelta(days=1))  # 设置足够长的超时时间,避免任务被中断
def listen_kafka_topic():
    from airflow.api.client.local_client import Client
    # 初始化Airflow本地客户端
    client = Client(None, None)
    
    # Set up the Kafka consumer
    consumer = KafkaConsumer(
        'test-topic',  # Topic name
        bootstrap_servers='hide:9092',  # Kafka broker address
        security_protocol='SASL_SSL',
        sasl_mechanism='PLAIN',
        sasl_plain_username='hide',
        sasl_plain_password='hide',
        group_id='my-group',  # Consumer group ID
        auto_offset_reset='earliest',  # Start from the earliest offset
        enable_auto_commit=False,  # 关闭自动提交,确保消息处理完成后再提交偏移量
        value_deserializer=lambda x: x.decode('utf-8')  # Decode message value
    )
    
    # 循环监听消息
    for message in consumer:
        message_from_topic = message.value
        print(message_from_topic)
        
        # 调用API触发下游DAG
        client.trigger_dag(
            dag_id='downstream_kafka_messages',
            conf={'message': message_from_topic},
            execution_date=datetime.now()
        )
        
        # 手动提交偏移量,确保消息已被处理
        consumer.commit()

方案2:使用Airflow异步Kafka传感器(推荐)

如果你的Airflow版本≥2.2.0,建议使用AsyncKafkaSensor实现无阻塞的消息监听,避免持续占用Worker资源:

修改后的DAG 1(使用异步传感器)

from airflow.decorators import dag, task
from airflow.sensors.kafka import AsyncKafkaSensor
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.utils.dates import days_ago
from datetime import datetime, timedelta

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

@dag(default_args=default_args, description='Streaming pipeline', start_date=days_ago(1), schedule_interval=None, catchup=False)
def job_offer_matching():
    # 异步监听Kafka消息
    kafka_sensor = AsyncKafkaSensor(
        task_id='listen_kafka_topic',
        kafka_config={
            'bootstrap_servers': 'hide:9092',
            'security_protocol': 'SASL_SSL',
            'sasl_mechanism': 'PLAIN',
            'sasl_plain_username': 'hide',
            'sasl_plain_password': 'hide',
            'group_id': 'my-group',
            'auto_offset_reset': 'earliest',
        },
        topics=['test-topic'],
        poke_interval=5,  # 检查间隔(秒)
        timeout=timedelta(days=7),
    )

    # 触发下游DAG的任务
    trigger_downstream = TriggerDagRunOperator(
        task_id='trigger_downstream_kafka_messages',
        trigger_dag_id='downstream_kafka_messages',
        conf={'message': '{{ task_instance.xcom_pull(task_ids="listen_kafka_topic") }}'},
        wait_for_completion=False
    )

    kafka_sensor >> trigger_downstream

job_offer_matching()

这种方式下,传感器会挂起等待Kafka消息,收到消息后将消息内容通过XCom传递给触发任务,再触发下游DAG。传感器会自动重复运行,持续监听新消息。

下游DAG优化(添加数据转换与MongoDB写入)

修改DAG 2的任务,实现数据转换和MongoDB写入逻辑:

from airflow.decorators import dag, task
from datetime import datetime, timedelta
from airflow.utils.dates import days_ago
from pymongo import MongoClient

default_args = {
    'owner': 'me',
}

@dag(default_args=default_args, description='Streaming pipeline Downstream', start_date=days_ago(1), schedule_interval=None, catchup=False)
def downstream_kafka_messages():
    
    @task
    def process_and_write_to_mongo(dag_run=None):
        # 获取上游传递的消息
        raw_message = dag_run.conf.get('message')
        
        # 数据转换逻辑示例
        transformed_data = {
            'received_at': datetime.now(),
            'content': raw_message,
            # 添加你的转换逻辑
        }
        
        # 连接MongoDB并写入数据
        mongo_client = MongoClient('mongodb://your-mongo-host:27017/')  # 替换为你的MongoDB地址
        db = mongo_client['your_database']
        collection = db['your_collection']
        collection.insert_one(transformed_data)
        
        print(f"Successfully processed and wrote message: {raw_message}")

    process_and_write_to_mongo()

downstream_kafka_messages()

注意事项

  1. 权限配置:确保运行Airflow Worker的用户有调用Airflow API触发DAG的权限。
  2. 偏移量管理:关闭自动提交偏移量,手动在消息处理完成后提交,避免消息丢失或重复处理。
  3. 资源占用:若使用方案1的长时循环任务,需确保Worker有足够资源,且任务超时时间设置合理;方案2的异步传感器资源占用更低,更适合生产环境。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 21:14:56