【发布时间】:2019-07-02 04:52:30
【问题描述】:
使用 Spring Boot,我正在尝试将我的 Kafka 消费者设置为批量接收模式:
@Bean
public ConcurrentKafkaListenerContainerFactory<String, GenericData.Record> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, GenericData.Record> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setMessageConverter(new StringJsonMessageConverter()); // I know this one won't work
factory.setBatchListener(true);
return factory;
}
@Bean
public ConsumerFactory<GenericData.Record, GenericData.Record> consumerFactory() {
Map<String, Object> dataRiverProps = getDataRiverProps();
dataRiverProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("bootstrap.servers"));
return new DefaultKafkaConsumerFactory<>(dataRiverProps);
}
这就是实际消费者的样子:
@KafkaListener(topics = "#{'${kafka.topics}'.split(',')}", containerFactory = 'kafkaListenerContainerFactory')
public void consumeAvro(List<GenericData.Record> list, Acknowledgment ack) {
messageProcessor.addMessageBatchToExecutor(list);
while (messageProcessor.getTaskSize() > EXECUTOR_TASK_COUNT_THRESHOLD) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
LOGGER_ERROR.error(ExceptionUtils.getStackTrace(e.getCause()));
}
}
}
我得到的异常如下所示:
nested exception is org.springframework.core.convert.ConverterNotFoundException: No converter found capable of converting from type [org.apache.avro.generic.GenericData$Record] to type [org.springframework.kafka.support.Acknowledgment]
at org.springframework.core.convert.support.ConversionUtils.invokeConverter(ConversionUtils.java:46)
at org.springframework.core.convert.support.GenericConversionService.convert(GenericConversionService.java:191)
at org.springframework.core.convert.support.GenericConversionService.convert(GenericConversionService.java:174)
at org.springframework.messaging.converter.GenericMessageConverter.fromMessage(GenericMessageConverter.java:66)
Kafka 消息是 AVRO 消息,我想将它们检索为 JSON 字符串。是否有可以插入 ConcurrentKafkaListenerContainerFactory 的 GenericData.Record 即用型 AVRO 转换器?谢谢!
【问题讨论】:
-
这就是我发现的。如果我从签名中删除确认,那么它会正常工作: public void consumeAvro(List
> list) { ...} 我希望能够手动提交,正如你从我的原始帖子我试图扼杀消费者。那么如果我在while循环后不通过调用ack.acknowledge()手动确认,是否意味着它会在while循环后自动确认,以及根据auto.commit.interval.ms设置定期确认?谢谢! -
有人能指点我一些例子来说明我可以为 GenericData.Record 实现消息转换器吗?谢谢!
-
你想转换成GenericData.Record POJO吗?
-
谢谢@Prabhakar!我真的不知道从哪里开始。如果有人甚至可以向我展示一个非常有帮助的示例!
-
我一定会创建一个例子。
标签: json spring-boot kafka-consumer-api avro