寻求Apache Kafka与MQTT双向桥接的Confluent MQTT连接器开源替代方案
开源免费的MQTT-Kafka双向桥接方案
当然有!完全不需要依赖Confluent的付费连接器,我给你整理了几个靠谱的开源免费方案,覆盖不同场景需求:
1. EMQX开源版的Kafka桥接插件
EMQX本身就是一款高性能的开源MQTT Broker,它自带官方维护的Kafka桥接插件,完美支持双向数据转发(MQTT → Kafka、Kafka → MQTT),非常适合中大规模的生产场景。
- 配置简单:安装EMQX开源版后,通过Web控制台或者配置文件就能快速启用桥接,只需要指定Kafka集群地址、主题映射规则(比如把MQTT的
sensor/#主题转发到Kafka的sensor-data主题,反之亦然),还能自定义消息格式(JSON、Protobuf等)。 - 优势:不用额外部署复杂组件,EMQX本身兼顾MQTT Broker和桥接功能,稳定性和性能都有保障,社区支持也很活跃。
2. 社区维护的Kafka Connect MQTT连接器(完全开源免费)
如果你的技术栈已经基于Kafka Connect,那可以用社区维护的开源连接器替代Confluent的付费版本,首推jcustenborder/kafka-connect-mqtt:
- 这是一个完全开源的项目,同时支持MQTT源连接器(从MQTT拉取数据到Kafka)和Sink连接器(从Kafka推送数据到MQTT),兼容Kafka 2.x及以上版本。
- 部署方式和普通Kafka Connect连接器一致:把编译好的jar包放到Connect的
plugin.path目录下,然后提交配置文件即可。 - 简单配置示例(MQTT源连接器):
{ "name": "mqtt-source-connector", "config": { "connector.class": "com.github.jcustenborder.kafka.connect.mqtt.MqttSourceConnector", "tasks.max": "1", "mqtt.server.uris": "tcp://your-mqtt-broker:1883", "mqtt.topics": "device/+/data", "kafka.topic": "mqtt-ingest-topic", "mqtt.qos": "1" } }
3. 自定义轻量桥接脚本(适合小型/定制化场景)
如果你的场景数据量不大,或者需要高度定制消息处理逻辑,自己写个简单的脚本是最灵活的选择,用Python、Node.js这类语言都能快速实现:
- 比如用Python的
paho-mqtt库连接MQTT Broker,kafka-python库连接Kafka,开两个线程分别处理双向数据流:import paho.mqtt.client as mqtt from kafka import KafkaProducer, KafkaConsumer import threading # MQTT → Kafka 转发线程 def forward_mqtt_to_kafka(): mqtt_client = mqtt.Client(client_id="mqtt-kafka-bridge") mqtt_client.connect("your-mqtt-broker", 1883, 60) mqtt_client.subscribe("sensor/#") kafka_producer = KafkaProducer(bootstrap_servers="your-kafka-broker:9092") def on_mqtt_message(client, userdata, msg): kafka_producer.send("mqtt-to-kafka-topic", msg.payload) print(f"Forwarded MQTT message to Kafka: {msg.topic} -> mqtt-to-kafka-topic") mqtt_client.on_message = on_mqtt_message mqtt_client.loop_forever() # Kafka → MQTT 转发线程 def forward_kafka_to_mqtt(): kafka_consumer = KafkaConsumer( "kafka-to-mqtt-topic", bootstrap_servers="your-kafka-broker:9092", group_id="kafka-mqtt-bridge-group" ) mqtt_client = mqtt.Client(client_id="kafka-mqtt-bridge") mqtt_client.connect("your-mqtt-broker", 1883, 60) for kafka_msg in kafka_consumer: mqtt_client.publish("kafka-output-topic", kafka_msg.value) print(f"Forwarded Kafka message to MQTT: kafka-to-mqtt-topic -> kafka-output-topic") # 启动双向转发线程 threading.Thread(target=forward_mqtt_to_kafka, daemon=True).start() threading.Thread(target=forward_kafka_to_mqtt, daemon=True).start() # 保持脚本运行 input("Press Enter to exit...\n") - 优势:完全可控,能根据需求添加消息过滤、格式转换、异常处理等逻辑,适合快速验证或小型项目。
4. Apache NiFi(可视化数据流编排)
如果你更倾向于无代码/低代码的解决方案,Apache NiFi是个不错的选择——它是一款开源的数据流编排工具,支持MQTT和Kafka的双向对接,还能做复杂的数据处理。
- 操作方式:通过拖拽组件就能搭建数据流:用
ConsumeMQTT拉取MQTT消息,经过转换/过滤后用PublishKafkaRecord发送到Kafka;反过来用ConsumeKafkaRecord消费Kafka消息,再用PublishMQTT推送到MQTT。 - 优势:可视化界面友好,无需编写代码,适合非开发人员或者需要复杂数据处理的场景,稳定性和扩展性都很强。
选择建议
- 大规模生产环境:优先选EMQX开源版或者社区版Kafka Connect连接器,稳定性和性能更有保障;
- 小型/定制化场景:自定义脚本最灵活;
- 可视化/低代码需求:Apache NiFi是首选。
内容的提问来源于stack exchange,提问作者Nitesh Ratnaparkhe
相关产品推荐
相关产品推荐

