Kafka Connect MQTT源连接器use.schema.id配置无效及参数疑问
问题解答
一、use.schema.id的作用
use.schema.id是Confluent Schema Registry相关转换器(如JsonSchemaConverter)的专属配置项,作用是强制指定使用Schema Registry中固定ID对应的Schema来处理数据的序列化/反序列化。设置该值后,转换器会直接从Schema Registry拉取对应ID的Schema,不再自动推导数据结构、注册新Schema,也不会默认使用最新版本的Schema。
二、配置不生效的问题解决
你遇到的核心问题是转换器配置项缺少正确前缀,导致use.schema.id和auto.register.schemas未被JsonSchemaConverter识别,依然使用默认值。
错误原因
在Kafka Connect的配置规则中,转换器的参数必须添加对应的前缀:
- 针对值转换器(
value.converter)的配置,必须以value.converter.作为前缀 - 你当前直接写
use.schema.id = 1和auto.register.schemas = false,这些配置会被当作连接器本身的参数,而非传递给JsonSchemaConverter,因此转换器没有应用你的自定义设置。
修正后的配置
给相关转换器配置添加value.converter.前缀,修改后的完整配置如下:
{ "name": "TEST12", "config": { "connector.class": "io.confluent.connect.mqtt.MqttSourceConnector", "value.converter.schema.registry.url": "http://13.127.201.183:30068", "confluent.topic.bootstrap.servers": "aepv6-base-shell-kafka:9092", "value.converter.auto.register.schemas": "false", "tasks.max": "2", "name": "TEST12", "kafka.topic": "TEST-KAFKA-12", "mqtt.topics": "TEST-MQTT", "mqtt.server.uri": "tcp://broker.hivemq.com:1883", "value.converter": "io.confluent.connect.json.JsonSchemaConverter", "confluent.license.topic.replication.factor": "1", "value.converter.use.schema.id": "1" } }
额外验证步骤
- 确认Schema Registry中ID=1的Schema确实存在,且结构与你预期一致
- 重启连接器,让新配置生效
- 查看连接器任务日志,确认转换器已正确加载指定ID的Schema
内容的提问来源于stack exchange,提问作者Prabhat Sharma
相关产品推荐
相关产品推荐

