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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:15:06