Quarkus对接Kafka报SRMSG00034背压错误的解决方法咨询
问题根因分析
这个SRMSG00034报错的核心是背压不匹配:上游消息生成速度远快于下游Kafka emitter的消息处理速度,emitter内部缓冲区耗尽后没有足够的请求槽位接收新消息。
你当前的实现存在两个核心问题:
invoke(emitter::send)的写法完全脱离了背压管控:invoke是同步操作,不会等待emitter.send()返回的异步Uni结果完成,会无节制地向下游提交发送请求,完全忽略Kafka的实际处理能力- 订阅逻辑不符合Mutiny的背压规则:
subscribe().with()的第一个入参是消费已处理元素的回调,你返回Uni.createFrom().voidItem()没有任何实际作用,也没有参与背压请求的管控
修复方案
第一步:替换invoke为call对齐异步发送速度
call操作符会等待每个异步操作(这里就是Kafka发送请求)完成后,再处理下一个元素,天然将上游的消息生成速度和下游Kafka的发送速度对齐,不需要修改默认的溢出配置即可解决绝大多数场景的问题。
如果需要提升吞吐量,可以搭配merge控制并行发送的并发数,避免并行请求过多压垮emitter。
第二步:调整订阅逻辑
去掉无意义的Uni返回,使用标准的订阅写法即可。
修改后的代码示例
// 基础版:串行发送,完全匹配下游速度,无溢出风险 Multi.createFrom().iterable(msgList) .onItem().transform(item -> { // 原有转换逻辑 }) // 用call等待send异步完成再处理下一条 .call(item -> emitter.send(item)) .subscribe().with( item -> {}, // 可添加发送成功的打点逻辑 Throwable::printStackTrace, () -> System.out.println("Done!") );
// 高吞吐量版:控制并行发送数,兼顾效率和背压 int maxParallelSend = 10; // 可根据Kafka集群性能调整 Multi.createFrom().iterable(msgList) .onItem().transform(item -> { // 原有转换逻辑 }) // 转成发送请求的Uni,控制最大并行数 .onItem().transformToUni(item -> emitter.send(item)) .merge(maxParallelSend) .subscribe().with( item -> {}, Throwable::printStackTrace, () -> System.out.println("Done!") );
极端场景的补充配置
如果你的消息峰值确实远大于Kafka的处理能力,再考虑添加@OnOverflow配置扩展缓冲区:
// 在emitter注入处添加配置,缓冲区大小设置为大于单批次最大消息量即可 @Inject @Channel("your-kafka-channel") @OnOverflow(value = OnOverflow.Strategy.BUFFER, bufferSize = 20000) Emitter<YourMessageType> emitter;
内容的提问来源于stack exchange,提问作者Celso Marques
相关产品推荐
相关产品推荐

