> ## 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.

# Timeplus Sink Connector

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

## 概述

Timeplus Sink Connector 从 Kafka Topic 读取记录，并通过 Timeplus v1beta2 HTTP 接口将记录值写入指定的 Timeplus Stream。它位于 Kafka 事件流与 Timeplus 实时分析存储之间，适合把日志、业务事件和其他持续产生的文本或 JSON 数据汇入一个目标 Stream。

Connector 将每条 Kafka 记录的值转换为字符串，并按记录顺序组成换行分隔的请求体。`raw` 模式把每个字符串写成一行文本；JSON 模式把字符串发送到流式 JSON 写入接口。Kafka 记录的 Key、Topic、Partition、Offset、时间戳和 Headers 不会自动映射为 Timeplus 列。

## 前置条件

* 准备可访问的 Timeplus Workspace 和具有写入权限的 API Key。默认配置要求预先创建与输入格式兼容的目标 Stream；若启用自动建流，API Key 还须具备建流权限，JSON 模式还需要 Schema 推断权限，且首条记录值须为 JSON 对象。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka、Timeplus Workspace 和目标 Stream，并确认网络连通和访问权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。下面的最小实用配置从一个 Kafka Topic 读取字符串，并以 `raw` 模式写入预先创建的 Timeplus Stream。

```properties theme={null}
connector.class=com.timeplus.kafkaconnect.TimeplusSinkConnector
topics=<input-topic>
value.converter=org.apache.kafka.connect.storage.StringConverter
timeplus.sink.address=<timeplus-address>
timeplus.sink.workspace=<workspace-id>
timeplus.sink.apikey=<timeplus-api-key>
timeplus.sink.stream=<stream-name>
```

将 `<input-topic>`、`<timeplus-address>`、`<workspace-id>`、`<timeplus-api-key>` 和 `<stream-name>` 替换为实际资源。目标 Stream 应预先存在并能接收每条记录值的字符串表示；`raw` 模式通常使用一个名为 `raw` 的 String 列。不要把 API Key 写入日志或提交到版本控制系统。

## 配置

### Timeplus 连接与认证

#### `timeplus.sink.address`

设置用于构造 Timeplus Stream、写入和 Schema 推断接口 URL 的基础地址。

* **类型**：`string`
* **默认值**：`https://dev.timeplus.cloud`
* **重要级别**：高
* **有效值 / 注意事项**：使用每个 Connect Worker 都能访问的 Timeplus 基础 URL。Connector 不校验 Scheme、Host 或连通性，也不会规范化尾部斜杠；对于其他 Timeplus Cloud 区域或自托管部署，应显式设置正确地址。

#### `timeplus.sink.workspace`

设置 Timeplus API URL 中使用的 Workspace 标识。

* **类型**：`string`
* **默认值**：`default`
* **重要级别**：高
* **有效值 / 注意事项**：使用目标地址中已存在的 Workspace ID。该值未经编码直接加入 URL，Connector 不校验是否为空或 Workspace 是否存在；API Key 必须有权访问该 Workspace。

#### `timeplus.sink.apikey`

设置发送到 Timeplus 请求 `X-API-KEY` Header 的 API Key。

* **类型**：`password`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：使用目标 Workspace 接受的 API Key。Connector 不校验长度、格式或非空；需要认证的部署必须显式提供有效密钥。启用自动建流时，密钥还须具备创建 Stream 和执行 Schema 推断的权限。该配置包含敏感信息，请使用安全的配置注入方式。

### 目标 Stream 与数据格式

#### `timeplus.sink.stream`

设置接收记录的 Timeplus Stream 名称。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **必填**：是
* **有效值 / 注意事项**：使用 Timeplus 接受的 Stream 名称。该值未经编码直接加入写入 URL；Connector 不校验空字符串或命名规则。`timeplus.sink.createStream=false` 时，目标 Stream 必须已经存在。

#### `timeplus.sink.dataFormat`

设置写入请求的数据格式及自动建流时的 Schema 生成路径。

* **类型**：`string`
* **默认值**：`raw`
* **重要级别**：高
* **有效值 / 注意事项**：推荐使用区分大小写的 `raw` 或 `json`。只有严格等于 `raw` 时才使用换行文本接口；启用自动建流时，该模式会创建仅含一个 `raw` String 列的 Stream。其他任意值都会进入流式 JSON 分支，但 Connector 不校验这些值是否有效。两种模式都发送转换后值的 `toString()` 结果，JSON 模式要求该字符串符合 Timeplus 流式 JSON 接口及目标 Stream Schema。

#### `timeplus.sink.createStream`

控制每个 Task 是否在首次处理非空记录批次时尝试创建目标 Stream。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：高
* **有效值 / 注意事项**：`false` 表示使用预先创建的 Stream。设为 `true` 时，`raw` 模式尝试创建一个包含 `raw` String 列的 Stream；JSON 模式使用首条记录值执行 Schema 推断后建流。Connector 不先检查 Stream 是否存在，每个 Task 都有独立的首次创建状态，创建失败也可能继续写入且不再由该 Task 重试创建。

### Kafka 输入订阅与任务

#### `connector.class`

选择 Timeplus Sink Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **必填**：是
* **有效值 / 注意事项**：使用 `com.timeplus.kafkaconnect.TimeplusSinkConnector`。

#### `topics`

指定 Connector 要读取的 Kafka Topic 列表。

* **类型**：`list`
* **默认值**：空列表 `[]`
* **重要级别**：高
* **必填**：与 `topics.regex` 二选一
* **有效值 / 注意事项**：使用逗号分隔的 Topic 名称。必须与非空的 `topics.regex` 二选一；两者同时设置或同时为空都会导致 Sink 配置校验失败。所有选中 Topic 的记录都会写入同一个 `timeplus.sink.stream`。

#### `topics.regex`

使用 Java 正则表达式选择 Connector 要读取的 Kafka Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **必填**：与 `topics` 二选一
* **有效值 / 注意事项**：使用非空且合法的 Java 正则表达式。必须与非空的 `topics` 二选一；所有匹配 Topic 的记录都会写入同一个 `timeplus.sink.stream`。

#### `tasks.max`

设置 Connector 可以创建的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须至少为 `1`。实际并行度受输入 Topic 分区数、Worker 数量和分区分配影响。每个 Task 独立发送 HTTP 请求；启用自动建流时，每个 Task 还可能独立尝试创建同一个 Stream。

### 记录值转换

#### `value.converter`

设置 Worker 用于反序列化 Kafka 记录值的 Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未设置时继承 Worker 级配置；显式设置时必须是可实例化的 `org.apache.kafka.connect.storage.Converter` 实现。Connector 对转换后的对象调用 `toString()`，因此应选择能产生预期文本行或有效 JSON 文本的 Converter。推荐使用 `org.apache.kafka.connect.storage.StringConverter` 保持 Kafka Value 的文本表示；Struct、Map、字节数组等对象的 `toString()` 结果不保证符合 JSON 格式。

## 最佳实践

### 将 JSON 事件写入预先创建的 Stream

**适用业务场景**：已经确定事件字段和类型，希望把 Kafka 中的 JSON 对象持续写入具有对应 Schema 的 Timeplus Stream，并由 Timeplus 按字段进行实时查询和分析。

**配置示例**：

```properties theme={null}
connector.class=com.timeplus.kafkaconnect.TimeplusSinkConnector
topics=<input-topic>
value.converter=org.apache.kafka.connect.storage.StringConverter
timeplus.sink.address=<timeplus-address>
timeplus.sink.workspace=<workspace-id>
timeplus.sink.apikey=<timeplus-api-key>
timeplus.sink.stream=<stream-name>
timeplus.sink.dataFormat=json
```

将占位符替换为实际资源，并预先创建与事件字段兼容的目标 Stream。输入 Topic 中每条记录的 Value 应是单个有效 JSON 对象的字符串；不要使用会把值转换为非 JSON `toString()` 表示的 Converter。

**关键说明**：预先创建 Stream 可以在接入前明确字段名称和类型，避免由首条事件决定推断结果。Connector 不进行 JSON 规范化或后续 Schema 演进；新增字段、类型变化和嵌套结构必须先确认与目标 Stream 兼容。

### 按命名规则扩展输入 Topic

**适用业务场景**：多个业务 Topic 使用统一命名规则，后续还会增加同类 Topic，希望 Connector 自动订阅匹配范围并将记录汇入同一个 Timeplus Stream。

**配置示例**：

```properties theme={null}
connector.class=com.timeplus.kafkaconnect.TimeplusSinkConnector
topics.regex=events\\..*
value.converter=org.apache.kafka.connect.storage.StringConverter
timeplus.sink.address=<timeplus-address>
timeplus.sink.workspace=<workspace-id>
timeplus.sink.apikey=<timeplus-api-key>
timeplus.sink.stream=<stream-name>
```

使用 `topics.regex` 后移除 `topics`，并将其他占位符替换为实际资源。确保所有匹配 Topic 的记录值都能由同一个目标 Stream 和数据格式接收。

**关键说明**：在 properties 文件中，`events\..*` 的正则转义反斜杠还需要按 Java Properties 规则转义，因此示例写为 `topics.regex=events\\..*`；解析后传给 Kafka Connect 的实际正则为 `events\..*`，只匹配以 `events.` 开头的 Topic。正则订阅简化了同类 Topic 的范围维护，但 Connector 不按 Topic 自动路由到不同 Stream，也不会把 Topic 名称写入目标列。修改表达式前应检查匹配范围，避免无关 Topic 被写入同一 Stream。

### 积压增长时增加 Task 并行度

**适用业务场景**：Connector 已稳定运行，输入 Topic 有多个分区且消费 Lag 持续增长，希望利用多个 Task 并行发送 HTTP 写入请求。

**配置示例**：

```properties theme={null}
connector.class=com.timeplus.kafkaconnect.TimeplusSinkConnector
topics=<input-topic>
tasks.max=3
value.converter=org.apache.kafka.connect.storage.StringConverter
timeplus.sink.address=<timeplus-address>
timeplus.sink.workspace=<workspace-id>
timeplus.sink.apikey=<timeplus-api-key>
timeplus.sink.stream=<stream-name>
```

将占位符替换为实际资源，先确认输入 Topic 至少有多个可分配分区，并使用预先创建的目标 Stream。逐步提高 `tasks.max`，同时观察消费 Lag、HTTP 延迟、Timeplus 限流和 Worker 资源使用情况。

**关键说明**：`tasks.max` 是上限，实际 Task 数量取决于分区分配。不同 Task 的 HTTP 请求可并发完成，Connector 不提供跨 Task 全局顺序或去重；增加并行度也不会拆分单次过大的请求批次。

## 监控

### 监控内容

监控 Kafka Connect Worker、Connector 和 Task 的健康状态及状态变化，关注输入吞吐、消费 Lag、处理延迟、Offset 提交、错误、重试和 Worker JVM 的 CPU、内存及垃圾回收信号；同时核对 Timeplus 实际写入量和 Worker 日志，因为部分 HTTP 写入失败可能只记录日志而不会使 Task 失败。仅在部署启用了相应 Kafka Connect 错误处理时关注 DLQ 活动。

### 导入 Grafana 大盘

确认 Kafka Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* Connector 对每条记录的 Value 调用 `toString()` 并追加换行，不读取 Kafka Key、Topic、Partition、Offset、时间戳、Headers 或 Connect Schema 来生成目标列；Struct、Map、字节数组等对象的字符串表示不保证为有效 JSON。
* 写入请求最终返回非成功状态或发生 `IOException` 时，Connector 可能只记录日志而不向 Kafka Connect 传播失败；Task 显示运行中或 Offset 已提交都不能单独证明对应数据已写入 Timeplus，因此不能依赖该版本获得至少一次交付保证。
* Timeplus 写入和 Kafka Offset 提交不是原子操作；写入成功后、Offset 提交前发生故障可能导致整批记录重放，而被 Connector 吞掉的写入失败会让 `put` 正常返回，Worker 随后可能推进并提交 Offset，形成数据缺口。Connector 不提供事务、幂等键、自动去重或 exactly-once 交付。
* 自动建流不检查 Stream 是否已经存在，且每个 Task 在启动或重启后都可能独立尝试创建；创建请求失败后，Task 仍可能继续写入并把本地状态标记为已尝试。JSON 自动建流只根据首条记录推断 Schema，不处理后续 Schema 演进。
* 一个 Connector 配置中的所有输入 Topic 和 Partition 都写入同一个目标 Stream；Connector 不提供按 Topic 或 Partition 路由多个 Stream 的配置，也不保证跨 Task 的全局顺序。
* Connector 将每次非空 `put` 收到的整批记录组成一个内存请求体，不提供批次大小、字节上限、拆包、异步队列或客户端超时配置；单条或整批数据超过 Timeplus 或代理限制时不会自动缩小批次。

## 常见问题

### Task 显示运行中，为什么 Timeplus 中没有新数据？

先检查 Worker 日志中的 Timeplus HTTP 失败状态或连接异常消息，再核对 `timeplus.sink.address`、`timeplus.sink.workspace`、`timeplus.sink.apikey`、`timeplus.sink.stream` 及目标 Stream Schema。该版本可能在写入失败后仍让 `put` 正常返回，因此还应直接比较 Kafka 输入量、已提交 Offset 和 Timeplus 实际行数；修正目标地址、权限、Stream 或数据格式后，再用新的测试记录确认写入恢复。

### JSON 写入被拒绝，如何检查记录值？

确认 `timeplus.sink.dataFormat=json`，并检查 `value.converter` 输出对象的 `toString()` 是否为单个有效 JSON 对象。使用 `StringConverter` 时，Kafka Value 本身应是 JSON 文本；使用 Struct、Map 或其他结构化对象时，不要假设其默认字符串表示符合 JSON。还要确认字段名称、类型和嵌套结构与目标 Stream Schema 兼容。

### 启用自动建流后，为什么 Stream 仍未正确创建？

确认 API Key 具有建流和 Schema 推断权限。JSON 模式下，首条记录值必须能解析为 JSON 对象；多个 Task 可能同时尝试创建同一个 Stream，已有 Stream 也不会被预先识别。由于创建失败可能只写日志且该 Task 不再重试，生产环境应优先预先创建并校验 Stream，然后保持 `timeplus.sink.createStream=false`。

### 为什么重启后可能出现重复数据或数据缺口？

Timeplus HTTP 写入和 Kafka Offset 提交分属两个步骤。若 Timeplus 已接收批次但 Offset 尚未提交，重启后可能重放该批；若写入失败被 Connector 记录后吞掉，`put` 仍正常返回，Worker 随后可能推进并提交 Offset，则该批可能不会自动重放。检查 Worker 日志、Offset 和 Timeplus 行数，并在业务侧使用稳定事件标识进行审计或去重；不要把该版本视为至少一次或恰好一次交付实现。

### Connector 启动时报 `topics` 和 `topics.regex` 配置错误，怎么办？

检查两项是否同时为空或同时设置。固定 Topic 场景只保留非空的 `topics`；按命名规则订阅时只保留非空且合法的 `topics.regex`。修改后重新提交配置，并确认所有选中的 Topic 数据都适合写入同一个 Timeplus Stream。
