Skip to main content

概述

Lightstreamer Sink Connector 消费 Kafka Topic 中的记录,将其转换为 Lightstreamer 字段更新,并通过 Lightstreamer Proxy Adapter 发送给订阅实时 Item 的 Web、移动端或其他客户端。它位于 Kafka Connect 与 Lightstreamer Server 之间:Kafka Connect 负责读取记录和管理 Offset,Connector 负责 Topic 到 Item 的路由以及记录到字段的映射,Lightstreamer Server 负责向客户端分发更新。 一个 Topic 可以映射到一个或多个静态 Item,也可以通过 Item 模板按记录 Key、Value、Header 或元数据把更新路由到参数化 Item。字段映射可提取 Key、Value、Header、Topic、分区、Offset 和时间戳,并将结果转换为字符串或 null。Connector 不提供初始快照;在普通主动连接模式下,只有存在匹配的客户端订阅时,记录才会产生 Item 更新。 该 Connector 适合把行情、设备状态、业务事件和实时看板数据从 Kafka 推送到 Lightstreamer 客户端。Lightstreamer Item 更新与 Kafka Offset 提交不是跨系统事务,故障恢复可能产生重复更新;业务需要在事件中保留稳定标识,并按需要在客户端或业务层处理重复。

前置条件

  • 部署 Lightstreamer Server 7.4.2 或更高版本,并在目标 Adapter Set 中配置可与 Connector 通信的 Proxy Data Adapter;其请求/响应端口、连接方向和可选认证信息需与 Connector 配置一致。
  • Lightstreamer 客户端需连接到正确的 Adapter Set 和 Data Adapter,并订阅 topic.mappingsitem.templates 定义的 Item 以及 record.mappings 定义的字段;没有匹配订阅时不会产生客户端更新。
  • 使用 KEYVALUE 路径提取时,所选 Kafka Connect Converter 必须生成与记录结构一致的非空 Connect Schema;使用 DLQ 时需准备 Connector 可写且不在输入订阅范围内的 Kafka Topic。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka Topic、Lightstreamer Server、Proxy Adapter 和订阅客户端,并确认网络连通与访问权限。集群准备和 Connector 管理操作参见管理 Connector。以下示例消费字符串消息,将每条记录的 Value 作为 message 字段发送到静态 Item events
将两处 <kafka-topic> 替换为同一个输入 Topic,并替换 Lightstreamer Proxy Adapter 地址。客户端应订阅对应 Adapter Set 中的 events Item 和 message 字段。若 Proxy Adapter 启用了认证,还需配置用户名和密码;凭证应通过部署环境的安全配置机制注入,不要写入版本库或共享日志。

配置

Connector 身份、输入与任务

connector.class

指定 Lightstreamer Sink Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 com.lightstreamer.kafka.connect.LightstreamerSinkConnector;所有运行 Task 的 Worker 均需能够加载该插件。
  • 必填:是

topics

指定需要消费的 Kafka Topic。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:多个名称用逗号分隔;与 topics.regex 二选一且必须配置其中一项。输入 Topic 还需匹配 topic.mappings 中的字面量或正则规则。启用 DLQ 时不能订阅 DLQ Topic。

topics.regex

通过正则表达式选择需要消费的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:使用 Java 正则语法;与 topics 互斥。配置 DLQ 时,该正则不能匹配 DLQ Topic。此项只决定 Kafka Connect 的输入订阅,不会自动启用 topic.mappings 的正则解释。

tasks.max

设置 Kafka Connect 允许创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1。该 Connector 始终只返回一个 Task 配置,因此提高此值不会增加单个 Connector 实例的任务并行度。

Proxy Adapter 连接与认证

lightstreamer.server.proxy_adapter.address

设置 Lightstreamer Proxy Adapter 的主机和请求/响应端口。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 host:port 格式,主机不能为空且端口必须为不带前导零的正整数。此项始终必填;启用连接反转时仍需提供,但其主机和端口不用于主动出站连接。方括号形式的 IPv6 地址不符合该配置的校验格式。
  • 必填:是

lightstreamer.server.proxy_adapter.socket.connection.setup.timeout.ms

设置主动连接模式下建立 Proxy Adapter Socket 的等待时间,单位毫秒。
  • 类型int
  • 默认值5000
  • 重要级别:低
  • 有效值 / 注意事项:必须大于或等于 00 表示不设置连接超时。仅在 connection.inversion.enable=false 时生效。

lightstreamer.server.proxy_adapter.socket.connection.setup.max.retries

设置主动连接模式下初始连接失败后的最大重试次数。
  • 类型int
  • 默认值1
  • 重要级别:中
  • 有效值 / 注意事项:必须大于或等于 00 表示不重试。此项只控制 Task 启动阶段的连接建立,不是记录级重试,也不会在运行期连接断开后持续重连。

lightstreamer.server.proxy_adapter.socket.connection.setup.retry.delay.ms

设置主动连接模式下两次初始连接尝试之间的等待时间,单位毫秒。
  • 类型long
  • 默认值5000
  • 重要级别:低
  • 有效值 / 注意事项:必须大于或等于 0;仅在 connection.inversion.enable=false 且最大重试次数大于 0 时生效。

lightstreamer.server.proxy_adapter.username

设置 Proxy Adapter 远程连接认证用户名。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 Proxy Adapter 启用了远程 Adapter 认证时配置;空字符串与未配置不同。主动连接和连接反转模式均可使用。

lightstreamer.server.proxy_adapter.password

设置 Proxy Adapter 远程连接认证密码。
  • 类型password
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:通常与 lightstreamer.server.proxy_adapter.username 一起配置。使用安全配置机制保存,不要在日志、文档或版本库中暴露真实值。

connection.inversion.enable

设置是否反转 Proxy Adapter 与 Connector 之间的连接建立方向。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项false 时 Connector 主动连接 lightstreamer.server.proxy_adapter.addresstrue 时 Connector 在 request_reply.port 监听,由 Proxy Adapter 主动连接。反转模式需同步配置 Proxy Adapter 的远程主机。

request_reply.port

设置连接反转模式下 Connector 监听的请求/响应端口。
  • 类型int
  • 默认值6661
  • 重要级别:低
  • 有效值 / 注意事项:仅在 connection.inversion.enable=true 时生效。配置校验只要求大于或等于 0,实际使用时应选择操作系统可绑定且位于有效 TCP 端口范围内的端口。

max.proxy.adapter.connections

设置连接反转模式下允许同时接入的 Proxy Adapter 连接数。
  • 类型int
  • 默认值1
  • 重要级别:低
  • 有效值 / 注意事项:必须大于或等于 1,仅在 connection.inversion.enable=true 时生效。同一批记录会发送给当时所有活动连接,但多个连接之间不提供事务性或顺序保证。

Topic 与 Item 路由

topic.mappings

定义 Kafka Topic 到 Lightstreamer 静态 Item 或 Item 模板的映射。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 topic:item1,item2;other-topic:item3 格式。每个 Topic 或正则键必须唯一,映射不能为空;同一映射内的重复 Item 会被去重。默认按 Topic 字面量匹配;启用 topic.mappings.regex.enable 后按 Java 正则完整匹配。item-template.<name> 引用必须在 item.templates 中定义。
  • 必填:是

topic.mappings.regex.enable

设置是否把 topic.mappings 中每个 Topic 键解释为 Java 正则表达式。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用后使用完整 Topic 名匹配;无效正则会在 Task 构建映射器时导致失败。此项与 topics.regex 相互独立,两处规则都需覆盖预期输入 Topic。

item.templates

定义根据记录内容筛选订阅的参数化 Lightstreamer Item 模板。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:使用 template-name:item-prefix-#{param=EXPRESSION};other-template:... 格式,模板名必须唯一且值不能为空。#{...} 是 Connector 端的模板定义;客户端需使用 item-prefix-[param=literalValue] 格式提供字面筛选值。模板参数只能使用标量值。topic.mappings 通过 item-template.<template-name> 引用模板。

记录转换与字段映射

key.converter

设置 Kafka 记录 Key 的 Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未指定时继承 Worker 设置。使用 KEY 提取表达式或基于 Key 的 Item 模板时,Converter 必须与消息编码一致并生成非空 Connect Schema;简单字符串 Key 可使用 org.apache.kafka.connect.storage.StringConverter

value.converter

设置 Kafka 记录 Value 的 Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未指定时继承 Worker 设置。使用 VALUE 提取表达式时,Converter 必须与消息编码一致并生成非空 Connect Schema;字符串消息可使用 org.apache.kafka.connect.storage.StringConverter。Converter 自身选项需按所选实现配置。

record.mappings

定义 Lightstreamer 字段名及其对应的 Kafka 记录提取表达式。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔的 field:#{EXPRESSION} 条目,字段名必须非空且唯一。表达式可从 KEYVALUEHEADERSTOPICPARTITIONOFFSETTIMESTAMP 提取数据。字段集合由此配置固定,输入记录中的其他字段不会自动发布;需要在表达式内使用逗号时还需遵守 properties 转义规则。
  • 必填:是

record.mappings.skip.failed.enable

设置单个字段提取失败时是否省略该字段并继续发送其余字段。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项false 时字段提取错误交给 record.extraction.error.strategy 处理;true 时只省略失败字段,其他字段仍可更新。此项不跳过 Item 模板参数提取失败,且全部字段失败时可能发送空字段 Map。

record.mappings.map.non.scalar.values.enable

设置是否允许把 Struct、Map、数组等非标量选择结果转换为字段文本。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项false 时字段映射必须选择标量值;true 时 Struct 转为不含 Schema 的 JSON 文本,其他复杂 Java 值按其字符串表示发送。此项只影响 record.mappings,不会放宽 Item 模板参数的标量要求。

提取错误与 DLQ

record.extraction.error.strategy

设置记录映射阶段发生提取错误时的处理策略。
  • 类型string
  • 默认值IGNORE_AND_CONTINUE
  • 重要级别:中
  • 有效值 / 注意事项:区分大小写,可选 IGNORE_AND_CONTINUEFORWARD_TO_DLQTERMINATE_TASK。忽略和成功写入 DLQ 都会推进该记录的处理位置;使用 FORWARD_TO_DLQ 隔离坏记录并继续处理后续记录时,必须同时配置非空的 errors.deadletterqueue.topic.nameerrors.tolerance=all。缺少 DLQ Topic 或保留默认容忍级别 none 都会使提取错误终止 Task。该策略只处理映射阶段的提取错误,不处理任意连接或通信异常。

errors.deadletterqueue.topic.name

设置用于接收提取失败记录的 Kafka DLQ Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:空字符串表示不启用 DLQ。使用 record.extraction.error.strategy=FORWARD_TO_DLQ 时必须配置非空值,并配合 errors.tolerance=all 才能在报告坏记录后继续处理后续记录;DLQ Topic 不能出现在 topics 中,也不能被 topics.regex 匹配。

errors.tolerance

设置 Kafka Connect 是否容忍记录处理失败并继续运行 Task。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 noneallnone 在首条失败记录超过容忍限制时终止 Task;all 允许容忍失败记录。对于 record.extraction.error.strategy=FORWARD_TO_DLQ,必须设置为 all,才能在将提取失败记录报告到已配置的 DLQ 后继续处理后续记录。此设置不会把 Proxy Adapter 连接、通信或其他任意错误转换为可容忍的 DLQ 记录。

errors.deadletterqueue.topic.replication.factor

设置 Kafka Connect 自动创建 DLQ Topic 时使用的副本因子。
  • 类型short
  • 默认值3
  • 重要级别:中
  • 有效值 / 注意事项:仅在配置的 DLQ Topic 不存在且由 Kafka Connect 创建时使用;该值必须适合目标 Kafka 集群的 Broker 数量和副本策略。

errors.deadletterqueue.context.headers.enable

设置是否在 DLQ 记录中附加 Kafka Connect 错误上下文 Header。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在实际产生 DLQ 记录时生效;启用后上下文 Header 使用 __connect.errors. 前缀,消费 DLQ 时应避免把这些内部诊断 Header 当作原始业务字段。

最佳实践

按 Topic 命名规则扩展输入范围

适用业务场景:首次接入已完成,后续会持续新增一组命名遵循统一规则的 Kafka Topic,希望新 Topic 自动进入同一个实时事件流,而不必反复修改固定 Topic 列表。 配置示例
关键说明:相对快速开始,此配置移除 topics,改用 topics.regex,并同时把 topic.mappings 的键切换为正则。两套规则相互独立,应保持范围一致;否则 Kafka Connect 可能已消费某个 Topic,但该记录没有匹配的 Item 路由。新增 Topic 前还应确认其消息编码和字段结构与现有映射兼容。

按业务实体细分客户端订阅

适用业务场景:客户端只需要某个订单、设备或账户的更新,不希望所有订阅者都接收同一 Topic 的全部记录。此时可用记录 Key 构造参数化 Item,让客户端按业务实体订阅。 配置示例
关键说明:相对快速开始,此配置增加 item.templates 和 Key Converter,并把静态 Item 改为模板引用。Connector 模板定义为 entity-#{id=KEY};客户端需使用字面筛选值订阅 entity-[id=42],Key 为 42 的记录才会匹配该订阅,不能使用 entity-42。从记录中提取的模板参数值及其 Connect Schema 必须与订阅参数匹配。业务实体 Key 应稳定且非空,避免因 Key 变化把同一实体的更新路由到不同 Item。

将字段提取失败记录隔离到 DLQ

适用业务场景:Connector 已持续运行,输入 Schema 演进或个别异常记录可能导致字段路径无法提取,希望保留问题记录供排查,同时继续处理后续有效记录。 配置示例
关键说明:相对快速开始,此配置使用带 Schema 的 JSON Value、嵌套字段路径和 FORWARD_TO_DLQ,并增加 errors.tolerance=all 与 DLQ Topic。all 是隔离坏记录并继续处理后续记录的必要条件;保留默认值 none 时,错误报告后仍会因超过容忍限制而终止 Task。<dlq-topic> 必须与输入 Topic 不同且不在订阅范围内。成功报告到 DLQ 后该记录的 Offset 会继续推进,因此应监控并消费 DLQ,修正数据或映射后再按业务规则重放;此组合只容忍映射阶段的提取失败,不捕获 Proxy Adapter 通信异常等任意错误。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 单个 Connector 实例始终只创建一个 Sink Task;提高 tasks.max 不会增加任务并行度,需要通过多个 Connector 实例划分输入 Topic 才能进行实例级扩展。
  • Connector 不提供初始快照,字段更新由当前客户端订阅驱动;普通主动连接模式下无订阅时不会持续推进 Connector 维护的安全 Offset,运行期间新增订阅不会自动收到此前在无匹配订阅时已消费的记录,重启后则可能从最后已提交位置重放。
  • 连接反转模式在没有活动 Proxy Adapter 或没有订阅时仍可能提交 Worker Offset,因此这些记录可能不会产生客户端更新,也不会依靠 Kafka Offset 自动补发。
  • 目标 Item 更新与 Kafka Offset 提交不是原子事务;故障恢复可能产生重复更新,Connector 不提供 exactly-once、目标侧去重或终端客户端接收确认。
  • KEYVALUE 路径提取要求 Converter 产生非空 Connect Schema;不能把无 Schema JSON Map 的字段导航视为已支持用法。
  • record.extraction.error.strategy 只处理映射阶段的提取错误。运行期连接中断和其他未捕获异常可能使 Task 失败,Connector 没有记录级重试、异步缓冲或主动连接的持续重连机制。

常见问题

为什么 Connector 和 Task 正常运行,但客户端收不到更新?

先确认客户端连接到正确的 Adapter Set 和 Data Adapter,并已订阅 topic.mappings 指向的静态 Item 或符合 item.templates 的参数化 Item。然后检查 Kafka Connect 的 topicstopics.regex 是否覆盖输入 Topic、topic.mappings 是否按正确的字面量或正则方式匹配,以及客户端请求的字段是否存在于 record.mappings。普通主动连接模式没有匹配订阅时不会产生更新;建立订阅后发送一条新记录进行验证,不要假设运行中新增订阅会自动回放此前记录。

为什么 Task 在读取某条记录后失败或记录没有字段更新?

检查 key.convertervalue.converter 与实际消息编码是否一致,并确认 KEYVALUE 路径对应非空 Schema 中的真实字段。字段缺失、数组越界、对标量继续取子字段或选择非标量值都可能导致提取错误。希望单个字段失败时保留其他字段,可评估 record.mappings.skip.failed.enable=true;希望隔离整条问题记录并继续处理后续记录,可配置 FORWARD_TO_DLQ、非空 DLQ Topic 和 errors.tolerance=all;需要立即阻止继续处理时使用 TERMINATE_TASK。该容忍设置只作用于符合错误处理路径的记录失败,不会忽略任意连接或通信错误。

为什么增加 tasks.max 后仍然只有一个 Task?

这是该 Connector 的任务模型。每个 Connector 实例与一个 Remote Adapter Task 对应,Connector 始终只生成一个 Task 配置。需要隔离业务流或扩展部署时,为不同 Topic 集合创建多个 Connector 实例,并确保各实例的输入订阅和 Lightstreamer 路由边界清晰,避免无意重复消费和重复更新。

为什么无法连接 Lightstreamer Proxy Adapter?

检查 lightstreamer.server.proxy_adapter.address 是否使用可解析的 host:port,端口是否与 Proxy Adapter 的请求/响应端口一致,防火墙是否允许连接,以及认证用户名和密码是否与 Proxy Adapter 配置匹配。如果启用了 connection.inversion.enable,还需确认 Connector 的 request_reply.port 可被绑定,并在 Proxy Adapter 侧配置指向 Connector 的远程主机。主动连接的重试参数只作用于 Task 启动阶段;运行期连接断开后,应检查 Task 状态和日志并按部署策略重启。