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

# Lightstreamer Sink Connector

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

## 概述

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.mappings` 或 `item.templates` 定义的 Item 以及 `record.mappings` 定义的字段；没有匹配订阅时不会产生客户端更新。
* 使用 `KEY` 或 `VALUE` 路径提取时，所选 Kafka Connect Converter 必须生成与记录结构一致的非空 Connect Schema；使用 DLQ 时需准备 Connector 可写且不在输入订阅范围内的 Kafka Topic。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka Topic、Lightstreamer Server、Proxy Adapter 和订阅客户端，并确认网络连通与访问权限。集群准备和 Connector 管理操作参见[管理 Connector](../manage-connectors)。以下示例消费字符串消息，将每条记录的 Value 作为 `message` 字段发送到静态 Item `events`。

```properties theme={null}
connector.class=com.lightstreamer.kafka.connect.LightstreamerSinkConnector
topics=<kafka-topic>
lightstreamer.server.proxy_adapter.address=<lightstreamer-host>:6661
topic.mappings=<kafka-topic>:events
record.mappings=message:#{VALUE}
value.converter=org.apache.kafka.connect.storage.StringConverter
```

将两处 `<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`
* **重要级别**：低
* **有效值 / 注意事项**：必须大于或等于 `0`；`0` 表示不设置连接超时。仅在 `connection.inversion.enable=false` 时生效。

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

设置主动连接模式下初始连接失败后的最大重试次数。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：中
* **有效值 / 注意事项**：必须大于或等于 `0`；`0` 表示不重试。此项只控制 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.address`；`true` 时 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}` 条目，字段名必须非空且唯一。表达式可从 `KEY`、`VALUE`、`HEADERS`、`TOPIC`、`PARTITION`、`OFFSET` 或 `TIMESTAMP` 提取数据。字段集合由此配置固定，输入记录中的其他字段不会自动发布；需要在表达式内使用逗号时还需遵守 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_CONTINUE`、`FORWARD_TO_DLQ` 或 `TERMINATE_TASK`。忽略和成功写入 DLQ 都会推进该记录的处理位置；使用 `FORWARD_TO_DLQ` 隔离坏记录并继续处理后续记录时，必须同时配置非空的 `errors.deadletterqueue.topic.name` 和 `errors.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`
* **重要级别**：中
* **有效值 / 注意事项**：可选 `none` 或 `all`。`none` 在首条失败记录超过容忍限制时终止 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 列表。

**配置示例**：

```properties theme={null}
connector.class=com.lightstreamer.kafka.connect.LightstreamerSinkConnector
topics.regex=events-.*
lightstreamer.server.proxy_adapter.address=<lightstreamer-host>:6661
topic.mappings=events-.*:events
topic.mappings.regex.enable=true
record.mappings=message:#{VALUE}
value.converter=org.apache.kafka.connect.storage.StringConverter
```

**关键说明**：相对快速开始，此配置移除 `topics`，改用 `topics.regex`，并同时把 `topic.mappings` 的键切换为正则。两套规则相互独立，应保持范围一致；否则 Kafka Connect 可能已消费某个 Topic，但该记录没有匹配的 Item 路由。新增 Topic 前还应确认其消息编码和字段结构与现有映射兼容。

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

**适用业务场景**：客户端只需要某个订单、设备或账户的更新，不希望所有订阅者都接收同一 Topic 的全部记录。此时可用记录 Key 构造参数化 Item，让客户端按业务实体订阅。

**配置示例**：

```properties theme={null}
connector.class=com.lightstreamer.kafka.connect.LightstreamerSinkConnector
topics=entity-events
lightstreamer.server.proxy_adapter.address=<lightstreamer-host>:6661
item.templates=entity-template:entity-#{id=KEY}
topic.mappings=entity-events:item-template.entity-template
record.mappings=message:#{VALUE}
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
```

**关键说明**：相对快速开始，此配置增加 `item.templates` 和 Key Converter，并把静态 Item 改为模板引用。Connector 模板定义为 `entity-#{id=KEY}`；客户端需使用字面筛选值订阅 `entity-[id=42]`，Key 为 `42` 的记录才会匹配该订阅，不能使用 `entity-42`。从记录中提取的模板参数值及其 Connect Schema 必须与订阅参数匹配。业务实体 Key 应稳定且非空，避免因 Key 变化把同一实体的更新路由到不同 Item。

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

**适用业务场景**：Connector 已持续运行，输入 Schema 演进或个别异常记录可能导致字段路径无法提取，希望保留问题记录供排查，同时继续处理后续有效记录。

**配置示例**：

```properties theme={null}
connector.class=com.lightstreamer.kafka.connect.LightstreamerSinkConnector
topics=<kafka-topic>
lightstreamer.server.proxy_adapter.address=<lightstreamer-host>:6661
topic.mappings=<kafka-topic>:events
record.mappings=message:#{VALUE.message}
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=true
record.extraction.error.strategy=FORWARD_TO_DLQ
errors.tolerance=all
errors.deadletterqueue.topic.name=<dlq-topic>
errors.deadletterqueue.context.headers.enable=true
```

**关键说明**：相对快速开始，此配置使用带 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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* 单个 Connector 实例始终只创建一个 Sink Task；提高 `tasks.max` 不会增加任务并行度，需要通过多个 Connector 实例划分输入 Topic 才能进行实例级扩展。
* Connector 不提供初始快照，字段更新由当前客户端订阅驱动；普通主动连接模式下无订阅时不会持续推进 Connector 维护的安全 Offset，运行期间新增订阅不会自动收到此前在无匹配订阅时已消费的记录，重启后则可能从最后已提交位置重放。
* 连接反转模式在没有活动 Proxy Adapter 或没有订阅时仍可能提交 Worker Offset，因此这些记录可能不会产生客户端更新，也不会依靠 Kafka Offset 自动补发。
* 目标 Item 更新与 Kafka Offset 提交不是原子事务；故障恢复可能产生重复更新，Connector 不提供 exactly-once、目标侧去重或终端客户端接收确认。
* `KEY` 和 `VALUE` 路径提取要求 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 的 `topics` 或 `topics.regex` 是否覆盖输入 Topic、`topic.mappings` 是否按正确的字面量或正则方式匹配，以及客户端请求的字段是否存在于 `record.mappings`。普通主动连接模式没有匹配订阅时不会产生更新；建立订阅后发送一条新记录进行验证，不要假设运行中新增订阅会自动回放此前记录。

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

检查 `key.converter`、`value.converter` 与实际消息编码是否一致，并确认 `KEY` 或 `VALUE` 路径对应非空 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 状态和日志并按部署策略重启。
