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

使用JMeter向Kafka Topic发送Avro消息遇Schema检索错误求助

解决JMeter发送Avro消息到Kafka时的Schema获取错误

错误根源

你遇到的Error retrieving Avro schema"string"错误,核心原因是:

  • 通过vars.get("MESSAGE")获取到的是普通字符串,但KafkaAvroSerializer要求传入的是符合预注册Avro Schema的GenericRecord/SpecificRecord对象,而非原始字符串。
  • 序列化器误将传入的字符串当作Schema标识符,尝试去Schema Registry查找名为"string"的Schema,自然找不到导致失败。

解决方案

方案1:将JSON字符串转换为Avro GenericRecord(适用于大多数场景)

如果你的MESSAGE是符合目标Topic Schema的JSON字符串,需要先将其解析为Avro对象再发送:

  1. 确保依赖齐全:JMeter的lib目录下必须包含以下jar包(版本需与Confluent Platform匹配):

    • confluent-kafka-avro-serializer
    • avro
    • kafka-clients
    • schema-registry-client
    • jackson-databind
  2. 修改JSR223 Sampler代码:

import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient;
import io.confluent.kafka.schemaregistry.client.SchemaRegistryClient;
import io.confluent.kafka.serializers.AbstractKafkaSchemaSerDeConfig;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.util.Map;

// 从用户变量获取配置
String brokers = vars.get("KAFKA_BROKERS");
String topic = vars.get("KAFKA_TOPIC");
String schemaRegistryUrl = vars.get("SCHEMA_REGISTRY_URL"); // 建议将此加入用户变量
String schemaRegistryAuth = vars.get("SCHEMA_REGISTRY_AUTH"); // 格式:username:password,加入用户变量
String msgJson = vars.get("MESSAGE");
String user = String.valueOf(ctx.getThreadNum() + 1);

// 初始化Schema Registry客户端
SchemaRegistryClient schemaRegistry = new CachedSchemaRegistryClient(
    schemaRegistryUrl,
    100,
    Map.of(
        AbstractKafkaSchemaSerDeConfig.BASIC_AUTH_CREDENTIALS_SOURCE, "USER_INFO",
        AbstractKafkaSchemaSerDeConfig.USER_INFO_CONFIG, schemaRegistryAuth
    )
);

// 获取Topic对应的最新Value Schema(默认命名规则:topic-name-value)
Schema schema = schemaRegistry.getLatestSchemaMetadata(topic + "-value").getSchema();

// 将JSON字符串转换为Avro GenericRecord
ObjectMapper mapper = new ObjectMapper();
Map<String, Object> msgMap = mapper.readValue(msgJson, Map.class);
GenericRecord avroRecord = new GenericData.Record(schema);
for (Schema.Field field : schema.getFields()) {
    avroRecord.put(field.name(), msgMap.get(field.name()));
}

// 配置Kafka Producer
Properties kafkaProps = new Properties();
kafkaProps.put("bootstrap.servers", brokers);
kafkaProps.put("schema.registry.url", schemaRegistryUrl);
kafkaProps.put("auto.register.schemas", "false");
kafkaProps.put("basic.auth.credentials.source", "USER_INFO");
kafkaProps.put("basic.auth.user.info", schemaRegistryAuth);
kafkaProps.put("security.protocol", "SASL_SSL");
kafkaProps.put("sasl.mechanism", "PLAIN");
kafkaProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
kafkaProps.put("value.serializer", "io.confluent.kafka.serializers.KafkaAvroSerializer");
kafkaProps.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username='<>' password='<>';");

// 发送消息
Producer<String, Object> producer = new KafkaProducer<>(kafkaProps);
try {
    producer.send(new ProducerRecord<>(topic, user, avroRecord)).get();
} finally {
    producer.close();
}

方案2:直接发送字符串(仅当Topic Schema为string类型时适用)

如果你的Topic的Value Schema本身就是{"type": "string"},可以:

  • 确认Schema Registry中已注册该Schema
  • 或者,若无需Avro序列化,直接将value.serializer改为org.apache.kafka.common.serialization.StringSerializer(但此时消息不再是Avro格式)

关键检查点

  • 确认auto.register.schemas=false时,目标Topic对应的Schema已在Registry中存在
  • 验证Schema Registry的认证信息(URL、用户名密码)正确,可通过curl测试连通性
  • 确保JMeter依赖包版本与Confluent Platform、Kafka集群版本兼容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 01:30:20