ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka消费者offset机制详解:从原理到librdkafka实践

Kafka消费者offset机制详解:从原理到librdkafka实践 1. 从“游标”说起offset到底是什么如果你刚把Kafka的生产者代码跑通开始琢磨消费者遇到的第一个绕不开的概念八成就是offset。我当年第一次看到“offset”这个词脑子里冒出来的是SQL里的OFFSET分页关键字以为就是一个简单的“跳过多少条”的偏移量。等真正写完消费者代码、跑起来之后才发现Kafka里的offset是一个更底层、更重要的东西——它是消费者在分区内的读取位置标记。简单做个类比Kafka的每个分区就像一本只允许顺序翻的账本每一页每一条消息都有一个页码这个页码就是offset。消费者读到哪里账本上就插一个书签。下次再读从书签位置往后继续翻就行。这个书签就是消费者自己维护的读取位置。在C/C场景下我们通常通过librdkafka来写消费者。librdkafka暴露出来的一堆回调函数、配置项里和offset相关的占了很大比重。理解不透彻的话写出来的消费者代码很容易出现两个极端要么重复消费大量消息要么把应该消费的消息悄悄丢掉。这两种问题在日志系统、订单处理这类场景里都是不可接受的。这篇文章我在学习日记系列里排在Kafka的第三篇专门讲消费者代码的第一部分。会重点拆解offset的语义、它存在哪里、怎么提交、有哪些坑。下一篇再展开讲消费者组的协调机制和负载均衡。2. offset的存储机制和历史演进2.1 旧版offset存放在ZooKeeperKafka早期的版本消费者消费完消息之后会把当前读到的offset提交到ZooKeeper的节点上。ZooKeeper充当了一个分布式协调存储的角色。但这种方式在高频提交的场景下暴露出了问题——ZooKeeper并不适合承载大量的写操作。你想象一下一个高吞吐的消费集群每秒可能提交成千上万次offset更新每次都去写ZooKeeper节点ZooKeeper的写入性能很快会成为瓶颈。我记得当时做运维的时候ZooKeeper集群的写请求量有一大半都是offset提交带来的非常被动。2.2 新版内置__consumer_offsets主题从Kafka 0.9版本开始官方把offset的存储挪到了Kafka自己内部的一个特殊主题里这个主题叫做__consumer_offsets。它本质上是一个普通的Kafka主题只不过有固定的50个分区默认值可以通过配置调整专门用来保存消费者组的提交记录。这样做的好处很明显Kafka本身就是为高吞吐、高并发写入设计的把offset提交变成普通的消息写入性能和可靠性反而都提升了。而且消费端的offset数据跟Kafka集群在一起天然支持复制和容灾。当某个broker挂了offset数据会从副本恢复不像以前还得担心ZooKeeper和Kafka集群之间数据一致性。__consumer_offsets是怎么决定一条提交记录去哪里的很简单按消费者组名做哈希再对50取模算出一个分区号。举个例子如果你要排查某个消费者组的提交情况可以先算出来它在哪个分区然后用命令行工具直接查看这个分区的消息内容。这是我排查offset相关问题时最常用的手段之一。2.3 消费者端的“两道offset”这里必须强调一个细节其实消费者端存在两个offset很容易搞混。第一个是已提交的offset存储在__consumer_offsets里是消费者组在broker侧的“官方记录”。它表示“这个组已经确认消费到哪里了”。第二个是当前消费位置保存在消费者客户端内存里表示“我这个消费者实例现在读到了哪条消息”。正常情况下你消费完一批消息、手动提交之后这两个值是一致的。但如果没来得及提交就宕机了或者自动提交的间隔还没到那这两个值就会出现差异。新消费者实例启动时broker返回的是已提交的offset而不是内存里的位置。这也解释了为什么消费者崩了之后重新上线会重复消费一部分消息。很多人问能不能不提交offset全靠内存位置续读不行。一旦消费者进程退出内存里的位置就没了。除非用seek接口在启动时手动指定位置但那等于重新造一套持久化机制属于自找麻烦。3. C/C开发者必须掌握的offset相关配置3.1 auto.offset.reset消费起点从哪开始新手最容易踩的坑就是auto.offset.reset这个配置项它决定了消费者在没有已提交offset可参考的时候应该从哪个位置开始读。这里要注意一个容易误会的点它只在没有已提交offset时生效。如果消费者组之前已经提交过offset那就从提交的位置继续跟这个配置没关系。什么情况下会走到“没有已提交offset”这条路常见的有三种消费者组第一次订阅这个主题之前从没消费过消费者组在服务端对应的offset记录被删除了比如超过offsets.retention.minutes保留期还没提交给消费者指定了一个全新的group.id在librdkafka里这个配置项有两个主要值配置值行为适用场景earliest从分区最早的可用消息开始读从头消费离线计算全量数据、数据修复、需要补历史数据latest从分区最新的消息开始读只消费新到的实时处理、日志监控、正常业务消费我自己在测试环境经常用earliest因为方便反复验证。生产环境的实时任务则用latest居多。但如果要求“从某个时间点开始消费”需要配合时间戳查询来定位offset单纯靠auto.offset.reset做不到。3.2 enable.auto.commit自动提交的陷阱librdkafka默认开启自动提交配置项是enable.auto.committrue。自动提交的逻辑是每隔auto.commit.interval.ms默认5000毫秒把当前消费到的位置提交上去。省事是真的省事但陷阱也很明显自动提交的时机是“当前读取到的位置”而不是“处理完成的位置”。假设你拉取了一批消息在内存里处理处理到一半还没结束5秒的间隔到了客户端把位置提交上去了。这时候如果进程崩溃重启后就会从提交的位置开始读那批你正在处理但还没处理完的消息就丢了。反过来也有问题。自动提交的间隔内如果频繁发生rebalance消费者组的分区重新分配新的消费者会从已提交的位置开始读未处理完的消息也会被重复消费。我现在的习惯是生产环境一律关闭自动提交改为手动提交。除非业务完全不关心重复消费和丢失比如一些纯日志抽样场景否则自动提交就是在赌运气。手动提交的代码量并不大可以在确认消息处理成功之后再提交把“至少一次”的语义握在自己手里。3.3 手动提交commitSync与commitAsync的取舍手动提交在librdkafka里主要有两个接口rd_kafka_commit()同步提交调用后等待broker确认提交成功才返回。rd_kafka_commit_async()异步提交调用后立即返回提交结果通过回调通知。同步提交的好处是确定性高失败能立刻知道。坏处是阻塞调用期间无法继续拉取和处理消息吞吐量会受影响。异步提交不阻塞但增加了编程复杂度——提交失败了怎么办要自己记日志、补偿重试。我在实践里用的策略是核心链路同步提交普通链路异步提交加失败重试。比如订单支付结果的处理要求不能丢那就同步提交宁可慢一点也不冒风险。像用户行为日志汇总这种丢失几条影响不大的用异步提交追求吞吐。还有一个小细节手动提交的时候可以提交全局位置也可以只提交某个分区的位置。librdkafka提供了rd_kafka_offsets_store()配合rd_kafka_commit()用的方式能把每个分区的offset分别管理按需提交。比如某个分区拉取到的数据量特别大其他分区处理完了可以先提交处理完的分区位置。这种精细化操作在高吞吐场景下比统一提交更高效也更安全。4. 消费者代码实操offset相关的核心API4.1 创建消费者并订阅主题在C/C环境里使用librdkafka写消费者第一步是创建配置对象、设置参数。我常用的配置代码如下#include librdkafka/rdkafka.h rd_kafka_t *create_consumer(const char *brokers, const char *group_id) { char errstr[512]; rd_kafka_conf_t *conf rd_kafka_conf_new(); rd_kafka_conf_set(conf, bootstrap.servers, brokers, errstr, sizeof(errstr)); rd_kafka_conf_set(conf, group.id, group_id, errstr, sizeof(errstr)); // 手动提交 rd_kafka_conf_set(conf, enable.auto.commit, false, errstr, sizeof(errstr)); // 没有已提交offset时从最早位置开始 rd_kafka_conf_set(conf, auto.offset.reset, earliest, errstr, sizeof(errstr)); // 关闭自动提交后需要设置提交回调 rd_kafka_conf_set_offset_commit_cb(conf, offset_commit_cb); rd_kafka_t *rk rd_kafka_new(RD_KAFKA_CONSUMER, conf, errstr, sizeof(errstr)); if (!rk) { fprintf(stderr, 创建消费者失败: %s\n, errstr); return nullptr; } return rk; }创建完消费者之后第二步就是订阅主题static void subscribe_topic(rd_kafka_t *rk, const char *topic) { rd_kafka_topic_partition_list_t *topics; topics rd_kafka_topic_partition_list_new(1); rd_kafka_topic_partition_list_add(topics, topic, RD_KAFKA_PARTITION_UA); // 订阅接口第二个参数是主题列表 rd_kafka_resp_err_t err rd_kafka_subscribe(rk, topics); if (err ! RD_KAFKA_RESP_ERR_NO_ERROR) { fprintf(stderr, 订阅主题失败: %s\n, rd_kafka_err2str(err)); } rd_kafka_topic_partition_list_destroy(topics); }这里有一个值得注意的细节RD_KAFKA_PARTITION_UA表示“未分配分区”UA是Unassigned的缩写。调用rd_kafka_subscribe()的时候我们只指定主题名不关心具体的分区分配具体哪个分区由谁负责是Kafka消费组协调器的事情。订阅只负责告诉你“我感兴趣的是哪些主题”实际分区分配在rebalance之后才会落定。4.2 拉取消息与查询当前位置订阅完成之后就是主循环拉消息。拉消息的同时我们经常需要查询当前消费位置。librdkafka提供了一组position相关的接口static void print_curr_position(rd_kafka_t *rk, rd_kafka_topic_partition_list_t *partitions) { rd_kafka_resp_err_t err rd_kafka_position(rk, partitions); if (err ! RD_KAFKA_RESP_ERR_NO_ERROR) { fprintf(stderr, 查询position失败: %s\n, rd_kafka_err2str(err)); return; } for (int i 0; i partitions-cnt; i) { rd_kafka_topic_partition_t *p partitions-elems[i]; printf(topic%s partition%d position% PRId64 \n, p-topic, p-partition, p-offset); } }主循环里poll的返回会带给我们消息和分区信息。我的习惯是每处理完一段消息就查询一次position用来做监控日志输出确认消费者没有“卡住”。4.3 手动提交offset的核心代码手动提交是整套offset操作里最关键的环节。我们来看一段完整的处理循环里面包含了正确的提交时机把握static void consume_loop(rd_kafka_t *rk) { rd_kafka_message_t *rkmessage; while (true) { rkmessage rd_kafka_consumer_poll(rk, 1000); if (!rkmessage) { // 超时没有消息继续等待 continue; } if (rkmessage-err RD_KAFKA_RESP_ERR_NO_ERROR) { // 处理业务数据 process_message(rkmessage); // 手动提交只提交当前这条消息所在的partition的offset1 rd_kafka_topic_partition_list_t *offsets rd_kafka_topic_partition_list_new(1); rd_kafka_topic_partition_t *p rd_kafka_topic_partition_list_add( offsets, rd_kafka_topic_name(rkmessage-rkt), rkmessage-partition); p-offset rkmessage-offset 1; rd_kafka_resp_err_t err rd_kafka_commit(rk, offsets, 0); if (err ! RD_KAFKA_RESP_ERR_NO_ERROR) { fprintf(stderr, 提交offset失败: %s\n, rd_kafka_err2str(err)); // 这里需要根据业务情况决定是否重试 } rd_kafka_topic_partition_list_destroy(offsets); } rd_kafka_message_destroy(rkmessage); } }关于提交的offset值这里必须说明一个关键点rkmessage-offset 1。Kafka的offset语义是下一条要读的消息的序号。如果你刚处理完offset为100的消息下次要从101开始那么提交的值是101不是100。很多新手在这里写错导致每次重启后多重复消费一条。实际上librdkafka在返回消息的时候rd_kafka_message_t里的offset字段就是当前消息的offset。提交时用offset 1是标准写法。如果是用rd_kafka_offsets_store()配合store和commit分离的方式store的时候也是存offset 1。有一个细节我踩过坑rd_kafka_commit()这个函数的第二个参数传NULL时提交的是全局所有分区的当前消费位置。但如果同时有多个线程在消费不同分区全局提交可能会把办线程还没消费完的分区位置也提交出去造成消息丢失。多线程消费的时候必须按分区粒度提交。4.4 seek操作主动调整offset除了被动跟随offset消费Kafka还允许消费者主动设置读取位置。librdkafka提供rd_kafka_seek()接口static void seek_to_offset(rd_kafka_t *rk, const char *topic, int32_t partition, int64_t offset) { rd_kafka_resp_err_t err rd_kafka_seek(rk, topic, partition, offset, 0); if (err ! RD_KAFKA_RESP_ERR_NO_ERROR) { fprintf(stderr, seek失败: %s\n, rd_kafka_err2str(err)); } }seek常用于几个场景想从头重新消费一遍、跳过某段坏数据继续消费、根据业务时间戳定位到某个时间点。常见的用法是先调rd_kafka_offsets_for_times()根据时间戳查到对应的offset再seek过去。要特别注意seek操作是消费者客户端本地生效的它不会同步修改broker上已提交的offset。也就是说你seek改的只是一个临时的读取位置下次重新启动消费者时它还是会从最近一次提交的位置开始除非你在seek之后再次调用rd_kafka_commit()把位置提交上去。实际操作中我一般在做“重新消费”类操作时会先seek然后立刻commit一次确保重启之后也能续上。5. 常见问题与排查技巧实录5.1 重复消费最典型的“追踪尾部消息”场景重复消费是offset相关报障里出现频率最高的。现象一般是这样消费者进程每跑一段时间就重启重启后总是把最后几条消息重新消费一遍。排查思路很清晰按下面几步走先确认是否是自动提交导致的。看配置里enable.auto.commit是不是true。如果是改成false手动提交。确认提交的时机。手动提交的情况要检查代码里是不是在rd_kafka_consumer_poll()返回后立刻commit了但业务逻辑还没处理完。检查提交的offset值是否正确。提交的是rkmessage-offset还是rkmessage-offset 1。用rd_kafka_commit()提交NULL时提交的其实是“最后poll到的消息位置”这个位置可能滞后于已处理完的消息。此时你看到消费者打印的日志里position每次重启都比上次多了几条不是减少了。这是“追尾”的典型症状。如果业务允许最稳妥的兜底方案是在消息处理逻辑里做幂等。幂等不解决重复问题但能消化重复问题带来的后果。数据库插入前先查重、Redis里用setnx、写日志用唯一ID去重这些都是成熟的办法。5.2 消息丢失offset提交太快了和重复消费相反的场景是消息丢失。表现为消费者没有报错业务数据处理也正常但最终发现某些数据没入库。排查方向聚焦在“提交offset先于数据处理完成”这个点。具体来说如果你用的不是单条消息处理完再提交而是一批一批拉取处理完一批再提交这一批的offset那么在这一批处理的过程中如果进程崩溃重启后就会从上一批的已提交位置开始这一批的中间数据会丢。另一个常见原因是用rd_kafka_consumer_poll()拉取消息后如果只把消息指针存起来、异步去处理同时又提交了offset那么处理线程还没跑完offset就已经提交了这时候进程退出任务没做完但是位置已经往前走了。标准解法是改成分区级别的“先处理、后提交”。如果一定需要异步化那就要引入自己的ack机制消息处理成功后才提交对应的offset。5.3 消费停滞position不前进消费者进程还活着日志也正常但就是感觉消费在停滞消息积压持续增多。这种情况要先看是不是position卡住了。我遇到过一次典型的卡住case代码里做了消息批量聚合攒够100条才提交一次offset。某一天消息量减小攒不够100条offset就迟迟不提交。消费者倒是还在拉消息但所有已拉取消息的处理进度被堵在聚合那一步position也不动积压不断增大。这个案例给的教训是提交逻辑不能依赖业务消息的到达频率。现在librdkafka支持定时提交和条件提交混合用哪怕一条消息都没有也应该定期把当前position提交一下避免异常情况下积压。排查position卡住用rd_kafka_position()打印每个分区的当前读取位置对比一下broker里该分区的最新offset用rd_kafka_query_watermark_offsets()可以查到分区的oldest和latest水位就知道消费端到底落后了多少。5.4 offset提交失败权限和版本兼容问题提交失败在C/C环境里通常有三种原因认证权限问题消费者没有写__consumer_offsets主题的权限。新版Kafka对内部主题也做ACL控制如果你的集群启用了SASL或者ACL需要给消费组对应的用户授予Group描述权限和内部主题的写权限。版本兼容问题客户端版本和broker版本差异过大导致提交协议不匹配。一般表现为rd_kafka_commit()返回UNKNOWN_TOPIC_OR_PART或者UNKNOWN_MEMBER_ID之类看起来莫名其妙的错误码。先检查客户端和broker的版本是否在兼容范围内。网络分区问题消费者和broker之间的连接断开提交请求发不出去。这种往往伴随着大量broker连接的transport错误日志。我的排查习惯是先看返回码用rd_kafka_err2str()转成可读字符串再看librdkafka的日志级别把log_level临时调成LOG_DEBUG它会打出每次提交请求的细节。实在不行就抓包看协议交互但一般调日志足够解决90%的问题。6. 一个在实际项目中沉淀的offset提交模式写了这么多最后分享一个我在项目里沉淀的提交模式适用于大多数C/C消费端场景。核心思想是消息处理与offset提交解耦用独立线程按条件触发提交。主消费线程负责poll消息、处理业务逻辑然后把当前每个分区已处理完成的最新位置记录到本地状态表里。提交线程每隔固定时间扫描本地状态表把有进展的分区位置提交到broker。同时在本地状态表里记录“已提交”和“已处理”两个水位。这样做的好处是主线程不被提交的RPC阻塞吞吐量更高。提交线程可以控制节奏批量提交减少broker压力。本地记录了两个水位能精确知道重启后可能重复多少数据便于业务侧做补偿。实现时要注意一点状态表的划分区间要按分区维度维护不能用全局变量存一个“最后处理位置”。Kafka的分区可能在不同消费者实例之间转移本地状态表必须记录的是“这个消费者实例当前持有的分区”的最新处理位置分区被转移走之后新的消费者实例会从提交水位重新拉取。这个模式做下来的直观收益是消费者稳定性明显提升线上很少再出现“重启丢消息”或“不知不觉重跑一大段”的情况。如果各位正在用librdkafka写生产级消费者建议把这套“双水位独立提交线程”的思路落到自己的代码里。过程中最值得花时间的部分不是写poll循环而是设计清楚每个分区什么时候算“处理成功”什么时候可以提交。把这件事想透了Kafka消费者的坑就少了一半。
返回列表