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

如何在多OSGi模块的Liferay项目中配置Kafka生产者?

Liferay作为Kafka生产者的配置实现方案

针对你在大型Liferay工作区搭建事件驱动通知系统的需求,以下是Liferay侧作为Kafka生产者的具体配置与实现步骤:

1. 准备OSGi兼容的Kafka依赖

Liferay基于OSGi容器,需确保Kafka客户端依赖能被OSGi识别:

  • 在目标OSGi模块的build.gradle中添加依赖:
    dependencies {
        compileOnly group: 'org.apache.kafka', name: 'kafka-clients', version: '匹配你的Kafka版本'
        compileOnly group: 'com.liferay.portal', name: 'com.liferay.portal.kernel', version: '匹配你的Liferay版本'
    }
    
  • 在模块的bnd.bnd中配置包导入,让OSGi容器能找到Kafka相关类:
    Import-Package: \
        org.apache.kafka.clients.producer,\
        org.apache.kafka.common.serialization,\
        com.liferay.portal.kernel.util,\
        *
    
    若官方kafka-clients不是OSGi bundle,可借助biz.aQute.bnd工具将其转换为OSGi兼容包。

2. 实现Kafka生产者OSGi组件

创建一个可被其他模块调用的Kafka生产者服务,通过OSGi注解声明:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import com.liferay.portal.kernel.log.Log;
import com.liferay.portal.kernel.log.LogFactoryUtil;
import org.osgi.service.component.annotations.Activate;
import org.osgi.service.component.annotations.Component;
import org.osgi.service.component.annotations.Deactivate;

import java.util.Properties;

@Component(immediate = true, service = KafkaEventProducer.class)
public class KafkaEventProducer {

    private static final Log _log = LogFactoryUtil.getLog(KafkaEventProducer.class);
    private KafkaProducer<String, String> producer;

    @Activate
    protected void activate() {
        Properties props = new Properties();
        // 配置Kafka broker地址(后续可外置化)
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        // 可选配置:acks=all、retries=3等,根据业务可靠性需求调整

        producer = new KafkaProducer<>(props);
    }

    @Deactivate
    protected void deactivate() {
        if (producer != null) {
            producer.close();
        }
    }

    // 对外提供的事件发送方法
    public void sendEvent(String topic, String key, String eventPayload) {
        ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, eventPayload);
        producer.send(record, (metadata, exception) -> {
            if (exception != null) {
                _log.error("Kafka事件发送失败", exception);
            } else {
                _log.debug("Kafka事件发送成功:topic=" + metadata.topic() + ", offset=" + metadata.offset());
            }
        });
    }
}

3. 触发事件发送的业务场景

根据你的业务需求,在Liferay中调用生产者服务发送事件:

场景1:监听Liferay实体CRUD事件

通过Liferay的ModelListener监听内置或自定义实体的生命周期事件:

import com.liferay.portal.kernel.model.BaseModelListener;
import com.liferay.portal.kernel.model.User;
import org.osgi.service.component.annotations.Component;
import org.osgi.service.component.annotations.Reference;

@Component(immediate = true)
public class UserModelListener extends BaseModelListener<User> {

    @Reference
    private KafkaEventProducer kafkaEventProducer;

    @Override
    public void onAfterCreate(User user) {
        // 构造JSON格式的事件 payload
        String payload = String.format(
            "{\"userId\": %d, \"action\": \"CREATE\", \"username\": \"%s\"}",
            user.getUserId(), user.getScreenName()
        );
        kafkaEventProducer.sendEvent("user-events", String.valueOf(user.getUserId()), payload);
    }
}

场景2:自定义业务逻辑主动触发

在自定义Portlet或Service中主动调用生产者:

import com.liferay.portal.kernel.portlet.bridges.mvc.MVCActionCommand;
import com.liferay.portal.kernel.portlet.bridges.mvc.BaseMVCActionCommand;
import javax.portlet.ActionRequest;
import javax.portlet.ActionResponse;
import org.osgi.service.component.annotations.Component;
import org.osgi.service.component.annotations.Reference;

@Component(
    property = {
        "javax.portlet.name=你的自定义Portlet名称",
        "mvc.command.name=sendCustomEvent"
    },
    service = MVCActionCommand.class
)
public class SendCustomEventMVCActionCommand extends BaseMVCActionCommand {

    @Reference
    private KafkaEventProducer kafkaEventProducer;

    @Override
    protected void doProcessAction(ActionRequest request, ActionResponse response) throws Exception {
        String customPayload = "{\"eventType\": \"CUSTOM_BIZ_ACTION\", \"data\": \"业务数据内容\"}";
        kafkaEventProducer.sendEvent("custom-business-events", "custom-key", customPayload);
    }
}

4. 配置优化与最佳实践

  • 配置外置化:避免硬编码Kafka地址,通过Liferay配置文件读取:
    在portal-ext.properties中添加:
    kafka.bootstrap.servers=kafka-broker-1:9092,kafka-broker-2:9092
    
    然后在生产者的activate方法中读取:
    import com.liferay.portal.kernel.util.PropsUtil;
    
    @Activate
    protected void activate() {
        String bootstrapServers = PropsUtil.get("kafka.bootstrap.servers");
        // 初始化props...
    }
    
  • 序列化优化:若事件为复杂对象,替换StringSerializer为JsonSerializer,并确保Spring Boot消费者使用对应反序列化器。
  • 性能调优:根据业务吞吐量调整batch.size、linger.ms、buffer.memory等参数,平衡延迟与发送效率。

5. 部署与测试

  • 将OSGi模块打包为jar,部署到Liferay的osgi/modules目录。
  • 启动Liferay与Kafka集群,触发对应业务事件(如创建用户),验证Kafka Topic是否收到消息,同时确认Spring Boot消费者正常处理。

内容的提问来源于stack exchange,提问作者Jithin M V

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 11:55:22