使用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对象再发送:
确保依赖齐全:JMeter的
lib目录下必须包含以下jar包(版本需与Confluent Platform匹配):confluent-kafka-avro-serializeravrokafka-clientsschema-registry-clientjackson-databind
修改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
相关产品推荐
相关产品推荐

