如何配置Kafka MQTT connector以订阅全部MQTT topic?
问题解决说明
Confluent MQTT Source Connector 本身支持使用#通配符订阅多个MQTT topic,你遇到的两类报错属于不同问题,分别对应如下解决方式:
1. 加引号配置时报通配符非法错误
配置mqtt.topics="#"或mqtt.topics="topic/#"时抛出的Invalid usage of multi-level wildcard错误,是因为配置文件中额外添加的引号会被连接器识别为topic字符串的一部分,不符合MQTT topic命名规范,同时导致通配符#无法被正常识别。
- 解决方法:
mqtt.topics配置项不需要加任何引号,直接填写通配符即可。
2. 去掉引号后报连接拒绝错误
Connection refused报错和通配符配置无关,属于连接器无法正常连接MQTT broker导致,按照以下顺序排查即可:
- 确认MQTT broker的8883端口正常监听,本地防火墙/安全策略没有拦截127.0.0.1的连接请求
- 确认配置的SSL证书路径、密码正确,证书已被MQTT broker信任
- 先用mosquitto_sub、MQTTX等本地客户端工具,使用和连接器完全相同的连接地址、用户名密码、SSL配置,测试订阅
#通配符是否正常,优先排除broker侧连通性、权限问题 - 若你的MQTT broker开启了topic权限控制,需要确认使用的账号有订阅根通配符
#的权限,多数生产环境broker会默认禁止普通账号订阅根通配符避免性能风险
正确配置示例
连通性问题解决后,mqtt.topics直接填写#即可:
name=mqttConnector tasks.max=1 connector.class=io.confluent.connect.mqtt.MqttSourceConnector mqtt.server.uri=ssl://<你的MQTT broker实际可访问地址>:8883 mqtt.topics=# mqtt.username=my-username mqtt.password=my-password confluent.topic.bootstrap.servers=localhost:9092 confluent.topic.replication.factor=1 # Auto topic Creation topic.creation.enable=true topic.creation.default.replication.factor=-1 topic.creation.default.partitions=-1 # SSL mqtt.ssl.trust.store.path=myStore mqtt.ssl.trust.store.password=my-password mqtt.ssl.key.store.path=myStore mqtt.ssl.key.store.password=my-password mqtt.ssl.key.password=
内容的提问来源于stack exchange,提问作者cylon86
相关产品推荐
相关产品推荐

