前言:之前博客里面提到本公司为物联网项目。项目中使用mqtt+kafka进行与设备端的通讯,之前的协议格式为json格式,现在改成字节数组byte[]格式进行通信。集成springboot后,具体的demo网上很多,接下来有时间会出一份kafka的demo。

  • 报错信息如下:
Failed to start bean 'org.springframework.kafka.config.internalKafkaListenerEndpointRegistry'; 
nested exception is org.apache.kafka.common.KafkaException:Failed to construct kafka consumer

  • 原因分析:
    之前json格式通信时候,构建kafka消费工厂的时候,其中ConcurrentMessageListenerContainer的key为String类型,而value现在为byte[]类型,所以构建消费者工厂的时候需要指定正确的value类型。

  • 代码如下:

	public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, byte[]>> kafkaListenerContainerByteFactory() {
    	ConcurrentKafkaListenerContainerFactory<String, byte[]> factory = new ConcurrentKafkaListenerContainerFactory<String, byte[]>();
    	factory.setConsumerFactory(consumerByteFactory());
    	factory.setConcurrency(concurrency);
    	factory.getContainerProperties().setPollTimeout(1500);
    	return factory;
    }

整体kafka生产者+kafka消费者的demo会在接下来的博客中陆续整理。

Logo

Kafka开源项目指南提供详尽教程,助开发者掌握其架构、配置和使用,实现高效数据流管理和实时处理。它高性能、可扩展,适合日志收集和实时数据处理,通过持久化保障数据安全,是企业大数据生态系统的核心。

更多推荐