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

寻求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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 14:42:55