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

Spring Kafka Consumer随机跳过Offset问题排查求助

问题:Spring Kafka Consumer随机跳过Offset

我使用Spring Kafka实现了一个Kafka Consumer,从test主题读取消息以处理业务逻辑,但发现监听器会随机跳过部分Offset。该主题包含20个分区,我的应用部署在4个节点上,每个节点配置5个线程来分配分区,消费组名称唯一且无重复。

Consumer代码

@KafkaListener(topics = "test", groupId = "testgroup", 
            concurrency = "5", containerFactory = "kafkaListenerContainerFactory", autoStartup = "true")
    public void handleMessage(@Payload String message,
                              @Header(value = KafkaHeaders.RECEIVED_MESSAGE_KEY, required = false) String messageKey,
                              @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition,
                              @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                              @Header(KafkaHeaders.OFFSET) String offset) {

Consumer配置

Map<String, Object> getConsumerProperties() {
            Map<String, Object> config = new HashMap<>();
            config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "bootstap server coordinates");
            config.put(ConsumerConfig.GROUP_ID_CONFIG, "consumerId");
            config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
            config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
            config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
            config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,"org.apache.kafka.common.serialization.StringDeserializer");
            config.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.RoundRobinAssignor");
            config.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "10000");
            config.put(ConsumerConfig.METADATA_MAX_AGE_CONFIG,30000);
            return config;
        }

已尝试的排查动作

  • 曾尝试将Consumer设置为read_uncommitted,怀疑是事务型Producer未提交某些Offset,但发现默认值即为read_uncommitted,无需显式配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 04:16:26