概述
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。下面的最小实用配置从一个 Kafka Topic 读取字符串,并以raw 模式写入预先创建的 Timeplus Stream。
<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时才使用换行文本接口;启用自动建流时,该模式会创建仅含一个rawString 列的 Stream。其他任意值都会进入流式 JSON 分支,但 Connector 不校验这些值是否有效。两种模式都发送转换后值的toString()结果,JSON 模式要求该字符串符合 Timeplus 流式 JSON 接口及目标 Stream Schema。
timeplus.sink.createStream
控制每个 Task 是否在首次处理非空记录批次时尝试创建目标 Stream。
- 类型:
boolean - 默认值:
false - 重要级别:高
- 有效值 / 注意事项:
false表示使用预先创建的 Stream。设为true时,raw模式尝试创建一个包含rawString 列的 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 按字段进行实时查询和分析。 配置示例:toString() 表示的 Converter。
关键说明:预先创建 Stream 可以在接入前明确字段名称和类型,避免由首条事件决定推断结果。Connector 不进行 JSON 规范化或后续 Schema 演进;新增字段、类型变化和嵌套结构必须先确认与目标 Stream 兼容。
按命名规则扩展输入 Topic
适用业务场景:多个业务 Topic 使用统一命名规则,后续还会增加同类 Topic,希望 Connector 自动订阅匹配范围并将记录汇入同一个 Timeplus Stream。 配置示例: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 写入请求。 配置示例: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,在 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。