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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 02:06:03