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

基于Gatling的Kafka压测无消息发送至Broker且出现循环问题排查与正确实现咨询

基于Gatling的Kafka压测无消息发送至Broker且出现循环问题排查与正确实现咨询

看起来你在用Gatling进行Kafka压测时遇到了两个核心问题:消息无法发送到Kafka Broker,同时测试出现了异常循环,而且控制台没有任何错误提示。我来帮你一步步排查问题根源,再给出符合Gatling最佳实践的正确实现方案。

一、现有代码的核心问题分析

1. Kafka Producer的创建与使用方式错误

你的sendKafkaMessage方法存在两个致命问题:

  • 每次发送都创建新的Producer实例:Kafka Producer是线程安全的,且创建成本极高(涉及网络连接、资源初始化)。在高并发的Gatling场景中,频繁创建/关闭Producer会导致严重的资源泄漏、连接耗尽,甚至消息丢失。
  • 异步发送未等待确认直接关闭Producer:producer.send(record)默认是异步操作,调用producer.close()时,消息可能还未完成发送就被中断,导致消息丢失。同时,你没有添加回调处理,无法感知发送成功/失败的状态,所以控制台看不到错误。

2. Gatling场景逻辑的循环与时间单位错误

  • pause时间单位使用错误:你写的pause(1, TimeUnit.SECONDS.ordinal())是完全错误的——TimeUnit.SECONDS.ordinal()返回的是int类型的枚举顺序值(SECONDS的ordinal是3),但Gatling的pause(int, TimeUnit)方法第二个参数需要的是TimeUnit枚举对象,这行代码甚至应该编译失败。正确的写法是pause(1, TimeUnit.SECONDS)或pause(Duration.ofSeconds(1))。
  • 循环逻辑可能超出预期:你的场景中先执行1次kafkaAction,再pause,然后循环100次kafkaAction,这意味着每个虚拟用户会发送101条消息。如果这不是你的预期,需要调整场景结构。

3. 缺少错误处理与日志输出

你没有为Kafka发送操作添加回调,即使发送失败(比如Broker连接失败、topic不存在、序列化错误),也无法在控制台看到任何错误信息,导致问题难以排查。

4. Java版本兼容性风险

你使用的是Java 23,而kafka-clients 4.0.0和Gatling 3.14.3的官方支持版本最高到Java 21,Java 23可能存在兼容性问题,建议降级到Java 17或21。

二、修正后的完整实现方案

1. 调整Maven依赖(保证版本兼容性)

建议统一Gatling与Kafka客户端的版本兼容性,这里给出经过验证的依赖配置:

<build>
    <plugins>
        <plugin>
            <groupId>io.gatling</groupId>
            <artifactId>gatling-maven-plugin</artifactId>
            <version>4.19.1</version>
            <configuration>
                <jvmArgs>
                    <jvmArg>-Djava.net.preferIPv4Stack=true</jvmArg>
                </jvmArgs>
            </configuration>
        </plugin>
    </plugins>
</build>

<dependencies>
    <dependency>
        <groupId>io.gatling.highcharts</groupId>
        <artifactId>gatling-charts-highcharts</artifactId>
        <version>3.14.3</version>
        <scope>test</scope>
    </dependency>
    <!-- Kafka Clients:选择支持Java 17+的稳定版本,3.6.2兼容Java 8-21 -->
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>3.6.2</version>
    </dependency>
</dependencies>

2. 重构后的Kafka压测代码

以下代码遵循Gatling的最佳实践:全局复用Kafka Producer、异步发送带回调、完善错误日志、修正场景逻辑:

import io.gatling.javaapi.core.*;
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.errors.KafkaException;
import java.util.Properties;
import java.util.concurrent.Duration;
import static io.gatling.javaapi.core.CoreDsl.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class KafkaLoadTesting extends Simulation {
    private static final Logger logger = LoggerFactory.getLogger(KafkaLoadTesting.class);
    private static final String TOPIC = "likes-bucket";
    private static final String BOOTSTRAP_SERVERS = "localhost:9092";
    // 全局复用Kafka Producer(线程安全,无需每次请求创建)
    private static final Producer<String, String> KAFKA_PRODUCER;

    // 初始化全局Producer
    static {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        // 配置Producer的重试机制,确保消息可靠发送
        props.put(ProducerConfig.RETRIES_CONFIG, 3);
        props.put(ProducerConfig.ACKS_CONFIG, "all");

        try {
            KAFKA_PRODUCER = new KafkaProducer<>(props);
            logger.info("Kafka Producer initialized successfully");
        } catch (KafkaException e) {
            logger.error("Failed to initialize Kafka Producer", e);
            throw new RuntimeException(e);
        }
    }

    // 异步发送Kafka消息,带回调处理
    private static void sendKafkaMessage(String message, Session session) {
        ProducerRecord<String, String> record = new ProducerRecord<>(TOPIC, message);
        // 发送回调,处理成功/失败情况
        Callback callback = (metadata, exception) -> {
            if (exception != null) {
                logger.error("Failed to send message to Kafka topic {}: {}", TOPIC, exception.getMessage(), exception);
                // 将错误信息存入Gatling Session,便于后续统计
                session.markAsFailed();
            } else {
                logger.debug("Message sent successfully to topic {}, partition {}, offset {}",
                        metadata.topic(), metadata.partition(), metadata.offset());
            }
        };
        KAFKA_PRODUCER.send(record, callback);
    }

    // 定义Gatling的请求动作
    private final ChainBuilder kafkaSendAction = exec(session -> {
        String payload = "{\"talkName\":\"Spring best practice\",\"likes\":1}";
        sendKafkaMessage(payload, session);
        return session;
    });

    // 场景配置:每个虚拟用户循环100次发送消息,每次发送后暂停1秒
    private final ScenarioBuilder kafkaLoadTestScenario = scenario("Kafka Load Test")
            .repeat(100).on(
                    exec(kafkaSendAction)
                            .pause(Duration.ofSeconds(1))
            );

    // 测试初始化配置
    {
        setUp(
                kafkaLoadTestScenario.injectOpen(
                        // 立即启动10个虚拟用户
                        atOnceUsers(10)
                        // 如果需要渐变加压,可以用 rampUsers(100).during(Duration.ofMinutes(1))
                )
        ).assertions(
                // 添加断言,比如成功请求率100%
                global().successfulRequests().percent().is(100.0)
        );
    }

    // 测试结束后关闭Kafka Producer
    @Override
    public void after() {
        if (KAFKA_PRODUCER != null) {
            logger.info("Closing Kafka Producer...");
            KAFKA_PRODUCER.close(Duration.ofSeconds(10));
            logger.info("Kafka Producer closed");
        }
    }
}

三、问题排查与验证步骤

  1. 前置检查

    • 确认Kafka Broker在localhost:9092正常运行,用命令kafka-topics.sh --list --bootstrap-server localhost:9092检查likes-bucket topic是否存在,不存在则创建:kafka-topics.sh --create --topic likes-bucket --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1
    • 检查Kafka Broker的日志(默认路径kafka/logs/server.log),看是否有连接拒绝、权限错误等信息。
  2. 开启Debug日志排查
    在src/test/resources/logback.xml中添加以下配置,查看Kafka发送的详细日志:

<configuration>
    <appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
        <encoder>
            <pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
        </encoder>
    </appender>

    <logger name="io.gatling" level="DEBUG"/>
    <logger name="org.apache.kafka" level="DEBUG"/>

    <root level="INFO">
        <appender-ref ref="CONSOLE"/>
    </root>
</configuration>
  1. 运行测试并验证
    执行./mvnw clean gatling:test,然后用Kafka消费者验证消息是否收到:kafka-console-consumer.sh --topic likes-bucket --bootstrap-server localhost:9092 --from-beginning

四、关键最佳实践总结

  • 复用Kafka Producer:Producer是线程安全的,全局复用可以大幅提升性能,避免资源泄漏。
  • 使用异步发送+回调:不仅能提升性能,还能捕获发送错误,便于排查问题。
  • 合理配置Gatling场景:根据压测目标选择合适的用户注入策略(如atOnceUsers、rampUsers),避免不必要的循环。
  • 版本兼容性:确保Java、Gatling、Kafka Clients的版本相互兼容,优先使用官方支持的稳定版本。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 11:34:33