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

Kafka 2.12-2.8.8消息重试致乱序问题求助

Kafka指定分区异步发送重试乱场景复现方案

问题核心原因

当前无法复现的关键问题:

  • 代码仅发送2条消息,无法观察到msg3-msg5先于重试的msg2到达的情况
  • retry.backoff.ms设为30000秒过长,重试触发前后续消息已全部发送完成
  • iptables拦截时机不准确,未精准命中msg2的首次发送请求,导致未触发重试

调整后的复现步骤及代码

1. 关键参数调整

修改生产者配置,缩短重试间隔,确保重试触发前后续消息有足够时间发送成功:

  • 将retry.backoff.ms改为5000
  • 延长delivery.timeout.ms至15000,避免重试前生产者超时退出

2. 补全5条消息发送逻辑

在代码中添加msg3-msg5的发送逻辑,并在回调中打印发送状态,便于观察顺序:

public class MyKafkaProducer1 {

    static {
        ILoggerFactory iLoggerFactory = LoggerFactory.getILoggerFactory();
        LoggerContext loggerContext = (LoggerContext) iLoggerFactory;
        Logger root = loggerContext.getLogger("root");
        root.setLevel(Level.ERROR);
    }

    public static void main(String[] args) throws InterruptedException {
        Properties props = new Properties();
        props.put("bootstrap.servers", "192.168.227.103:9092");
        props.put("acks", "1");
        props.put("delivery.timeout.ms", 15000);
        props.put("retries", 2);
        props.put("retry.backoff.ms", 5000);
        props.put("batch.size", 16384);
        props.put("request.timeout.ms", 5000);
        props.put("linger.ms", 1);
        props.put("buffer.memory", 33554432);
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("max.in.flight.requests.per.connection", 5); // 保持大于1,允许并发请求导致乱序

        Producer<String, String> producer = new KafkaProducer<>(props);
        Thread.sleep(5000);

        // 发送msg1
        producer.send(new ProducerRecord<>("testTopic", 0, "key" + 1, "msg1"), (metadata, exception) -> {
            if (exception != null) {
                System.err.println("msg1发送失败: " + exception.getMessage());
            } else {
                System.out.println("msg1发送成功,offset: " + metadata.offset());
            }
        });

        // 等待msg1发送完成后,执行iptables拦截
        Thread.sleep(2000);
        // 发送msg2(此时已拦截,会触发重试)
        producer.send(new ProducerRecord<>("testTopic", 0, "key" + 2, "msg2"), (metadata, exception) -> {
            if (exception != null) {
                System.err.println("msg2发送失败,触发重试: " + exception.getMessage());
            } else {
                System.out.println("msg2发送成功,offset: " + metadata.offset());
            }
        });

        // 立即放开iptables拦截,让msg3-msg5正常发送
        Thread.sleep(1000);
        // 发送msg3
        producer.send(new ProducerRecord<>("testTopic", 0, "key" + 3, "msg3"), (metadata, exception) -> {
            if (exception != null) {
                System.err.println("msg3发送失败: " + exception.getMessage());
            } else {
                System.out.println("msg3发送成功,offset: " + metadata.offset());
            }
        });

        // 发送msg4
        producer.send(new ProducerRecord<>("testTopic", 0, "key" + 4, "msg4"), (metadata, exception) -> {
            if (exception != null) {
                System.err.println("msg4发送失败: " + exception.getMessage());
            } else {
                System.out.println("msg4发送成功,offset: " + metadata.offset());
            }
        });

        // 发送msg5
        producer.send(new ProducerRecord<>("testTopic", 0, "key" + 5, "msg5"), (metadata, exception) -> {
            if (exception != null) {
                System.err.println("msg5发送失败: " + exception.getMessage());
            } else {
                System.out.println("msg5发送成功,offset: " + metadata.offset());
            }
        });

        producer.flush();
        Thread.sleep(15000); // 等待所有重试完成
        producer.close();
    }
}

3. 精准iptables拦截操作流程

  1. 确保Kafka集群正常运行,testTopic已创建且包含至少1个分区
  2. 启动生产者程序,控制台打印msg1发送成功后,立即执行拦截命令:
    iptables -A OUTPUT -d 192.168.227.103 -p tcp --dport 9092 -j DROP
    
  3. 等待1秒(确保msg2的首次请求被拦截),执行放开命令:
    iptables -D OUTPUT -d 192.168.227.103 -p tcp --dport 9092 -j DROP
    
  4. 观察生产者控制台打印,最终会看到msg3-msg5先于重试的msg2完成发送,消费testTopic即可验证分区内消息顺序为msg1、msg3、msg4、msg5、msg2

原理说明

当max.in.flight.requests.per.connection大于1时,Kafka生产者允许同时发送多个请求到Broker。msg2首次发送失败进入重试队列后,后续的msg3-msg5会正常发送并成功写入分区,msg2重试完成后会追加到分区末尾,从而产生乱序。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 21:37:08