Liferay集成Kafka遇消息接收失败问题寻求解决方案
Liferay集成Kafka问题排查与实现方案
需求与问题描述
我在网络上未找到Liferay集成Kafka的参考资料,需要实现两个核心需求:
- 向Kafka Topic推送消息
- 从Kafka Topic拉取并接收消息
我尝试了以下配置与代码,但从Kafka终端推送消息后无法接收消息。
已尝试的依赖配置
compileInclude "org.springframework.kafka:spring-kafka:2.9.2" compileInclude "org.apache.kafka:kafka-streams:3.3.1"
已尝试的代码实现
Kafka配置类
import org.apache.kafka.common.serialization.Serdes; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.annotation.EnableKafkaStreams; import org.springframework.kafka.annotation.KafkaStreamsDefaultConfiguration; import org.springframework.kafka.config.KafkaStreamsConfiguration; import java.util.HashMap; import java.util.Map; import static org.apache.kafka.streams.StreamsConfig.APPLICATION_ID_CONFIG; import static org.apache.kafka.streams.StreamsConfig.BOOTSTRAP_SERVERS_CONFIG; import static org.apache.kafka.streams.StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG; import static org.apache.kafka.streams.StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG; @Configuration @EnableKafka @EnableKafkaStreams public class KafkaConfig { @Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) KafkaStreamsConfiguration kStreamsConfig() { Map<String, Object> props = new HashMap<>(); props.put(APPLICATION_ID_CONFIG, "streams-app"); props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.put(DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); return new KafkaStreamsConfiguration(props); } }
Kafka消息接收器
import com.liferay.portal.kernel.log.Log; import com.liferay.portal.kernel.log.LogFactoryUtil; import org.osgi.service.component.annotations.Component; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.scheduling.annotation.Async; import org.springframework.transaction.annotation.Transactional; @Component(immediate = true) public class KafkaMessageReceiver { public static Log _log = LogFactoryUtil.getLog(KafkaMessageReceiver.class); @Async @KafkaListener( topics = "liferay-topic", concurrency = "2" ) @Transactional public void handleMessage(String payload) { _log.info(payload); } }
问题排查与修复方案
1. 适配Liferay的OSGI容器环境
Liferay基于OSGI运行,仅使用OSGI的@Component无法让Spring容器正确识别@KafkaListener注解。需要给接收器类同时添加Spring的@Component注解,确保Spring和OSGI容器都能识别该组件:
import org.springframework.stereotype.Component; // ... @Component @org.osgi.service.component.annotations.Component(immediate = true) public class KafkaMessageReceiver { // ... }
2. 补充Kafka消费者核心配置
当前配置仅提供了Kafka Streams的配置,但@KafkaListener依赖独立的消费者工厂和容器工厂。在KafkaConfig中添加以下Bean:
import org.apache.kafka.common.serialization.StringDeserializer; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import static org.apache.kafka.clients.consumer.ConsumerConfig.*; // ... @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(GROUP_ID_CONFIG, "liferay-consumer-group"); props.put(KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; }
3. 移除不必要的注解依赖
若未在Liferay中配置Spring异步和事务管理器,@Async和@Transactional会导致监听方法无法执行。暂时移除这两个注解,待环境配置完善后再添加:
@KafkaListener( topics = "liferay-topic", concurrency = "2" ) public void handleMessage(String payload) { _log.info(payload); }
4. 补充消息推送实现
若需要实现消息推送功能,在KafkaConfig中添加生产者配置:
import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import static org.apache.kafka.clients.producer.ProducerConfig.*; // ... @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); configProps.put(KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); }
然后创建消息发送组件:
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Component; @Component @org.osgi.service.component.annotations.Component(immediate = true) public class KafkaMessageSender { @Autowired private KafkaTemplate<String, String> kafkaTemplate; public void sendMessage(String topic, String message) { kafkaTemplate.send(topic, message); } }
5. 验证步骤
- 确认Kafka服务正常运行,
liferay-topic已通过命令行或Kafka管理工具创建 - 启动Liferay后,查看日志确认消费者是否成功连接Kafka集群
- 使用Kafka命令行推送消息:
kafka-console-producer.sh --broker-list localhost:9092 --topic liferay-topic,输入消息后检查Liferay日志是否打印消息内容
内容的提问来源于stack exchange,提问作者Kiran Kumar
相关产品推荐
相关产品推荐

