概述
Azure IoT Hub Sink Connector 从 Kafka Topic 读取云到设备消息,并将每条记录转换为一条 Azure IoT Hub 消息发送到记录指定的设备。它位于 Kafka 事件流和 IoT Hub 设备消息通道之间,适合将业务系统产生的设备指令、通知或配置更新转发到 IoT Hub。 记录值可以使用 Struct 或 JSON 字符串表示。消息至少包含messageId、message 和 deviceId;目标设备由 deviceId 选择。Struct 中可使用 expiry 指定消息过期时间,JSON 字符串中对应的属性名为 expiryTime。Connector 不负责设备反馈、直接方法或其他云到设备通信模式。
前置条件
- Azure IoT Hub 中须存在接收消息的目标设备,连接字符串须具有向 IoT Hub 发送云到设备消息所需的权限;如果使用消息过期时间,还须确认目标 Azure IoT Hub 和设备端的业务规则符合预期。
授权许可
使用 MIT License。快速开始
提前准备 Connect Cluster、Kafka 和 Azure IoT Hub,确认 Worker 可以访问输入 Topic 及 Azure IoT Hub,并准备目标设备。集群与 Connector 的管理操作见管理 Connector。下面的配置读取一个 Topic 中的 JSON 字符串,并将每条消息发送到其deviceId 指定的设备。
<input-topic> 替换为输入 Topic,将 <iot-hub-connection-string> 替换为 Azure IoT Hub 连接字符串。不要把连接字符串写入日志或提交到版本控制系统。使用该示例时,Topic 中每条消息的字符串值应为可解析的 JSON,例如:
expiryTime 必须是可解析的 ISO-8601 Instant;省略该属性表示不设置过期时间。输入记录缺少必需字段、字段类型不匹配或 JSON 无法解析时,记录转换会失败。
配置
连接与消息
IotHub.ConnectionString
设置 Azure IoT Hub 连接字符串。
- 类型:
string - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:必须提供可用于创建 Azure IoT Hub ServiceClient 的连接字符串。配置定义不会校验非空、格式或可达性;错误值可能在 Connector 启动或实际发送时失败。该配置包含敏感凭证,请使用安全的配置注入方式,不要在日志中输出。
IotHub.MessageDeliveryAcknowledgement
设置写入 Azure IoT Hub Message 的消息传递确认模式。
- 类型:
string - 默认值:
None - 重要级别:高
- 有效值 / 注意事项:只能使用区分大小写的
None、Full、PositiveOnly或NegativeOnly。该值作为每条 Azure 消息的确认元数据发送;它不是设备已接收、已处理或已反馈的证明,也不会生成 Kafka 确认记录。
输入订阅与任务
connector.class
选择 Azure IoT Hub Sink Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:使用
com.microsoft.azure.iot.kafka.connect.sink.IotHubSinkConnector。
topics
指定要读取的 Kafka Topic 列表。
- 类型:
list - 默认值:空列表
[] - 重要级别:高
- 必填:与
topics.regex二选一 - 有效值 / 注意事项:使用逗号分隔的 Topic 名称。必须与非空的
topics.regex二选一;两者同时设置或同时为空都会导致 Sink 配置校验失败。
topics.regex
使用正则表达式订阅 Kafka Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 必填:与
topics二选一 - 有效值 / 注意事项:使用 Java 正则表达式,并确保表达式非空。必须与非空的
topics二选一;两者同时设置或同时为空都会导致 Sink 配置校验失败。
tasks.max
设置 Connector 可创建的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:必须至少为
1。实际并行度受 Kafka 分区数量和分配结果限制;增加该值不会保证按设备串行,也不会建立跨分区或跨 Task 的全局顺序。
数据转换
key.converter
将 Kafka 记录键转换为 Kafka Connect 数据。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时使用 Worker 级别的键转换器;如果设置,必须是可实例化的
org.apache.kafka.connect.storage.Converter实现。Connector 使用记录值选择设备,不读取记录键来映射目标设备。
value.converter
将 Kafka 记录值转换为 Kafka Connect 数据,供 Connector 映射为 Azure IoT Hub 消息。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时使用 Worker 级别的值转换器;如果设置,必须是可实例化的
org.apache.kafka.connect.storage.Converter实现。转换后的值须符合 Connector 支持的 String 或 Struct 路径。String 值应为包含messageId、message、deviceId和可选expiryTime的 JSON;Struct 值须包含同名的messageId、message、deviceId字符串字段,并可包含字符串类型的expiry字段。
最佳实践
按正则订阅一组输入 Topic
适用业务场景:首次接入多个命名规则一致的消息 Topic,或后续会持续创建同类 Topic,希望 Connector 自动纳入匹配的 Topic,而不需要每次修改固定 Topic 列表。 配置示例:<iot-hub-connection-string> 替换为 Azure IoT Hub 连接字符串。使用 topics.regex 后移除 topics;不要同时配置两个订阅方式。表达式匹配的每个 Topic 中的字符串值都应遵循相同的 JSON 消息结构。
关键说明:使用正则订阅便于统一管理同类 Topic,但 Topic 命名规则会直接决定消费范围。变更命名规则前先确认新旧 Topic 是否会同时匹配,避免意外扩大或缩小输入范围。
按分区增加 Task 并行度
适用业务场景:Connector 已能稳定发送消息,输入 Topic 有多个分区且单个 Task 的逐条同步发送吞吐不足,希望通过增加 Task 并行处理分区。 配置示例:<input-topic> 和 <iot-hub-connection-string> 替换为实际值。先确认输入 Topic 至少有两个可分配的分区,再逐步提高 tasks.max,并观察 Task 状态、发送延迟和消费 Lag。
关键说明:tasks.max 是上限,不是实际并行 Task 数量;Kafka Connect 会根据分区分配结果创建和分配任务。每个 Task 独立发送消息,同一设备跨分区或跨 Task 的处理顺序不受 Connector 保证,因此需要顺序时应在上游设计合适的分区键。
监控
监控内容
监控 Kafka Connect Worker、Connector 和 Task 的健康状态及状态变化,关注输入吞吐、发送延迟、消费 Lag、Offset 提交、错误与重试情况,并同时观察 Worker JVM 的 CPU、内存和垃圾回收信号。由于 Connector 对每条记录同步发送,发送延迟升高可能拖慢 Task 的处理;若启用了 Kafka Connect 错误处理或死信队列,再关注对应的错误处理计数和 DLQ 活动。不要把IotHub.MessageDeliveryAcknowledgement 当作设备反馈指标。
导入 Grafana 大盘
从下载 Connect Cluster Grafana 大盘获取 Dashboard JSON,准备已采集 Kafka Connect 指标的 Prometheus 数据源,并确认标签名称与数据源配置匹配,然后在 Grafana 中导入该 JSON 并选择对应数据源。限制条件
- Connector 只处理云到设备消息,不提供设备反馈、确认结果消费、直接方法或其他 Azure IoT Hub 通信模式。
- 每条输入记录单独转换并同步发送,没有 Connector 级批量请求、异步发送、限速或可配置的发送中消息数量。
topics和topics.regex必须且只能配置一个非空值,不能同时使用或同时省略。- 消息发送成功与 Kafka Offset 提交不是同一个原子操作;进程在发送完成但 Offset 尚未提交时失败,恢复后可能再次发送记录。
- Connector 不提供跨分区、跨 Task 或按设备的全局顺序,也不提供幂等去重或 exactly-once 交付保证。
- Struct 输入要求
messageId、message和deviceId为字符串字段;String 输入必须是可反序列化为消息对象的 JSON,错误的字段或格式会导致记录转换失败。
常见问题
Connector 启动时报 topics 和 topics.regex 配置错误,怎么办?
检查两个配置是否同时为空或同时设置。固定 Topic 场景只保留非空的 topics;按命名规则订阅时只保留非空的 topics.regex。修改后重新提交 Connector 配置,并确认正则表达式是合法的 Java 正则。
消息转换失败,如何检查输入记录?
先查看 Task 日志中的转换错误,再确认值转换器输出的是 String 或 Struct。String 值应为 JSON,并包含字符串类型的messageId、message 和 deviceId;Struct 值应包含这些字段,并使用字符串类型的 expiry 表示可选过期时间。JSON 中的过期字段名是 expiryTime,不要把它与 Struct 字段名 expiry 混用。