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

如何配置Airflow的Kafka Consumer以消费Topic全部消息?

问题

我希望能够查看Kafka Topic中的所有消息。目前运行以下Airflow代码时,Consumer仅能消费当前DAG执行时Producer写入的消息,无法消费上一次DAG运行时已写入Topic的历史消息,请问该如何配置实现?

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow_provider_kafka.operators.consume_from_topic import ConsumeFromTopicOperator
from airflow_provider_kafka.operators.produce_to_topic import ProduceToTopicOperator
from datetime import datetime, timedelta
import json
import logging

default_args = {
    "owner": "airflow",
    "depend_on_past": False,
    "start_date": datetime(2021, 7, 20),
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}


fruits_test = ["Apple", "Pear", "Peach", "Banana"]
def producer_function():
    for i in fruits_test:
        yield (json.dumps(i), json.dumps(i + i))


consumer_logger = logging.getLogger("airflow")
def consumer_function(message, prefix=None):
    key = json.loads(message.key())
    value = json.loads(message.value())
    consumer_logger.info(f"{prefix} {message.topic()} @ {message.offset()}; {key} : {value}")
    return


with DAG(
    "kafka_DAG",
    default_args=default_args,
    description="KafkaOperators",
    schedule_interval=None,
    start_date=datetime(2021, 1, 1),
    catchup=False,
    tags=["Test_DAG"],
) as dag:

    t1 = ProduceToTopicOperator(
        task_id="produce_to_topic",
        topic="topictest",
        producer_function=producer_function,
        kafka_config={"bootstrap.servers": ":9092"},
    )


    t2 = ConsumeFromTopicOperator(
        task_id="consume_from_topic",
        topics=["topictest"],
        apply_function=consumer_function,
        apply_function_kwargs={"prefix": "consumed:::"},
        consumer_config={
            "bootstrap.servers": ":9092",
            "group.id": "test-consumer-group",
            "enable.auto.commit": False,
            "auto.offset.reset": "earliest",
        },
        commit_cadence="end_of_batch",
        max_messages=10,
        max_batch_size=2,
    )
解决方案

核心问题是你的消费者组已经提交过偏移量记录——哪怕设置了auto.offset.reset: earliest,只要Kafka里存在这个消费者组的已提交偏移,就会从上次的位置继续消费,不会从头读取历史消息。要实现消费所有历史消息,有两种可靠的配置方式:

方法一:使用动态生成的消费者组ID

每次运行DAG时使用不同的消费者组ID,Kafka会把它当成全新的消费者,触发auto.offset.reset规则,直接从Topic最开始的位置消费所有消息。修改consumer_config里的group.id即可:

"group.id": f"test-consumer-group-{datetime.now().strftime('%Y%m%d%H%M%S')}",

方法二:禁止偏移量提交

如果想固定用同一个消费者组ID,那就完全禁止提交偏移量,这样每次启动消费者时都没有已记录的偏移,会从头消费。把commit_cadence改成never就行:

commit_cadence="never",

修改后的完整代码(方法一示例)

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow_provider_kafka.operators.consume_from_topic import ConsumeFromTopicOperator
from airflow_provider_kafka.operators.produce_to_topic import ProduceToTopicOperator
from datetime import datetime, timedelta
import json
import logging

default_args = {
    "owner": "airflow",
    "depend_on_past": False,
    "start_date": datetime(2021, 7, 20),
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}


fruits_test = ["Apple", "Pear", "Peach", "Banana"]
def producer_function():
    for i in fruits_test:
        yield (json.dumps(i), json.dumps(i + i))


consumer_logger = logging.getLogger("airflow")
def consumer_function(message, prefix=None):
    key = json.loads(message.key())
    value = json.loads(message.value())
    consumer_logger.info(f"{prefix} {message.topic()} @ {message.offset()}; {key} : {value}")
    return


with DAG(
    "kafka_DAG",
    default_args=default_args,
    description="KafkaOperators",
    schedule_interval=None,
    start_date=datetime(2021, 1, 1),
    catchup=False,
    tags=["Test_DAG"],
) as dag:

    t1 = ProduceToTopicOperator(
        task_id="produce_to_topic",
        topic="topictest",
        producer_function=producer_function,
        kafka_config={"bootstrap.servers": ":9092"},
    )

    t2 = ConsumeFromTopicOperator(
        task_id="consume_from_topic",
        topics=["topictest"],
        apply_function=consumer_function,
        apply_function_kwargs={"prefix": "consumed:::"},
        consumer_config={
            "bootstrap.servers": ":9092",
            # 动态生成消费者组ID,每次运行都是新组
            "group.id": f"test-consumer-group-{datetime.now().strftime('%Y%m%d%H%M%S')}",
            "enable.auto.commit": False,
            "auto.offset.reset": "earliest",
        },
        commit_cadence="end_of_batch",
        max_messages=10,
        max_batch_size=2,
    )

t1 >> t2

额外提醒

  • 原代码漏了datetime、timedelta、json、logging的导入,已经补上,不然运行会报错。
  • 还要确认你的Kafka Topic保留策略允许历史消息存在(默认保留7天),如果旧消息已经被清理,再怎么配置也读不到。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 00:40:38