基于Jakarta EE与JBoss/WildFly的Kafka消息服务示例代码及文档请求
在Jakarta EE + WildFly中集成Kafka处理海量消息
一、环境依赖与WildFly配置
1. 项目依赖(Maven)
在pom.xml中添加匹配Kafka集群版本的客户端依赖:
<dependencies> <!-- Jakarta EE API --> <dependency> <groupId>jakarta.platform</groupId> <artifactId>jakarta.jakartaee-api</artifactId> <version>10.0.0</version> <scope>provided</scope> </dependency> <!-- Kafka客户端 --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.6.1</version> <scope>compile</scope> </dependency> </dependencies>
2. WildFly模块配置(可选)
若需将Kafka客户端作为WildFly全局模块(避免重复打包),在$WILDFLY_HOME/modules/system/layers/base/org/apache/kafka/main下创建module.xml:
<?xml version="1.0" encoding="UTF-8"?> <module xmlns="urn:jboss:module:1.9" name="org.apache.kafka"> <resources> <resource-root path="kafka-clients-3.6.1.jar"/> <resource-root path="lz4-java-1.8.0.jar"/> <resource-root path="snappy-java-1.1.10.3.jar"/> </resources> <dependencies> <module name="javax.api"/> <module name="javax.transaction.api"/> <module name="org.slf4j"/> </dependencies> </module>
然后在项目的jboss-deployment-structure.xml中声明依赖:
<?xml version="1.0" encoding="UTF-8"?> <jboss-deployment-structure> <deployment> <dependencies> <module name="org.apache.kafka" export="true"/> </dependencies> </deployment> </jboss-deployment-structure>
二、Kafka生产者示例(CDI实现)
用CDI管理生产者实例,搭配批量发送优化海量消息吞吐量:
import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; 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 java.util.Properties; @ApplicationScoped public class KafkaBatchProducer { private final KafkaProducer<String, String> producer; @Inject org.eclipse.microprofile.config.Config config; public KafkaBatchProducer() { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, config.getValue("kafka.bootstrap.servers", String.class)); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 批量发送优化 props.put(ProducerConfig.BATCH_SIZE_CONFIG, "32768"); // 32KB批量阈值 props.put(ProducerConfig.LINGER_MS_CONFIG, "10"); // 等待10ms凑齐批量 props.put(ProducerConfig.ACKS_CONFIG, "1"); // 平衡可靠性与性能 props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); // 幂等性防止重复发送 this.producer = new KafkaProducer<>(props); } // 异步单条发送 public void send(String topic, String key, String value) { producer.send(new ProducerRecord<>(topic, key, value), (metadata, e) -> { if (e != null) { // 异常处理:日志记录、重试逻辑 e.printStackTrace(); } }); } // 批量发送 public void sendBatch(String topic, Iterable<ProducerRecord<String, String>> records) { records.forEach(producer::send); producer.flush(); // 强制提交批量 } // 应用 shutdown 时关闭生产者 public void close() { producer.close(); } }
在src/main/resources/META-INF/microprofile-config.properties配置Kafka地址:
kafka.bootstrap.servers=localhost:9092
三、Kafka消费者示例(Singleton EJB实现)
用@Startup EJB启动消费线程,结合批量拉取处理海量消息:
import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import jakarta.ejb.Singleton; import jakarta.ejb.Startup; import jakarta.inject.Inject; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @Singleton @Startup public class KafkaBatchConsumer { private KafkaConsumer<String, String> consumer; private ExecutorService executor; @Inject org.eclipse.microprofile.config.Config config; @PostConstruct public void init() { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, config.getValue("kafka.bootstrap.servers", String.class)); props.put(ConsumerConfig.GROUP_ID_CONFIG, "wildfly-jakarta-ee-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 批量消费优化 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "1000"); // 单次拉取1000条 props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, "65536"); // 64KB才触发拉取 props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // 手动提交偏移量 consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("massive-data-topic")); // 启动独立线程消费,避免阻塞EJB初始化 executor = Executors.newSingleThreadExecutor(); executor.submit(this::consume); } private void consume() { while (!Thread.currentThread().isInterrupted()) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200)); if (!records.isEmpty()) { processBatch(records); // 异步提交偏移量,不阻塞消费 consumer.commitAsync((offsets, e) -> { if (e != null) { // 提交失败处理:日志、重试 e.printStackTrace(); } }); } } } private void processBatch(ConsumerRecords<String, String> records) { // 自定义批量业务逻辑:如批量写入DB、调用业务服务 records.forEach(record -> { System.out.printf("处理消息:key=%s, value=%s, offset=%d%n", record.key(), record.value(), record.offset()); }); } @PreDestroy public void shutdown() { executor.shutdownNow(); consumer.close(); } }
四、海量消息处理关键优化
- 生产者端:
- 启用批量发送+延迟等待,减少网络IO次数
- 开启幂等性与事务(如需Exactly-Once语义)
- 合理配置
acks参数,避免过度等待确认
- 消费者端:
- 关闭自动提交偏移量,手动提交确保消息不丢失
- 增大拉取批量,减少集群交互
- 用线程池异步处理消息,避免阻塞消费线程
- WildFly配置:
- 在
standalone.xml中调整managed-executor-service线程池大小 - 调大JVM堆内存(如
-Xms4G -Xmx8G),避免OOM - 配置Kafka客户端日志级别,减少不必要的日志输出
- 在
内容的提问来源于stack exchange,提问作者EbrahimFarahmand
相关产品推荐
相关产品推荐

