ARTICLE DETAIL

资讯详情

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

Kafka延迟队列实战:三种方案对比与生产落地

Kafka延迟队列实战:三种方案对比与生产落地 说实话第一次意识到Kafka做延迟队列这件事有点拧巴是在一个订单超时关闭的业务里。业务方说得很简单订单创建后30分钟未支付就关掉消息用Kafka发就行。结果我去翻Kafka官方文档发现压根没有延迟消息、定时消息这种概念。你可以在Kafka里看到delayed produce、delayed fetch这些内部机制但原生API就是不给普通用户开通到点可见的能力。后来我把Kafka源码里那套处理延迟任务的时间轮机制研究了一遍又在生产环境里实际落地过好几种方案才彻底想明白Kafka做延迟队列不是不行关键是得在外面自己搭一套机制把延迟这个概念补进去。这篇文章就把我实践过的方案完完整整拆开讲从最低成本的时间戳轮询到复制Kafka内部的TimingWheel时间轮设计再到生产环境最稳的多级Topic转投方案最后附上我踩过的坑和延迟精度的实测思路。适合已经有Kafka基础、想在自己项目里实现可靠延迟队列的读者。1. 先搞清楚Kafka为什么没有延迟消息这种一等公民能力1.1 从Kafka的设计哲学看延迟队列的本质冲突Kafka从设计之初就不打算做定时送达这件事。它的核心定位是高吞吐、持久化、可回溯的分布式日志消息一旦被写入所有消费者组都可以根据自己记录的offset自由消费Kafka完全不关心这条消息应该在什么时间点被看见。延迟队列的本质需求恰恰相反它要求的是到时间才可见。这意味着系统需要给每条消息附加一个可见性时间戳并且要有一个调度机制在某个精确时刻把消息从不可见状态切换为可见状态。这两者的冲突点非常本质Kafka保证的是消息不会丢、分区有序、写进去就能被拉取延迟队列要求的是写进去先藏着到点才能被拉取。用快递柜来做类比可能更直观——Kafka是一个24小时营业的快递柜它负责保管和自取但晚上八点整才允许你取件这个规则快递柜本身不负责需要额外的机制去实现。1.2 消费者拉取模型决定了到点可见不是默认能力再往底层看KafkaConsumer的工作模型是持续不断的poll轮询。消费者把消息从broker拉回来之后这批消息就从这个消费者的视角里出现了处理完之后提交offset代表消费完成。这个模型里有个很微妙的点消息一旦被poll出来Kafka集群层面就认为它已经被订阅了如果你不处理又不提交offset这批消息的唯一去处就是你自己的本地内存或者磁盘。所谓的延迟就变成了一件纯粹由客户端自己负责的事——把消息藏在自己手里到点再处理。这就带来了延迟队列的基础问题延迟逻辑落在消费者客户端而不在broker端。所以要实现延迟队列本质上只有两条路要么在客户端引入一种延迟触发机制让消息在本地或单独的服务里乖乖等待要么在消费者拉取那一刻做一个判断过滤不满足时间条件的消息先不处理等下一轮poll再判断。1.3 Kafka内部其实一直在处理延迟Purgatory是现成的参考这里有个特别有意思的反差Kafka对外没有提供延迟队列但它在系统内部处理延迟请求的机制成熟得惊人。这套机制叫DelayedOperationPurgatory中文可以理解成延迟操作炼狱专门管理那些条件暂时不满足、需要等一会儿再处理的请求。举几个常见的内部场景。Producer发送消息时如果acksall就要等所有副本都写入成功那些副本还没跟上来的写入请求会被挂进Purgatory延迟一段时间再检查。Consumer发来fetch请求时如果当前分区没有新消息这个fetch请求也会被挂起等新消息到了再唤醒。还有消费组协调器在等待成员加入时也需要延迟判断是否超时。管理这些延迟任务的核心数据结构就是Kafka源码里的TimingWheel时间轮。它是从Netty的HashedWheelTimer借鉴演化来的任务插入、删除、到期触发都极其高效。官方没有把这个能力开放给用户但这是一个绝佳的参考实现。想用Kafka做可靠的延迟队列借鉴这一套时间轮的设计比自己在客户端里写一堆乱糟糟的判断逻辑要靠谱得多。2. 方案A时间戳过滤轮询——成本最低但要小心三个坑2.1 思路与消息结构设计方案A是所有方案里实现成本最低的思路一句话就能说清生产者发送消息的时候在消息头里带上一个期望投递时间字段消费者每次poll到消息后先判断时间到没到没到就等下一轮再判断。具体到消息结构最简单的做法是利用Kafka消息的headers属性。比如定义两个header_deliverAt代表期望投递的时间戳毫秒_delayMs代表延迟时长毫秒。实际使用中我更推荐只存_deliverAt因为消费者判断的时候直接拿当前时间和它做比较就可以不用再额外算一次。生产者在构建消息的时候注意一点_deliverAt必须用字符串序列化因为Kafka headers的value是字节数组。写入时指定为System.currentTimeMillis() delayMs即可。2.2 可直接运行的消费者骨架下面这段代码就是方案A的完整消费端核心逻辑我基于Kafka 3.x的API写了一个骨架基本可以直接拿去用。// 假设有一个本地延迟队列用于暂存未到期的消息 PriorityBlockingQueueDelayedRecord delayedQueue new PriorityBlockingQueue(); while (true) { // 第一步先把本地队列里到期的消息处理掉 while (delayedQueue.peek() ! null System.currentTimeMillis() delayedQueue.peek().getDeliverAt()) { DelayedRecord record delayedQueue.poll(); process(record); // 真正的业务处理逻辑 } // 第二步拉取Kafka里的新消息 ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(200)); for (ConsumerRecordString, String record : records) { long deliverAt Long.parseLong( new String(record.headers().lastHeader(_deliverAt).value())); if (System.currentTimeMillis() deliverAt) { process(record); // 已到期的消息直接处理 } else { delayedQueue.add(new DelayedRecord(record, deliverAt)); } } }这段代码从一个新的角度回答了Kafka生产消费命令启动一次会一直运行吗这个问题。消费者必须保持while(true)持续运行因为poll是一个被动的拉取动作一旦进程退出本地延迟队列里那些未到期的消息就全没了Kafka却已经认为这些消息被拉取过了。所以方案A对消费者进程的稳定性要求极高直接决定了消息不丢的底线。初看这个方案很简单但它最大的问题是隐藏得很深。netstat看一眼就明白了消费者进程其实是把Kafka当成了一个消息源真正holding住延迟消息的是消费者自己的内存PriorityBlockingQueue这就是方案A最容易踩坑的地方。2.3 这个方案为什么只适合小规模OOM、空轮询、精度漂移我在生产环境真正用过方案A处理一个内部通知场景消息量大概每分钟几千条延迟时长1到5分钟结果下面几个坑全部踩了一遍。第一个坑是内存爆炸。假设你的消息有延迟10分钟的档位前10分钟内的所有消息都会堆在本地PriorityBlockingQueue里如果消息量大而且每条消息的body都比较重JVM堆内存飙升到GC告警甚至OOM是很快的事。Kafka是落盘存储的消息放在broker上几乎不占客户端内存但如果用方案A等于是把Kafka的存储能力废掉了把压力全部转嫁到消费者JVM。这也是为什么说方案A天然不适用于大批量延迟场景。第二个坑是重复消费和offset管理的复杂性。延迟期间消息已经进了本地队列但offset还没提交一旦消费者宕机或者触发rebalance重启后会从上次提交的offset重新拉取这批消息本地队列里那些还没处理的消息再次出现处理逻辑如果没有幂等设计业务数据就乱了。你说先提交offset吧本地队列消费到一半宕机那些消息就永远丢了。这是一个无解的取舍。第三个坑是延迟精度完全看poll间隔的脸色。方案A里poll(200)意味着最多200毫秒扫描一次本地队列看似精度还行但实际情况往往是业务代码处理耗时不稳定导致poll间隔被拉长5秒档的消息可能实际延迟到6秒甚至10秒。热搜词里那个kafka消息延迟高的问题在方案A里非常典型——不是Kafka本身延迟高而是消费者的本地延迟队列拖累了整体节奏。所以我的结论很直接方案A够简单但只能用在延迟消息量小、延迟精度要求不高、能接受消息偶尔丢失和重复的原型验证或者内部小工具场景。真要上核心业务它撑不住。3. 方案B时间轮算法——借鉴Kafka Purgatory的底层解法3.1 先看Kafka源码里TimingWheel怎么组织时间方案A是消息堆积在消费者本地内存方案B则是换一种思路用一个专门的时间轮数据结构来管理延迟任务让消息本身不被消费者长期占用。要理解时间轮最直接的办法就是翻开Kafka源码看TimingWheel的实现。它的核心结构由这么几个要素组成tickMs一个槽位代表的时间跨度Kafka集群中默认是1毫秒wheelSize一整圈的槽位数量源码里默认是20startMs时间轮创建时的时间戳buckets一个数组每个元素是一个TimerTaskList双向链表存放所有到期时间落在同一个槽位里的延迟任务overflowWheel当任务延迟时间超过当前时间轮的覆盖范围时就提升到更高一层的时间轮形成层级结构第一层时间轮覆盖0到20毫秒第二层每个槽位代表20毫秒并覆盖0到400毫秒第三层每个槽位代表400毫秒并覆盖0到8000毫秒以此类推。层级结构让时间轮既能处理毫秒级任务又能无压力的处理分钟级甚至小时级任务。3.2 任务插入与bucket降级机制向时间轮里插入一个延迟任务时会先计算这个任务距离当前时间有多少个tick从而定位到对应的槽位。如果延迟范围超出了当前层的覆盖范围就把任务转交给overflowWheel递归插入到合适的高层槽位。每个槽位存储的不是单个任务而是一个TimerTaskList双向链表。这样做有一个很重要的原因如果每个任务单独创建一个定时器大量任务的插入和取消会非常昂贵。但用时间轮链表结构添加、删除任务都是O(1)的操作到期时整个bucket被取出来交给后台线程逐条执行回调即可。这里有个关键的技术概念叫bucket降级。高层时间轮在tick推进过程中会逐渐把任务流转到低层时间轮的bucket里。比如一个延迟1分钟的任务最开始被放在覆盖小时级范围的高层槽位里随着时间推进当距离到期只有几百毫秒时它会被重新映射到低层时间轮的合适槽位。这个过程是Kafka时间轮高性能的核心它保证每个任务的到期检查都在最细粒度的时间刻度上进行。3.3 用现成的HashedWheelTimer把时间轮接入Kafka生产消费链路你当然可以自己照着Kafka源码撸一个时间轮但在Java生态里Netty的HashedWheelTimer就是现成的高性能时间轮实现Kafka早期版本也参考过它。直接把HashedWheelTimer集成到Kafka的消费链路里写法非常简单。// 创建一个tick为10ms、槽位512个的HashedWheelTimer Timer timer new HashedWheelTimer( r - new Thread(r, delay-timer), 10, TimeUnit.MILLISECONDS, 512 ); // 消费者poll到消息后不直接处理而是注册到时间轮 timer.newTimeout(timeout - { // 到期后把消息转发到真正的业务处理流程 kafkaProducer.send(new ProducerRecord(realBizTopic, messageBody)); }, delayMs, TimeUnit.MILLISECONDS);对比方案A方案B的进步非常明显。消费者poll到消息时马上把消息转存到时间轮里自己继续poll下一批没有长期占用内存的积压。到期时间到了回调线程自动把消息转发到业务处理链路。消费者进程的poll节奏不会被延迟任务拖累。但时间轮方案有个致命的特性需要清醒认识它是一个JVM内存里的状态容器。假如服务进程突然崩溃时间轮中所有未触发的延迟任务会全部丢失。所以生产环境用时间轮必须配套设计补偿机制。我实践过的做法是原始消息仍然先落盘到Kafka时间轮里的任务只作为内存态的延迟触发指针服务重启后通过重新消费Kafka分区并对比消息里的到期时间把还没到期的任务重新注册进新启动的时间轮。用Kafka的持久化替时间轮兜底。3.4 时间轮的边界条件与注意事项时间轮在实际落地时还有几个容易被忽视的细节。第一个是层级覆盖范围的问题。即使时间轮支持overflowWheel层级扩展你也不能把一个延迟一年的任务直接塞进默认配置的时间轮除非预先估算好最大延迟范围合理配置tickMs和wheelSize。比如一个tick为10ms、512个槽位的时间轮覆盖最大延迟就是5.12秒不够的话会自动扩展层级但扩展层级过多会影响性能曲线所以预先规划比事后补救靠谱。第二个是多实例部署的一致性。时间轮是单机内存态的分布式环境下需要保证同一个业务键比如同一个订单ID的延迟任务始终被同一台实例处理否则判断和调度就会乱。好在Kafka消费组机制天然保证了分区分配的一致性只要路由键用的是同一个就能把同一条消息固定到同一个消费者实例上。用时间轮方案我还是那句话对延迟精度要求高、可以接受额外补偿逻辑、团队有能力写代码处理崩溃恢复的场景它是很好的选择。但要论生产环境的稳妥程度多级Topic转投方案更让人放心。4. 方案C多级Topic 定时消息转投——生产环境最稳妥的可落地方案4.1 延迟级别Topic设计不是随便建几个Topic就行方案C的思路完全绕开了客户端持有延迟消息这个坑改用一个独立的消息流转链路生产者先把消息发到一个延迟Topic由专门的转投服务消费延迟Topic到时间后再把消息转发到真正的业务Topic业务消费者只盯着自己的业务Topic消费。延迟Topic的规划设计是这个方案的基石。我比较推荐按固定时间档位建立一组Topic比如app-delay-5sapp-delay-30sapp-delay-5mapp-delay-1happ-delay-1d为什么不直接对每个延迟值建一个Topic因为Topic数量是集群层面的资源每个Topic都会带来额外的元数据管理和网络开销建几百个Topic会让broker和运维都很难受。固定档位的代价是最多牺牲一点延迟精度换来的是Topic数量和运维复杂度完全可控。生产者侧的逻辑也随之简化。延迟5秒的消息route到app-delay-5s延迟30秒的消息route到app-delay-30s以此类推。消息的headers里带上原始业务Topic名称body带上真正的业务消息内容。这里的关键设计是延迟Topic的消息体本身不执行业务逻辑它就是一个待转投任务。4.2 转投服务的poll pause seek实现转投服务是这个方案的大脑它承担一个非常核心的职责消费延迟Topic里那些还没到期的消息不处理也不丢到点了再投递出去。放在Kafka消费模型里这里的难点是消息还没到期时消费者不能一直卡在同一批消息上但也不能简单地把这批消息当垃圾扔掉。对于没到期的消息正确姿势是不提交offset把对应的分区暂停pause然后把消费者位置重置seek回那条未到期消息的offset等时间到了再恢复resume分区。这样消费者不会一直空转轮询也不会把未到期消息积压在本地内存里。粗暴地说就是把等待的压力通过与Kafka分区的交互操作返还给了Kafka侧的存储。下面给出一个转投服务的核心骨架代码。while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(200)); // 记录每条分区里最早未到期消息的投递时间 MapTopicPartition, Long earliestDeliverAt new HashMap(); for (ConsumerRecordString, String record : records) { long deliverAt Long.parseLong( new String(record.headers().lastHeader(_deliverAt).value())); if (System.currentTimeMillis() deliverAt) { // 到期了转发到业务Topic producer.send(new ProducerRecord( record.headers().lastHeader(_bizTopic).toString(), record.key(), record.value() )).get(3, TimeUnit.SECONDS); // 这条消息处理完成提交offset consumer.commitSync(Collections.singletonMap( new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() 1) )); } else { // 未到期记录分区最早到期时间后续统一pauseseek TopicPartition tp new TopicPartition(record.topic(), record.partition()); earliestDeliverAt.merge(tp, deliverAt, Math::min); } } // 统一暂停所有还有未到期消息的分区并恢复到未到期位置 for (Map.EntryTopicPartition, Long entry : earliestDeliverAt.entrySet()) { TopicPartition tp entry.getKey(); consumer.pause(Collections.singleton(tp)); // 从最早未到期的那条offset重新开始 if (latestOffsets.containsKey(tp)) { consumer.seek(tp, latestOffsets.get(tp)); } // 到点后resume分区 delayResume(tp, entry.getValue()); } }这段代码里有一个非常重要的细节enable.auto.commit必须设为false手工管理offset。否则消费者在poll到未到期消息时自动提交offset会把它们标记为已消费消息在到点前就丢了。这个方案里最关键的问题就是谁来在到达到期时间之后恢复分区我的实现是用一个后台的ScheduledThreadPoolExecutor扫描各分区对应的到期时间到点后调用consumer.resume(tp)。每一轮poll里可能同时存在多个不同到期时间的未到期分区但到期恢复的粒度最细也就做到分区级别这是Kafka消费者模型的一个天然边界因为一个分区内消息按offset有序不方便对单条消息单独做pause。4.3 重复消费、消息丢失和幂等保护多级Topic方案里消息从生产到业务消费至少经过两次Kafka投递第一次生产者写入延迟Topic第二次转投服务把消息转发到业务Topic。Kafka的语义是at-least-once两个环节叠加重复消费的风险几乎翻倍。这一点必须直面不能指望运气。我在这套方案里做的兜底措施有三层。第一层是生产端幂等。Kafka生产者在配置enable.idempotencetrue之后同一会话内对同一分区的消息写入不会产生重复虽然它不能覆盖跨会话级别的重复但已经能挡掉最典型的网络重试导致的重复。Kafka 3.0之后这个参数默认是开启的如果你在更早版本上记得手动打开。第二层是转投过程的重复控制。转投服务成功发送到业务Topic之后如果同步提交offset前进程崩溃重启后会重新消费同一条消息并再次转发到业务topic导致业务Topic出现一模一样的两条消息。这个只能从业务侧兜底消息里必须携带唯一的messageId比如订单号或UUID业务消费者拿到之后用Redis的SETNX或者数据库唯一键做去重。第三层是延迟漂移告警。监听转投服务实际投递时间和期望投递时间的差值如果某条消息晚投太久说明转投链路有积压或调度异常直接上报告警而不是默默放行避免延迟任务累积演变成雪崩。这三个层次配合起来多级Topic方案才算是把消息不丢、尽量不重、异常可见这几个核心诉求全部覆盖住了。5. 三个方案怎么选精度、成本、运维维度全对比5.1 一张表看懂差异到这儿三个方案都讲完了说实话光看单点细节很容易犯迷糊先放一张我平时给团队选的对比表对比维度方案A 时间戳轮询方案B 时间轮方案C 多级Topic转投延迟精度依赖poll间隔秒级且不稳定受tick影响可到毫秒级依赖调度精度秒级可保证内存占用高消息积压在消费者本地中时间轮持有未到期任务低未到期消息留在Kafka磁盘消息丢失风险高宕机时本地缓存消失且offset杂乱中需配套补偿机制低持久化Kafka兜底正确处理offset则不丢实现复杂度低中高要处理恢复逻辑中高需要额外转投服务和Topic规划吞吐能力低受单消费者内存限制中高高可以灵活横向扩展消费者实例运维友好度低问题难排查中需要实时盯着JVM内存高链路清晰Lag指标肉眼可见适用场景原型验证、小规模内部工具对精度要求高且接受补偿核心业务、长期运营、不丢不重从表里可以看到三个方案并没有哪一个全面胜出。方案A赢在简单方案B赢在精度方案C赢在稳妥和可运维性。5.2 从业务场景反推选型聊选型的时候我一直跟团队说不要光盯着延迟队列四个字要先问三个问题这个延迟任务的量级是多少允许丢失或重复吗能不能接受引入额外的开发运维成本如果是订单超时关闭、优惠券过期这类交易核心场景我的建议是直接方案C。延迟精度几百毫秒或一两秒完全可以接受但消息丢失是零容忍的而且这类场景需要能随时查这条订单消息现在积压在哪、什么时候会触发链路越清晰越好多级Topic方案天然满足。如果是RPC调用重试、缓存延迟刷新这类服务内部任务方案B更合适。延迟任务的数据密集且生命周期短时间轮回调的效率远高于Kafka消息转投而且这些任务失败了对业务影响有限配合数据库表做补偿就能兜住。如果只是搭个Demo演示Kafka的消费能力方案A十分钟就能跑通没必要为一个演示去建转投服务。另外多说一句如果团队愿意引入新的中间件RocketMQ和云上提供的定时消息产品也都很成熟但这意味着你的技术栈要从Kafka分裂出一套新的消息体系对运维和人力都是额外负担。既然场景已经全面围绕Kafka展开方案C是最平滑的落地路径。5.3 关于集群参数和可视化监控的补充配置不管你选哪个方案延迟队列对Kafka集群本身的可靠性要求都要高于普通业务。我的建议是延迟Topic统一设置acksall配合min.insync.replicas2避免写入副本未同步就返回成功。查看延迟Topic里的数据可以直接用命令行kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic app-delay-5s --from-beginning生产环境我更推荐直接上可视化工具Kafka UI和Kafdrop都是不错的选择能看到每个延迟Topic的分区数、消息总量和消费组Lag。对延迟队列来说Lag指标是所有监控项里最需要盯的。正常的链路里业务Topic的Lag应该长时间维持在0延迟Topic的Lag代表未到期的消息积压量这是合理的但如果你发现某个延迟Topic的Lag在持续快速上涨且转投服务的消费速率跟不上说明延迟任务积压了得赶紧扩容消费者或者排查处理性能。6. 生产环境踩坑记录延迟精度、消费堆积和集群参数6.1 消费组rebalance为什么会打乱延迟顺序方案C落地初期我遇到过延迟消息乱序的问题。最终定位到根因是转投服务在处理一批消息时耗时超过了max.poll.interval.ms的默认值5分钟触发了消费组rebalance。rebalance期间消费者实例会被暂停分区会被重新分配原本按照offset有序排列的延迟消息被分到不同实例处理瞬时产生了乱序。这个问题有两个层面的解法。第一个是参数层面转投服务独立使用一个消费组把这个组的max.poll.interval.ms调整到10分钟以上同时把heartbeat.interval.ms保持在3秒以内确保心跳一直在线。第二个是设计层面如果转投逻辑里有一条消息处理时间不可控那就不要把复杂逻辑放在poll线程里poll只负责拉消息和提交offset真正的转发动作丢给后面的线程池异步处理用异步屏障保证最终一致性。方案A里也有类似的rebalance隐患。消费者本地队列里存了未到期消息rebalance时该分区被分给其他实例新实例会重新消费这批消息但老实例的内存队列里还留着旧的消息。两边同时处理不重复才怪。所以方案A对rebalance几乎零容忍消费者进程不能随意重启。6.2 延迟精度偏差的根源poll间隔、fetch参数和时钟很多人以为延迟队列的精度只取决于调度触发的时间点实际上延迟精度的偏差是层层叠加的。我实测过一组数据延迟5秒档的消息配置poll(100)和fetch.max.wait.ms100时P99实际延迟在5.2秒左右但把fetch.max.wait.ms改到500之后P99延迟直接跳到5.8秒。原因很简单消费者在长时间没有新消息时会阻塞在拉取请求上这个阻塞时间本身就加进了延迟链路。所以做延迟队列时转投服务的poll参数和fetch.max.wait.ms别设太大。如果你的延迟精度要求是秒级fetch.max.wait.ms控制在100毫秒以内比较稳妥。同时转投服务里所有定时调度器统一用ScheduledThreadPoolExecutor不要用Timer——Timer的调度完全基于系统时钟参考和单线程执行一旦有任务执行超时后续所有任务都会整体漂移这在延迟场景里是很要命的。还有一个容易被忽略的是多实例机器之间的时钟同步。方案B的时间轮和方案C的到期判断都依赖本机系统时钟如果两台消费者的机器时钟差了好几秒同样的延迟任务可能在两台实例上判断出不同的到期时间。生产环境给所有机器配上NTP同步这是最基础但最容易被忽视的步骤。6.3 压测与验证如何证明你的延迟队列是合格的延迟队列上线前我一直坚持做一轮完整的压测验证不能只看功能跑通就完事。这里给出一套我实际用过的验证方案。首先构造一批带唯一业务ID的延迟消息设置不同的延迟档位5s、30s、1h等生产者发送时记录本地时间戳T0。业务侧消费者收到转投后的消息时记录接收时间T1那么这条消息的实际延迟DELTA T1 - T0。统计所有消息的P50、P95、P99延迟值同时统计最大偏离度——比如5秒档消息实际最大延迟到了6.2秒偏离1.2秒这个数据如果超出业务容忍范围就要回头查poll间隔和调度线程是否合理。然后验证不丢不重。发送的总消息数和业务侧实际接收并处理成功的消息数做比对必须相等。每条消息带唯一的messageId业务侧处理时把messageId写入Redis或数据库唯一索引跑完压测统计有没有重复插入冲突。最后是宕机恢复演练。在压测中直接kill掉转投服务进程等30秒再启动观察三件事延迟Topic里未提交offset的消息是否被重新消费业务Topic里是否出现了重复消息延迟消息的整体延迟是否因为宕机而发生了大规模漂移。这轮演练的意义在于把最坏情况的处理流程提前走通而不是等到线上真出问题再临时救火。我第一次做方案C压测时发现5秒档的P99延迟达到了8秒排查了一圈才发现是转投服务的消费者线程里顺便做了一次数据库批量插入把poll线程卡了将近3秒。把数据库操作挪到异步线程池之后P99延迟才回到5.5秒以内。这种隐蔽的性能问题只有真实压测才暴露得出来。说实话延迟队列这个需求不管业务方怎么描述落到Kafka技术上始终绕不开客户端持有时间状态这件事。我最后在团队里落地的是方案C加外部Redis去重原因也很朴素链路上每一步都看得见、算得清每条消息都有Kafka的offset作为追溯依据出了问题能查能回溯对我来说这就是最实用的一套方案。时间轮更多是用在了JVM内部的重试场景和短时延任务里。方案没有绝对的高下之分结合你的消息量级、精度要求、团队运维能力选定一套并吃透它的边界就不会失控。
返回列表