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

Ballerina中kafka消息发送至topic失败报错该如何解决?

Ballerina Kafka 投递超时报错排查方案

报错本质

你遇到的Expiring 1 record(s) for example-0:120010 ms has passed since batch creation是Kafka生产者的标准投递超时错误:消息进入生产者的批量缓冲区后,超过默认120秒仍未成功写入Broker,被主动丢弃触发报错。

可能原因及对应解决方法

  • 1. 生产者配置缺失核心参数

你当前的ProducerConfiguration仅配置了三个基础参数,未适配生产场景的发送逻辑:

  • 默认deliveryTimeoutMs为120000ms(即报错里的120秒),和你设置的3次重试逻辑冲突,重试还没跑完就触发了超时
  • 默认batchSize为16KB,如果单条消息大小超过这个阈值会被阻塞发送
    解决方法:补充配置参数即可:
kafka:ProducerConfiguration producerConfiguration = {
    clientId: "basic-producer",
    acks: "all",
    retryCount: 3,
    // 新增以下配置
    deliveryTimeoutMs: 300000, // 投递超时调整为5分钟,预留充足重试时间
    lingerMs: 5, // 小批量消息聚合等待时长,降低Broker请求压力
    batchSize: 163840 // 批次大小调至160KB,适配更大的单条消息体积
};
  • 2. Kafka Broker 连接异常

你代码中使用的kafka:DEFAULT_URL默认指向localhost:9092,需要确认连接有效性:

  • 检查Kafka Broker的server.properties配置,listeners和advertised.listeners是否对外暴露localhost:9092地址,若配置为其他主机名/IP会导致连接失败
  • 先用Kafka自带命令行工具验证连通性:
    # 测试生产消息
    bin/kafka-console-producer.sh --topic example --bootstrap-server localhost:9092
    # 新开终端测试消费消息
    bin/kafka-console-consumer.sh --topic example --from-beginning --bootstrap-server localhost:9092
    

若命令行也无法正常生产消费,调整Broker监听配置后重启即可。

  • 3. Topic 分区不可用

你创建topic时设置的replication-factor=1,每个分区只有1个副本,若该副本所在Broker宕机,分区会直接进入不可用状态,生产者写入请求会持续重试直到超时:
解决方法:查看topic状态确认分区可用性:

bin/kafka-topics.sh --describe --topic example --bootstrap-server localhost:9092

若输出中Leader列为-1、Isr列表为空,说明分区副本损坏,可重启对应Broker或删除topic后重新创建。

  • 4. 代码初始化逻辑问题

你当前把kafka:Producer的初始化放在了main函数外部,若初始化时Broker未启动、连接失败,不会立即抛出错误,直到send操作时才会触发超时:
解决方法:将生产者初始化逻辑移入main函数内,新增关闭逻辑:

public function main() returns error? {
    kafka:ProducerConfiguration producerConfiguration = {
        clientId: "basic-producer",
        acks: "all",
        retryCount: 3,
        deliveryTimeoutMs: 300000,
        lingerMs: 5,
        batchSize: 163840
    };
    kafka:Producer kafkaProducer = check new (kafka:DEFAULT_URL,producerConfiguration);
    string message = "Hello World, Ballerina";
    check kafkaProducer->send({
                                topic: "example",
                                value: message.toBytes()});
    check kafkaProducer->'flush();
    kafkaProducer->close();
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 14:15:00