概述
Azure IoT Hub Source Connector 从 Azure IoT Hub 的 Event Hub 兼容终结点读取设备遥测数据,并将消息写入 Kafka 主题。Connector 按 IoT Hub 分区分配任务,每条消息转换为固定的结构化记录,包含设备标识、Event Hub offset、入队时间、序列号、原始内容以及系统属性和应用属性。设备发送的内容保留在content 字段中,适合下游按统一元数据格式消费不同设备的遥测数据。
前置条件
- Azure IoT Hub 必须已启用可读取事件的 Event Hub 兼容终结点,并准备好对应的事件端点、兼容名称、分区数、使用的消费者组和共享访问策略。
- 用于连接的共享访问策略必须具有读取 IoT Hub 事件流所需的权限;将主键作为敏感信息管理,不要写入日志或提交到代码仓库。
授权许可
使用 MIT License。快速开始
准备好 Connect Cluster、Kafka,以及可访问的 Azure IoT Hub 后,确认网络连通和访问权限,然后按照 管理 Connector 中的步骤提交配置。IotHub.ConsumerGroup 时使用 $Default;未指定起始时间或起始 offset 时,首次启动从流的开头开始。Connector 启动后会将读取到的记录写入 Kafka.Topic。
配置
Azure IoT Hub 连接
IotHub.EventHubCompatibleName
IoT Hub 事件终结点的 Event Hub 兼容名称。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填。可在 Azure Portal 的 IoT Hub、终结点、事件中获取。
IotHub.EventHubCompatibleEndpoint
IoT Hub 事件终结点的 Event Hub 兼容地址。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填,必须是可解析的终结点 URI。
IotHub.AccessKeyName
用于访问事件流的共享访问策略名称。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填,例如 Azure IoT Hub 中已配置的
service策略名称。
IotHub.AccessKeyValue
所选共享访问策略的主键。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填。该值属于凭据,应通过安全的配置管理方式提供。
Azure IoT Hub 消费
IotHub.ConsumerGroup
读取 IoT Hub 事件流时使用的消费者组。
- 类型:
string - 默认值:
$Default - 重要级别:中
- 有效值 / 注意事项:应使用已在 IoT Hub 事件终结点中创建的消费者组。为独立的消费应用使用单独的消费者组,避免互相移动读取位置。
IotHub.Partitions
IoT Hub 的分区数量。
- 类型:
int - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填,必须与实际 IoT Hub 分区数一致。该值用于生成任务的分区分配。
起始位置
IotHub.StartTime
从指定 UTC 时间开始读取消息。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:可选,使用 ISO-8601 时间戳,例如
2026-09-18T00:00:00Z。设置后优先于IotHub.Offsets作为初始位置;已提交的 Kafka Connect offset 仍优先使用。
IotHub.Offsets
为每个 IoT Hub 分区指定初始 Event Hub offset,多个 offset 使用逗号分隔。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:可选,按分区顺序提供 offset。设置
IotHub.StartTime后该配置会被忽略。空项从对应分区的流开头开始。
轮询
BatchSize
每次从 IoT Hub 请求的消息数量上限。
- 类型:
int - 默认值:
100 - 重要级别:中
- 有效值 / 注意事项:应使用正整数。增大该值可以减少轮询次数,但会增加单次轮询返回的数据量。
ReceiveTimeout
接收消息时等待数据的最长时间,单位为秒。
- 类型:
int - 默认值:
60 - 重要级别:中
- 有效值 / 注意事项:该值会传递给 Event Hub 接收器。根据消息到达频率和延迟要求调整。
Kafka 输出
Kafka.Topic
接收 IoT Hub 消息的 Kafka 主题。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填。所有已分配分区读取到的记录都会写入该主题。
最佳实践
首次接入时明确数据起点
适用业务场景:首次将已有 IoT Hub 遥测接入 Kafka,需要从某个时间点建立数据基线,并继续读取之后到达的新消息。 配置示例:IotHub.StartTime 只决定没有已保存 source offset 时的初始位置。为该 Connector 使用独立消费者组,便于把它的读取进度与其他消费应用隔离。时间值使用 UTC,并根据业务需要选择基线时间;时间越早,首次读取的数据量可能越大。
按 IoT Hub 分区扩展任务
适用业务场景:IoT Hub 有多个分区,单个任务无法满足吞吐或延迟要求,需要增加并行读取能力。 配置示例:IotHub.Partitions 应与实际分区数一致,tasks.max 可以设置为不超过分区数的并行任务数。Connector 以轮询方式将分区分配给任务,一个任务可能负责多个分区;任务数超过可分配分区数不会产生额外的分区读取能力。
监控
监控内容
监控 Kafka Connect Worker、Connector 和 Task 的运行状态,关注任务重启、吞吐量、端到端延迟、source offset 提交、错误与重试,以及 Worker JVM 的堆内存、线程和 GC 指标。若部署启用了错误处理和死信队列,再关注死信队列写入和积压。导入 Grafana 大盘
从 下载地址 获取大盘,使用已采集 Kafka Connect 指标的 Prometheus 数据源,并确保指标标签与大盘变量匹配,然后在 Grafana 中导入 JSON 文件。限制条件
- 一个 Kafka topic 接收该 Connector 所有已分配 IoT Hub 分区的记录,Connector 不会按设备或消息 schema 自动拆分到不同主题。
IotHub.StartTime和IotHub.Offsets只用于没有已保存 source offset 时的初始定位;已提交的 Kafka Connect offset 优先级更高。- Connector 依赖 Azure IoT Hub 的 Event Hub 兼容接口和有效的共享访问策略,不能在没有这些外部资源的环境中独立产生源记录。
常见问题
为什么 Connector 启动后没有读取到消息?
检查 Event Hub 兼容名称、终结点、共享访问策略和消费者组是否属于同一个 IoT Hub,并确认IotHub.Partitions 与实际分区数一致。若配置了 IotHub.StartTime,确认该时间使用 UTC 且仍在可读取的数据保留范围内;若 Kafka Connect 已保存 source offset,Connector 会从已保存位置继续读取。
为什么设置了 IotHub.Offsets 却没有从这些位置开始?
如果同时设置了 IotHub.StartTime,offset 配置会被忽略。此外,已有的 Kafka Connect source offset 会优先于这两个初始位置配置。检查当前 Connector 是否使用了预期的消费者组和 offset 存储,再根据需要清理或重新建立 Connector 的读取状态。
如何选择 tasks.max?
先确认 IoT Hub 分区数,再将 tasks.max 设置为不超过分区数的值。任务会按分区分配工作;超过分区数的任务不会增加实际读取并行度。调整后同时观察 Task 状态、吞吐和延迟。