> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Azure IoT Hub Source Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Azure IoT Hub Source Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors) 中的步骤提交配置。

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.source.IotHubSourceConnector
tasks.max=1
Kafka.Topic=<kafka-topic>
IotHub.EventHubCompatibleName=<eventhub-compatible-name>
IotHub.EventHubCompatibleEndpoint=<eventhub-compatible-endpoint>
IotHub.AccessKeyName=<shared-access-policy-name>
IotHub.AccessKeyValue=<shared-access-policy-primary-key>
IotHub.Partitions=<iot-hub-partition-count>
```

将尖括号中的值替换为实际资源。未指定 `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，需要从某个时间点建立数据基线，并继续读取之后到达的新消息。

**配置示例**：

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.source.IotHubSourceConnector
tasks.max=1
Kafka.Topic=<kafka-topic>
IotHub.EventHubCompatibleName=<eventhub-compatible-name>
IotHub.EventHubCompatibleEndpoint=<eventhub-compatible-endpoint>
IotHub.AccessKeyName=<shared-access-policy-name>
IotHub.AccessKeyValue=<shared-access-policy-primary-key>
IotHub.ConsumerGroup=<dedicated-consumer-group>
IotHub.Partitions=<iot-hub-partition-count>
IotHub.StartTime=2026-09-18T00:00:00Z
```

**关键说明**：`IotHub.StartTime` 只决定没有已保存 source offset 时的初始位置。为该 Connector 使用独立消费者组，便于把它的读取进度与其他消费应用隔离。时间值使用 UTC，并根据业务需要选择基线时间；时间越早，首次读取的数据量可能越大。

### 按 IoT Hub 分区扩展任务

**适用业务场景**：IoT Hub 有多个分区，单个任务无法满足吞吐或延迟要求，需要增加并行读取能力。

**配置示例**：

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.source.IotHubSourceConnector
tasks.max=4
Kafka.Topic=<kafka-topic>
IotHub.EventHubCompatibleName=<eventhub-compatible-name>
IotHub.EventHubCompatibleEndpoint=<eventhub-compatible-endpoint>
IotHub.AccessKeyName=<shared-access-policy-name>
IotHub.AccessKeyValue=<shared-access-policy-primary-key>
IotHub.ConsumerGroup=<dedicated-consumer-group>
IotHub.Partitions=4
BatchSize=100
ReceiveTimeout=60
```

**关键说明**：`IotHub.Partitions` 应与实际分区数一致，`tasks.max` 可以设置为不超过分区数的并行任务数。Connector 以轮询方式将分区分配给任务，一个任务可能负责多个分区；任务数超过可分配分区数不会产生额外的分区读取能力。

## 监控

### 监控内容

监控 Kafka Connect Worker、Connector 和 Task 的运行状态，关注任务重启、吞吐量、端到端延迟、source offset 提交、错误与重试，以及 Worker JVM 的堆内存、线程和 GC 指标。若部署启用了错误处理和死信队列，再关注死信队列写入和积压。

### 导入 Grafana 大盘

从 [下载地址](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json) 获取大盘，使用已采集 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 状态、吞吐和延迟。
