Skip to main content

概述

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
  • 重要级别:高
  • 有效值 / 注意事项:推荐使用区分大小写的 rawjson。只有严格等于 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 按字段进行实时查询和分析。 配置示例
将占位符替换为实际资源,并预先创建与事件字段兼容的目标 Stream。输入 Topic 中每条记录的 Value 应是单个有效 JSON 对象的字符串;不要使用会把值转换为非 JSON 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 写入请求。 配置示例
将占位符替换为实际资源,先确认输入 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,在 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.addresstimeplus.sink.workspacetimeplus.sink.apikeytimeplus.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 启动时报 topicstopics.regex 配置错误,怎么办?

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