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

如何在kafka-python中实现带回调的异步提交?

完善Kafka异步提交回调函数实现

针对你的需求,commit_async()的回调函数需要接收两个参数:成功提交时的偏移量字典,以及提交失败时的异常对象(成功时为None)。下面是修改后的完整代码,包含了提交结果的日志记录:

import logging as log

from kafka import KafkaConsumer, TopicPartition
from message_handler_impl import MessageHandlerImpl


def on_commit(offsets: dict[TopicPartition, int], exception: Exception | None):
    if exception:
        # 提交失败时记录错误日志,包含异常详情和尝试提交的偏移量
        log.error("异步提交偏移量失败,异常信息: %s,目标偏移量: %s", 
                  str(exception), offsets, exc_info=True)
    else:
        # 提交成功时可按需记录日志(高流量场景可关闭以减少日志量)
        log.info("异步提交偏移量成功,已提交偏移量: %s", offsets)


class KafkaMessageConsumer:
    def __init__(self, bootstrap_servers: str, topic: str, group_id: str):
        self.bootstrap_servers = bootstrap_servers
        self.topic = topic
        self.group_id = group_id
        # 单分区场景可配置参数提升消费效率
        self.consumer = KafkaConsumer(
            topic,
            bootstrap_servers=bootstrap_servers,
            group_id=group_id,
            enable_auto_commit=False,
            auto_offset_reset='latest',
            # 调大单次拉取字节数,减少网络请求次数
            fetch_max_bytes=52428800,  # 50MB
            max_poll_records=100
        )

    def consume_messages(self, message_handler: MessageHandlerImpl = MessageHandlerImpl()):
        try:
            while True:
                try:
                    msg_pack = self.consumer.poll(timeout_ms=1000)
                    if not msg_pack:
                        continue
                    
                    # 单分区场景简化遍历逻辑
                    for _, messages in msg_pack.items():
                        message_handler.process_messages(messages)

                    # 传入回调执行异步提交
                    self.consumer.commit_async(callback=on_commit)
                except Exception as e:
                    log.error("消费/处理消息出错: %s", e, exc_info=True)

        finally:
            log.info("消费者即将关闭,执行最终同步提交确保偏移量持久化")
            # 关闭前执行同步提交,避免异步提交未完成导致偏移量丢失
            try:
                self.consumer.commit()
            except Exception as e:
                log.error("最终同步提交偏移量失败: %s", e, exc_info=True)
            self.consumer.close()
            log.info("消费者已关闭")


if __name__ == "__main__":
    kafka_consumer = KafkaMessageConsumer("localhost:9092", "test-topic", "test-group")
    kafka_consumer.consume_messages()

关键说明:

  • 回调函数参数:on_commit必须接收offsets和exception两个参数,exception不为None时代表提交失败,需记录异常和对应偏移量用于排查。
  • 单分区优化:因Topic仅一个分区,可调整fetch_max_bytes和max_poll_records参数,减少与Kafka集群的交互次数,提升消费效率。
  • 关闭前同步提交:异步提交可能存在未完成请求,关闭前执行一次同步提交,最大程度避免偏移量丢失。
  • 日志控制:高流量场景下可关闭提交成功的日志,只保留失败日志,减少磁盘IO开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 21:45:29