You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 15:45:42