基于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"); } } }
三、问题排查与验证步骤
前置检查
- 确认Kafka Broker在
localhost:9092正常运行,用命令kafka-topics.sh --list --bootstrap-server localhost:9092检查likes-buckettopic是否存在,不存在则创建:kafka-topics.sh --create --topic likes-bucket --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1 - 检查Kafka Broker的日志(默认路径
kafka/logs/server.log),看是否有连接拒绝、权限错误等信息。
- 确认Kafka Broker在
开启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>
- 运行测试并验证
执行./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
相关产品推荐
相关产品推荐

