相关疑难解决方法(0)

如何配置 spring-kafka 以忽略格式错误的消息?

我们的 Kafka 主题之一存在问题,该主题由此处描述的DefaultKafkaConsumerFactory&ConcurrentMessageListenerContainer组合与工厂使用的 a组合使用。不幸的是,有人有点热情,并在该主题上发布了一些无效的消息。似乎 spring-kafka 默默地无法处理这些消息中的第一个。是否可以让 spring-kafka 记录错误并继续?查看记录的错误消息,Apache kafka-clients 库似乎应该处理在迭代一批消息时其中一个或多个消息可能无法解析的情况?JsonDeserializer

以下代码是说明此问题的示例测试用例:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.junit.ClassRule;
import org.junit.Test;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
import org.springframework.kafka.listener.MessageListener;
import org.springframework.kafka.listener.config.ContainerProperties;
import org.springframework.kafka.support.SendResult;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.kafka.test.rule.KafkaEmbedded;
import org.springframework.kafka.test.utils.ContainerTestUtils;
import org.springframework.util.concurrent.ListenableFuture;

import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertThat; …
Run Code Online (Sandbox Code Playgroud)

error-handling apache-kafka spring-kafka

5
推荐指数
1
解决办法
6607
查看次数