> ## 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 Sink Connector

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

## 概述

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](../manage-connectors)。下面的配置读取一个 Topic 中的 JSON 字符串，并将每条消息发送到其 `deviceId` 指定的设备。

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.sink.IotHubSinkConnector
topics=<input-topic>
value.converter=org.apache.kafka.connect.storage.StringConverter
IotHub.ConnectionString=<iot-hub-connection-string>
```

将 `<input-topic>` 替换为输入 Topic，将 `<iot-hub-connection-string>` 替换为 Azure IoT Hub 连接字符串。不要把连接字符串写入日志或提交到版本控制系统。使用该示例时，Topic 中每条消息的字符串值应为可解析的 JSON，例如：

```json theme={null}
{"messageId":"msg-1001","message":"reboot","deviceId":"device-001"}
```

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 列表。

**配置示例**：

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.sink.IotHubSinkConnector
topics.regex=devices\..*\.commands
value.converter=org.apache.kafka.connect.storage.StringConverter
IotHub.ConnectionString=<iot-hub-connection-string>
```

将 `<iot-hub-connection-string>` 替换为 Azure IoT Hub 连接字符串。使用 `topics.regex` 后移除 `topics`；不要同时配置两个订阅方式。表达式匹配的每个 Topic 中的字符串值都应遵循相同的 JSON 消息结构。

**关键说明**：使用正则订阅便于统一管理同类 Topic，但 Topic 命名规则会直接决定消费范围。变更命名规则前先确认新旧 Topic 是否会同时匹配，避免意外扩大或缩小输入范围。

### 按分区增加 Task 并行度

**适用业务场景**：Connector 已能稳定发送消息，输入 Topic 有多个分区且单个 Task 的逐条同步发送吞吐不足，希望通过增加 Task 并行处理分区。

**配置示例**：

```properties theme={null}
connector.class=com.microsoft.azure.iot.kafka.connect.sink.IotHubSinkConnector
topics=<input-topic>
tasks.max=2
value.converter=org.apache.kafka.connect.storage.StringConverter
IotHub.ConnectionString=<iot-hub-connection-string>
```

将 `<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 大盘](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)获取 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` 混用。

### 发送调用成功后，为什么 Kafka 记录可能再次发送？

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

### 增加 `tasks.max` 后处理顺序发生变化，正常吗？

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