如何在多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:9092activate方法中读取: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
相关产品推荐
相关产品推荐

