【发布时间】:2021-06-15 01:07:22
【问题描述】:
这是我的配置代码。
import java.util.function.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.springframework.boot.autoconfigure.kafka.ConcurrentKafkaListenerContainerFactoryConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler;
import org.springframework.web.client.RestTemplate;
@Configuration
public class StreamConfiguration {
@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
ConcurrentKafkaListenerContainerFactoryConfigurer factoryConfigure,
ConsumerFactory<Object, Object> kafkaConsumerFactory) {
ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factoryConfigure.configure(factory, kafkaConsumerFactory);
factory.setBatchListener(true);
factory.setBatchErrorHandler(new SeekToCurrentBatchErrorHandler() {
@Override
public void handle(Exception thrownException, ConsumerRecords<?, ?> data, Consumer<?, ?> consumer,
MessageListenerContainer container) {
Config.this.ehException = thrownException;
super.handle(thrownException, data, consumer, container);
}
});
return factory;
}
@Bean
public RestTemplate restTemplate() {
return new RestTemplate();
}
}
这是我的消费者代码
@KafkaListener(id = "#{'${spring.kafka.listener.id}'}", topics = "#{'${spring.kafka.consumer.topic}'}")
public void getTopics(@RequestBody List<Request> model) {
streamProcessor.runParallel(model.parallelStream());
}
我在处理异常时遇到错误,指出 Consumer 类型的参数数量不正确;它不能用 arguments 参数化。 而且我对导入要导入的配置感到困惑(Config.this.ehException = throwException;),因为有两个选项 apache common 和 apache client admin。
请帮助我无法设置批处理错误处理程序,并且由于反序列化错误而处于无限循环中:(((。 我正在使用 Java8
【问题讨论】: