Skip to main content

概述

Azure IoT Hub Sink Connector 从 Kafka Topic 读取云到设备消息,并将每条记录转换为一条 Azure IoT Hub 消息发送到记录指定的设备。它位于 Kafka 事件流和 IoT Hub 设备消息通道之间,适合将业务系统产生的设备指令、通知或配置更新转发到 IoT Hub。 记录值可以使用 Struct 或 JSON 字符串表示。消息至少包含 messageIdmessagedeviceId;目标设备由 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,例如:
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
  • 重要级别:高
  • 有效值 / 注意事项:只能使用区分大小写的 NoneFullPositiveOnlyNegativeOnly。该值作为每条 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 值应为包含 messageIdmessagedeviceId 和可选 expiryTime 的 JSON;Struct 值须包含同名的 messageIdmessagedeviceId 字符串字段,并可包含字符串类型的 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 级批量请求、异步发送、限速或可配置的发送中消息数量。
  • topicstopics.regex 必须且只能配置一个非空值,不能同时使用或同时省略。
  • 消息发送成功与 Kafka Offset 提交不是同一个原子操作;进程在发送完成但 Offset 尚未提交时失败,恢复后可能再次发送记录。
  • Connector 不提供跨分区、跨 Task 或按设备的全局顺序,也不提供幂等去重或 exactly-once 交付保证。
  • Struct 输入要求 messageIdmessagedeviceId 为字符串字段;String 输入必须是可反序列化为消息对象的 JSON,错误的字段或格式会导致记录转换失败。

常见问题

Connector 启动时报 topicstopics.regex 配置错误,怎么办?

检查两个配置是否同时为空或同时设置。固定 Topic 场景只保留非空的 topics;按命名规则订阅时只保留非空的 topics.regex。修改后重新提交 Connector 配置,并确认正则表达式是合法的 Java 正则。

消息转换失败,如何检查输入记录?

先查看 Task 日志中的转换错误,再确认值转换器输出的是 String 或 Struct。String 值应为 JSON,并包含字符串类型的 messageIdmessagedeviceId;Struct 值应包含这些字段,并使用字符串类型的 expiry 表示可选过期时间。JSON 中的过期字段名是 expiryTime,不要把它与 Struct 字段名 expiry 混用。

发送调用成功后,为什么 Kafka 记录可能再次发送?

Connector 的发送调用完成和 Kafka Offset 提交分属两个步骤。若进程在 Azure 服务接受发送请求后、Offset 提交前失败,恢复后可能重新读取并发送该记录。检查 Task 状态、Offset 提交情况和消费 Lag,并根据消息的业务幂等设计处理潜在重复;不要把确认模式解释为跨系统 exactly-once 保证。

增加 tasks.max 后处理顺序发生变化,正常吗?

正常。Task 数量增加后,Kafka 可能将不同分区分配给不同 Task;Connector 只保持单个 Task 在一次处理调用中收到的记录迭代顺序,不保证跨分区、跨 Task 或按设备的全局顺序。需要顺序时,应让相关消息进入同一分区并控制并行度。