写点什么

【kafka 异常】使用 Spring-kafka 遇到的坑

  • 2022 年 9 月 29 日
    江西
  • 本文字数:3495 字

    阅读完需:约 11 分钟

作者石臻臻,CSDN 博客之星 Top5Kafka Contributornacos Contributor华为云 MVP,腾讯云 TVP,滴滴 Kafka 技术专家 KnowStreaming


KnowStreaming 是滴滴开源的Kafka运维管控平台, 有兴趣一起参与参与开发的同学,但是怕自己能力不够的同学,可以联系我,当你导师带你参与开源!

CORRUPT_MESSAGE

这个错误一般是压缩策略为 cleanup.policy=compact 的情况下,key 不能为空

o.a.k.c.p.i.Sender 595 [WARN] [Producer clientId=producer-1] Got error produce response with correlation id 131 on topic-partition SHI_TOPIC1-0, retrying (2147483521 attempts left). Error: CORRUPT_MESSAGE
复制代码

查看一下压缩策略

bin/kafka-topics.sh --describe --zookeeper xxxx:2181 --topic SHI_TOPIC1
Topic:SHI_TOPIC1 PartitionCount:1 ReplicationFactor:1 Configs:cleanup.policy=compact Topic: SHI_TOPIC1 Partition: 0 Leader: 0 Replicas: 0 Isr: 0 
复制代码

Configs:cleanup.policy=compact :

然后再检查一下自己发送消息的时候是不是没有传 key

参考链接

问题堆栈信息

org.springframework.kafka.listener.ListenerExecutionFailedException: invokeHandler Failed; nested exception is java.lang.IllegalStateException: No Acknowledgment available as an argument, the listener container must have a MANUAL AckMode to populate the Acknowledgment.;  nested exception is java.lang.IllegalStateException:  No Acknowledgment available as an argument, the listener container must have a MANUAL AckMode to populate the Acknowledgment.

复制代码

问题原因


解决方案


问题堆栈信息

 Failed to start bean 'org.springframework.kafka.config.internalKafkaListenerEndpointRegistry'; nested exception is java.lang.IllegalStateException: Consumer cannot be configured for auto commit for ackMode MANUAL_IMMEDIATE
复制代码

问题原因

不能再配置中既配置kafka.consumer.enable-auto-commit=true 自动提交; 然后又在监听器中使用手动提交

例如:

kafka.consumer.enable-auto-commit=true
复制代码


    @Autowired    private ConsumerFactory consumerFactory;        @Bean    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaManualAckListenerContainerFactory() {        ConcurrentKafkaListenerContainerFactory<Integer, String> factory =                new ConcurrentKafkaListenerContainerFactory<>();        factory.setConsumerFactory(consumerFactory);        //设置提交偏移量的方式 当Acknowledgment.acknowledge()侦听器调用该方法时,立即提交偏移量        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);        return factory;    }


    /**     * 手动ack  提交记录     * @param data     * @param ack     * @throws InterruptedException     */    @KafkaListener(id = "consumer-id2",topics = "SHI_TOPIC1",concurrency = "1",            clientIdPrefix = "myClientId2",containerFactory = "kafkaManualAckListenerContainerFactory")    public void consumer2(String data, Acknowledgment ack)  {        log.info("consumer-id2-手动ack,提交记录,data:{}",data);        ack.acknowledge();
    }
复制代码

解决方法:

将自动提交关掉,或者去掉手动提交;如果你想他们都同时存在,某些情况自动提交;某些情况手动提交; 那你创建 一个新的consumerFactory 将它的是否自动提交设置为 false;比如


@Configuration@EnableKafkapublic class KafkaConfig {
    @Autowired    private KafkaProperties properties;
    /**     * 创建一个新的消费者工厂     * 创建多个工厂的时候 SpringBoot就不会自动帮忙创建工厂了;所以默认的还是自己创建一下     * @return     */    @Bean    public ConsumerFactory<Object, Object> kafkaConsumerFactory() {        Map<String, Object> map =  properties.buildConsumerProperties();        DefaultKafkaConsumerFactory<Object, Object> factory = new DefaultKafkaConsumerFactory<>( map);        return factory;    }
    /**     * 创建一个新的消费者工厂     * 但是修改为不自动提交     *     * @return     */    @Bean    public ConsumerFactory<Object, Object> kafkaManualConsumerFactory() {        Map<String, Object> map =  properties.buildConsumerProperties();        map.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,false);        DefaultKafkaConsumerFactory<Object, Object> factory = new DefaultKafkaConsumerFactory<>( map);        return factory;    }

    /**     * 手动提交的监听器工厂 (使用的消费组工厂必须 kafka.consumer.enable-auto-commit = false)     * @return     */    @Bean    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaManualAckListenerContainerFactory() {        ConcurrentKafkaListenerContainerFactory<Integer, String> factory =                new ConcurrentKafkaListenerContainerFactory<>();        factory.setConsumerFactory(kafkaManualConsumerFactory());        //设置提交偏移量的方式 当Acknowledgment.acknowledge()侦听器调用该方法时,立即提交偏移量        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);        return factory;    }
    /**     * 监听器工厂 批量消费     * @return     */    @Bean    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> batchFactory() {        ConcurrentKafkaListenerContainerFactory<Integer, String> factory =                new ConcurrentKafkaListenerContainerFactory<>();        factory.setConsumerFactory(kafkaConsumerFactory());        factory.setBatchListener(true);        return factory;    }}
复制代码

消费者监听的时候 指定对应的 容器工厂就行了kafkaManualAckListenerContainerFactory

    /**     * 手动ack  提交记录     * @param data     * @param ack     * @throws InterruptedException     */    @KafkaListener(id = "consumer-id2",topics = "SHI_TOPIC1",concurrency = "1",            clientIdPrefix = "myClientId2",containerFactory = "kafkaManualAckListenerContainerFactory")    public void consumer2(String data, Acknowledgment ack)  {        log.info("consumer-id2-手动ack,提交记录,data:{}",data);        ack.acknowledge();
    }

复制代码



问题堆栈信息

[WARN] Error registering AppInfo mbean javax.management.InstanceAlreadyExistsException: kafka.consumer:type=app-info,id=myClientId-3
复制代码

问题原因

官网描述The client.id property (if set) is appended with -n where n is the consumer instance that corresponds to the concurrency. This is required to provide unique names for MBeans when JMX is enabled.意思是这个 id 在 JMX 中注册需要 id 名唯一;不要重复了;

解决方法:

将监听器的 id 修改掉为唯一值 或者 消费者的全局配置属性中不要知道 client-id ;则系统会自动创建不重复的 client-id

发布于: 刚刚阅读数: 3
用户头像

关注公众号: 石臻臻的杂货铺 获取最新文章 2019.09.06 加入

进高质量滴滴技术交流群,只交流技术不闲聊 加 szzdzhp001 进群 20w字《Kafka运维与实战宝典》PDF下载请关注公众号:石臻臻的杂货铺

评论

发布
暂无评论
【kafka异常】使用Spring-kafka遇到的坑_Kafk_石臻臻的杂货铺_InfoQ写作社区