# MQTT 到 Kafka Bridge 实验

这个仅绑定本机回环地址的实验展示了一条可验证的 MQTT 到 Kafka 接入边界，使用：

- Eclipse Mosquitto 2.0.22 提供 MQTT；
- Redpanda 25.1.9 提供兼容 Kafka 的事件日志和 HTTP Proxy；
- MQTT.js 5.15.2 与 Node.js `fetch` 实现 Bridge。

实验不需要真实凭据或生产端点。

## 运行实验

在 RunMQTT 仓库根目录执行：

```bash
docker compose -f public/examples/mqtt-kafka-bridge/compose.yaml up -d
pnpm mqtt:kafka-bridge:self-test
pnpm mqtt:kafka-bridge
```

等待 Bridge 输出 `bridge_ready`。另开一个终端，发布一条合成的有效事件：

```bash
docker compose -f public/examples/mqtt-kafka-bridge/compose.yaml exec mosquitto \
  mosquitto_pub -h 127.0.0.1 -t lab/device-042/telemetry -q 1 \
  -m '{"eventId":"evt-001","recordedAt":"2026-01-01T00:00:00Z","value":21.4}'
```

Bridge 会校验载荷，从 MQTT Topic 推导 `deviceId`，添加 `schemaVersion` 和 `ingestedAt`，再以 `device-042` 为 Kafka Record Key 写入 `iot.telemetry`。查看一条记录：

```bash
docker compose -f public/examples/mqtt-kafka-bridge/compose.yaml exec redpanda \
  rpk topic consume iot.telemetry --num 1 --brokers redpanda:9092
```

向同一个 MQTT Topic 发布 `{}`，再消费 `iot.telemetry.dlq`，即可验证 Schema 拒绝路径。DLQ Envelope 会记录失败原因、字节数和 SHA-256 摘要，但不会复制被拒绝的载荷。

按 Ctrl+C 停止 Bridge，然后删除这些一次性容器：

```bash
docker compose -f public/examples/mqtt-kafka-bridge/compose.yaml down -v
```

## 示例验证了什么

- MQTT Topic 可以规范化为规模更小的 Kafka Topic 分类；
- 使用已授权 Topic 路径而非不可信载荷字段提供 Partition Key；
- 有效事件与校验失败会走可观测、相互独立的路径；
- HTTP Proxy 瞬时失败最多尝试四次，并采用指数退避；
- 同一设备的事件使用相同 Kafka Record Key，因此进入同一 Partition。

## 生产限制

这个实验刻意没有实现生产 Connector：

- Mosquitto 允许匿名访问，因为它只绑定 `127.0.0.1`。生产环境必须认证每个发布方，并实施设备级 Topic ACL。
- MQTT 投递确认与异步 Kafka 写入相互独立。Bridge 在两者之间崩溃可能丢失记录。如果该边界不允许丢失，请使用受支持的 Connector 或持久本地 Spool/Outbox。
- Bridge 是单进程，不包含 Leader Election、流量准入限制、Schema Registry 集成、持久重试存储或密钥管理。
- 如果 Kafka 本身不可用，示例也无法写入 Kafka DLQ。生产环境需要独立的持久失败路径和运维告警。
- DLQ 不是重试循环。请限制访问、设置保留时间、明确负责人，并在修复根因后再重放。
- HTTP 端点与匿名测试 Broker 只绑定本机回环地址。不要把此 Compose 项目暴露到共享网络。

完整的身份、Topic、Schema、Partition、重试、DLQ、可观测性与指令返回边界，请参阅 [MQTT 到 Kafka 架构指南](/zh-Hans/mqtt-vs-kafka)。
