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

寻求支持Google Protobuf的JMeter/Gatling Kafka性能测试插件

针对Spring Boot + Kafka + Protobuf的性能测试工具方案

一、扩展Gatling Kafka插件支持Protobuf

Gatling的gatling-kafka-plugin本身没有内置Protobuf支持,但可以通过自定义序列化逻辑快速实现:

  • 实现KafkaMessageSerializer接口,在序列化方法中调用Protobuf对象的toByteArray()方法转换为字节数组:
import io.gatling.kafka.serialization.KafkaMessageSerializer
import com.your.protobuf.generated.MessageProto.MyMessage

class ProtobufSerializer extends KafkaMessageSerializer[MyMessage] {
  override def serialize(topic: String, data: MyMessage): Array[Byte] = {
    data.toByteArray
  }
}
  • 在Gatling测试脚本中指定自定义序列化器:
kafka("Protobuf Message Send")
  .send(
    KafkaMessage(
      topic = "test-topic",
      payload = MyMessage.newBuilder().setField1("test").build(),
      key = Some("key1"),
      serializer = new ProtobufSerializer()
    )
  )

二、用Locust + 自定义Kafka客户端

Locust的灵活性很高,适合快速编写自定义测试逻辑:

  • 依赖confluent-kafka和Protobuf生成的Python类,编写Locust任务:
from locust import HttpUser, task, between
from confluent_kafka import Producer
import message_pb2

class KafkaProtobufUser(HttpUser):
    wait_time = between(0.1, 0.5)

    def on_start(self):
        self.producer = Producer({"bootstrap.servers": "localhost:9092"})

    @task
    def send_protobuf_message(self):
        msg = message_pb2.MyMessage()
        msg.field1 = "test-load"
        serialized_msg = msg.SerializeToString()
        self.producer.produce("test-topic", value=serialized_msg)
        self.producer.flush()

三、修复KLoadGen的使用问题

如果KLoadGen不能正常工作,大概率是配置项错误,检查以下几点:

  • 确保Kafka生产者的key.serializer和value.serializer都设置为org.apache.kafka.common.serialization.ByteArraySerializer
  • 在KLoadGen的配置中,正确指定:
    • messageClassName:Protobuf生成类的全限定名(比如com.your.protobuf.generated.MessageProto$MyMessage)
    • protoPath:.proto文件所在的目录路径
    • 若需要动态生成测试数据,可通过messageSchema配置Protobuf字段的生成规则

四、自定义Java测试脚本 + 并发调度

如果以上工具都不适用,可以写一个简单的原生Java脚本,结合ExecutorService实现并发发送:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import com.your.protobuf.generated.MessageProto.MyMessage;

import java.util.Properties;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class ProtobufKafkaLoadTest {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");

        ExecutorService executor = Executors.newFixedThreadPool(50);
        KafkaProducer<byte[], byte[]> producer = new KafkaProducer<>(props);

        for (int i = 0; i < 10000; i++) {
            executor.submit(() -> {
                MyMessage msg = MyMessage.newBuilder().setField1("test").build();
                producer.send(new ProducerRecord<>("test-topic", msg.toByteArray()));
            });
        }

        executor.shutdown();
        producer.close();
    }
}

可以用多实例运行脚本的方式,模拟高并发负载。

内容的提问来源于stack exchange,提问作者tosi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 01:34:54