如何配置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
相关产品推荐
相关产品推荐

