This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 66333d02d8 [Docs][Connector-V2] Improve ActiveMQ Web3j Sls Socket and
EdgeSocket docs (#11575)
66333d02d8 is described below
commit 66333d02d84a92a85781455f1cf95dc902bb3aae
Author: Daniel Carter <[email protected]>
AuthorDate: Tue Aug 11 22:39:39 2026 +0800
[Docs][Connector-V2] Improve ActiveMQ Web3j Sls Socket and EdgeSocket docs
(#11575)
Co-authored-by: DanielCarter-stack
<[email protected]>
---
docs/en/connectors/sink/Activemq.md | 33 ++++++++++++++++++++
docs/en/connectors/sink/Sls.md | 50 ++++++++++++++++++++++++++++++-
docs/en/connectors/sink/Socket.md | 20 +++++++++----
docs/en/connectors/source/EdgeSocket.md | 45 ++++++++++++++++++++++++++++
docs/en/connectors/source/Socket.md | 40 +++++++++++++++++++++++--
docs/en/connectors/source/Web3j.md | 53 +++++++++++++++++++++++++++++++++
docs/zh/connectors/sink/Activemq.md | 32 ++++++++++++++++++++
docs/zh/connectors/sink/Sls.md | 50 +++++++++++++++++++++++++++++--
docs/zh/connectors/sink/Socket.md | 17 +++++++----
docs/zh/connectors/source/EdgeSocket.md | 45 ++++++++++++++++++++++++++++
docs/zh/connectors/source/Socket.md | 36 ++++++++++++++++++++--
docs/zh/connectors/source/Web3j.md | 51 +++++++++++++++++++++++++++++++
12 files changed, 455 insertions(+), 17 deletions(-)
diff --git a/docs/en/connectors/sink/Activemq.md
b/docs/en/connectors/sink/Activemq.md
index 3d5924eb1c..072febf334 100644
--- a/docs/en/connectors/sink/Activemq.md
+++ b/docs/en/connectors/sink/Activemq.md
@@ -31,6 +31,7 @@ a sink-only connector; SeaTunnel does not provide an ActiveMQ
source connector.
| dispatch_async | boolean | no | -
| Whether the broker dispatches messages asynchronously.
|
| nested_map_and_list_enabled | boolean | no | -
| Whether structured message properties and `MapMessage` entries can contain
nested `Map` and `List` objects.
|
| warn_about_unstarted_connection_timeout | int | no | -
| Timeout in milliseconds before ActiveMQ warns that a connection was not
started correctly. Set a value less than `0` to disable the warning in the
ActiveMQ client. |
+| consumer_expiry_check_enabled | boolean | no | -
| Whether the ActiveMQ client checks message expiration in each
`MessageConsumer` before dispatching messages.
|
## Notes
@@ -75,6 +76,38 @@ sink {
}
```
+In streaming mode, the sink keeps the same broker connection open and writes
each row as it
+arrives. Username/password can also be embedded in the `uri`, for example
+`tcp://admin:admin@localhost:61616`:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ id = int
+ name = string
+ }
+ }
+ rows = [
+ { kind = INSERT, fields = [1, "Alice"] }
+ ]
+ }
+}
+
+sink {
+ ActiveMQ {
+ uri = "tcp://admin:admin@localhost:61616"
+ queue_name = "testQueue"
+ }
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/en/connectors/sink/Sls.md b/docs/en/connectors/sink/Sls.md
index c1809f6a87..3c18322d45 100644
--- a/docs/en/connectors/sink/Sls.md
+++ b/docs/en/connectors/sink/Sls.md
@@ -46,11 +46,15 @@ Maven central repository.
- The configured RAM user must have permission to write logs to the target
project and logstore.
- The sink writes data as soon as `write` is called. It does not provide
exactly-once commit semantics.
+ In streaming mode the connector flushes row by row; rely on checkpointing
only for downstream state,
+ not for the SLS writes themselves.
+- Each row is serialized as a JSON object and stored under the `content` key
of an SLS log item. The
+ remaining row fields are not mapped to separate log keys.
- Do not print `access_key_secret` in logs or job descriptions.
## Task Example
-### Write Rows to SLS
+### Write Rows to SLS (Batch)
```hocon
env {
@@ -89,6 +93,50 @@ sink {
}
```
+### Write Rows to SLS (Streaming)
+
+In streaming mode the connector keeps the SLS producer connection open and
writes each row as it
+arrives. Configure `checkpoint.interval` to make downstream state recoverable,
but keep in mind that
+each `PutLogs` call is independent and may retry only within the open client
session.
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 30000
+}
+
+source {
+ FakeSource {
+ row.num = 10
+ map.size = 10
+ array.size = 10
+ bytes.length = 10
+ string.length = 10
+ schema = {
+ fields = {
+ id = "int"
+ name = "string"
+ description = "string"
+ weight = "string"
+ }
+ }
+ }
+}
+
+sink {
+ Sls {
+ endpoint = "cn-hangzhou.log.aliyuncs.com"
+ project = "project1"
+ logstore = "logstore1"
+ access_key_id = "xxxxxxxxxxxxxxxxxxxxxxxx"
+ access_key_secret = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
+ source = "seatunnel-streaming"
+ topic = "fake-source"
+ }
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/en/connectors/sink/Socket.md
b/docs/en/connectors/sink/Socket.md
index e2c5b105dc..3d61cb7f32 100644
--- a/docs/en/connectors/sink/Socket.md
+++ b/docs/en/connectors/sink/Socket.md
@@ -16,7 +16,15 @@ import ChangeLog from '../changelog/connector-socket.md';
## Description
-Used to send data to a socket server in streaming or batch mode. Each
SeaTunnel row is serialized as one JSON line.
+Used to send data to a socket server in streaming or batch mode. Each
SeaTunnel row is serialized to a
+JSON object via `JsonSerializationSchema` and written to the configured TCP
port. **The connector does
+not append any delimiter at all** — neither a newline, nor any other separator
between records. Multiple
+records therefore travel as one undelimited, continuous TCP byte stream of
concatenated JSON objects
+(for example `{"a":1}{"a":2}{"a":3}`). The output is explicitly *not*
line-framed JSON, so the peer
+must handle framing itself: parse consecutive JSON values with a streaming
JSON parser (such as
+Jackson's `MappingIterator`) rather than a line-oriented parser. Tools like
`nc -l` only echo the raw
+concatenated bytes, so they are useful for a quick single-row check but cannot
split records on their
+own.
> For example, if the data from upstream is [`age: 12, name: jared`], the
> content send to socket server is the following: `{"name":"jared","age":17}`
@@ -26,12 +34,13 @@ Used to send data to a socket server in streaming or batch
mode. Each SeaTunnel
|----------------|---------|----------|---------|-----------------------------------------------------------------------------------------------------------------|
| host | String | Yes | | socket server host
|
| port | Integer | Yes | | socket server port
|
-| max_retries | Integer | No | 3 | The number of retries to
send record failed
|
+| max_retries | Integer | No | 3 | The number of retries to
send record failed. Set to `-1` to retry indefinitely, or `0` to fail
immediately. |
| common-options | | No | - | Sink plugin common
parameters, please refer to [Sink Common
Options](../common-options/sink-common-options.md) for details |
:::tip
-Socket sink is mainly used for local debugging and simple integrations. It
reconnects and retries failed writes according to `max_retries`, but it does
not provide exactly-once delivery.
+Socket sink is mainly used for local debugging and simple integrations. It
reconnects and retries failed writes according to `max_retries`, but it does
not provide exactly-once delivery. The TCP client
+opens one connection per writer; `host`/`port` are the *server* endpoint that
this client connects to.
:::
@@ -61,6 +70,7 @@ sink {
Socket {
host = "localhost"
port = 9999
+ max_retries = 3
}
}
```
@@ -73,10 +83,10 @@ nc -l -v 9999
* Start a SeaTunnel task
-* Socket Server Console print data
+* Socket Server Console print data. No delimiter is appended, so multiple rows
arrive as concatenated JSON objects in the raw byte stream (line breaks shown
here only for readability):
```text
-{"name":"jared","age":17}
+{"name":"jared","age":17}{"name":"jared","age":18}...
```
## Changelog
diff --git a/docs/en/connectors/source/EdgeSocket.md
b/docs/en/connectors/source/EdgeSocket.md
index b56c5d469a..1c7f4dbfc5 100644
--- a/docs/en/connectors/source/EdgeSocket.md
+++ b/docs/en/connectors/source/EdgeSocket.md
@@ -91,6 +91,51 @@ sink {
}
```
+### Common Downstream Pattern
+
+EdgeSocket is a `STRING` (or schema-typed) source. The most common production
pattern is to
+chain a `Sql` transform and a database sink so each batch becomes one or more
rows in a target
+table. The example below mirrors what the bundled e2e tests use:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 3000
+}
+
+source {
+ EdgeSocket {
+ port = 10091
+ auth_type = "TOKEN"
+ token = "edge-token"
+ packet_mode = "RAW"
+ max_retries = 3
+ reconnect_interval_ms = 2000
+ accept_timeout_ms = 5000
+ }
+}
+
+transform {
+ Sql {
+ query = "SELECT CONCAT(value, '_transformed') AS value_text FROM
source_table"
+ }
+}
+
+sink {
+ Jdbc {
+ url =
"jdbc:mysql://mysql-e2e:3306/seatunnel?useSSL=false&allowPublicKeyRetrieval=true"
+ driver = "com.mysql.cj.jdbc.Driver"
+ username = "root"
+ password = "mysqlpw"
+ query = "insert into edge_socket_sink(value_text) values (?)"
+ }
+}
+```
+
+Other database sinks such as `Postgres`, `ClickHouse`, or `Kafka` work the
same way; the source has
+no first-class binding to a specific sink, so choose whichever target your
collectors need.
+
## Schema Mode
By default the source emits one STRING field (value) containing the raw
payload line.
diff --git a/docs/en/connectors/source/Socket.md
b/docs/en/connectors/source/Socket.md
index 2493801027..3ed34c0baf 100644
--- a/docs/en/connectors/source/Socket.md
+++ b/docs/en/connectors/source/Socket.md
@@ -21,7 +21,16 @@ import ChangeLog from '../changelog/connector-socket.md';
## Description
-Used to read newline-delimited text data from a socket server. Each line
received from the socket becomes one SeaTunnel row.
+Used to read newline-delimited text data from a socket server. Each line
received from the socket
+becomes one SeaTunnel row of type `STRING`. In streaming mode the source stays
connected to the
+socket and reads lines as they arrive; in batch mode the reader performs a
single read of whatever
+data is currently available on the socket, emits any complete
newline-terminated lines from that read
+(plus any trailing partial line as a final row), and then finishes — it does
not wait for the
+connection to close and there is no read-timeout setting.
+
+The connector uses a single split (source parallelism is fixed at 1). `host`
and `port` refer to the
+*server* endpoint that SeaTunnel connects to; configure a sink, transformer,
or peer like `nc -l`
+on the other side.
## Data Type Mapping
@@ -41,7 +50,8 @@ Socket source reads each incoming line as a string record.
:::tip
-Socket source is mainly used for local debugging and simple text streams. It
does not checkpoint socket-server offsets, so it should not be used when
replayable, exactly-once reads are required.
+Socket source is mainly used for local debugging and simple text streams. It
does not checkpoint socket-server offsets, so it should not be used when
replayable, exactly-once reads are required. Each line
+is treated as one record. Empty lines produce a row with an empty-string
payload; they are not skipped.
:::
@@ -101,6 +111,32 @@ spark
[spark]
```
+### Streaming Mode
+
+In streaming mode the source keeps the socket open and reads new lines
continuously. Pair it with a
+downstream sink that can buffer events or checkpoint them:
+
+```bash
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ Socket {
+ host = "localhost"
+ port = 9999
+ }
+}
+
+sink {
+ Console {
+ parallelism = 1
+ }
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/en/connectors/source/Web3j.md
b/docs/en/connectors/source/Web3j.md
index 959fd513fa..e71d2f68a0 100644
--- a/docs/en/connectors/source/Web3j.md
+++ b/docs/en/connectors/source/Web3j.md
@@ -28,6 +28,10 @@ string that contains `blockNumber` and the read timestamp.
In batch mode, the source emits one row and then finishes. In streaming mode,
it keeps polling the
provider and emits the latest observed block number.
+The connector uses a single split and does not support parallelism. Each row
produced contains the
+result of one HTTP `eth_blockNumber` call, so the effective polling rate
follows the response time
+of the configured provider.
+
## Source Options
| Name | Type | Required | Default | Description |
@@ -46,8 +50,20 @@ The JSON stored in `value` has this shape:
{"blockNumber":19525949,"timestamp":"2024-03-27T13:28:45.605Z"}
```
+## Notes
+
+- The `url` must point to a JSON-RPC compatible Web3 provider, such as Infura,
Alchemy, or a
+ self-hosted Ethereum node. HTTPS is recommended; the connector does not
perform additional
+ authentication, so put the API key directly into the URL when the provider
requires one.
+- The connector exposes only a single row type with the `value` field. Use a
SQL transform or JSON
+ path downstream to extract `blockNumber` or `timestamp` for further
processing.
+- In streaming mode the connector keeps the HTTP connection open and emits the
latest block
+ number observed on each poll; pair it with checkpointing only if downstream
sinks require it.
+
## Example
+In batch mode, the source emits one row and then finishes:
+
```hocon
env {
parallelism = 1
@@ -75,6 +91,43 @@ Then you will get data similar to the following:
{"value":"{\"blockNumber\":19525949,\"timestamp\":\"2024-03-27T13:28:45.605Z\"}"}
```
+In streaming mode, the connector keeps polling the provider and emits a row on
every poll containing
+the latest observed block number:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ Web3j {
+ url = "https://mainnet.infura.io/v3/xxxxx"
+ plugin_output = "web3j"
+ }
+}
+
+sink {
+ Assert {
+ plugin_input = "web3j"
+ rules {
+ field_rules = [
+ {
+ field_name = value
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
+```
+
## Changelog
<ChangeLog />
diff --git a/docs/zh/connectors/sink/Activemq.md
b/docs/zh/connectors/sink/Activemq.md
index 41aa57dd79..fbb1cea36a 100644
--- a/docs/zh/connectors/sink/Activemq.md
+++ b/docs/zh/connectors/sink/Activemq.md
@@ -31,6 +31,7 @@ Sink,SeaTunnel 目前没有提供 ActiveMQ Source 连接器。
| dispatch_async | boolean | 否 | - | Broker
是否异步分发消息。
|
| nested_map_and_list_enabled | boolean | 否 | - | 是否允许结构化消息属性和
`MapMessage` 条目中包含嵌套的 `Map`、`List` 对象。
|
| warn_about_unstarted_connection_timeout | int | 否 | - |
连接没有正确启动时,ActiveMQ 客户端发出警告前等待的毫秒数。设置为小于 `0` 的值可以关闭这个警告。
|
+| consumer_expiry_check_enabled | boolean | 否 | - | 是否在每个
`MessageConsumer` 分发消息前检查消息是否已经过期。
|
## 注意事项
@@ -75,6 +76,37 @@ sink {
}
```
+在流式模式下,Sink 会保持与 Broker 的连接持续打开,每来一行数据就写入一条。用户名 / 密码也可以直接
+写在 `uri` 中,例如 `tcp://admin:admin@localhost:61616`:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ id = int
+ name = string
+ }
+ }
+ rows = [
+ { kind = INSERT, fields = [1, "Alice"] }
+ ]
+ }
+}
+
+sink {
+ ActiveMQ {
+ uri = "tcp://admin:admin@localhost:61616"
+ queue_name = "testQueue"
+ }
+}
+```
+
## 变更日志
<ChangeLog />
diff --git a/docs/zh/connectors/sink/Sls.md b/docs/zh/connectors/sink/Sls.md
index 215ecd3116..e398948446 100644
--- a/docs/zh/connectors/sink/Sls.md
+++ b/docs/zh/connectors/sink/Sls.md
@@ -44,12 +44,14 @@ JSON,然后作为 SLS 日志项写入,日志内容的 key 为 `content`。
## 注意事项
- 配置的 RAM 用户需要有向目标 project 和 logstore 写入日志的权限。
-- sink 在收到数据时立即写入,不提供精确一次提交语义。
+- sink 在收到数据时立即写入,不提供精确一次提交语义。流处理模式下连接器按行写入;checkpoint
+ 只对下游状态有用,并不能保证 SLS 端的写入语义。
+- 每条数据都会被序列化为 JSON,并写入 SLS 日志项 `content` 字段,不会映射到其它日志 key。
- 不要在日志或任务说明里打印 `access_key_secret`。
## 任务示例
-### 写入数据到 SLS
+### 写入数据到 SLS(批处理)
```hocon
env {
@@ -88,6 +90,50 @@ sink {
}
```
+### 写入数据到 SLS(流处理)
+
+流处理模式下,连接器会保持 SLS Producer 的连接持续打开,每来一行数据就写入一条。
+可以配置 `checkpoint.interval` 保护下游状态,但需要清楚每条 `PutLogs` 调用互相独立,
+重试只在 Producer 会话内进行,不会跨重启。
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 30000
+}
+
+source {
+ FakeSource {
+ row.num = 10
+ map.size = 10
+ array.size = 10
+ bytes.length = 10
+ string.length = 10
+ schema = {
+ fields = {
+ id = "int"
+ name = "string"
+ description = "string"
+ weight = "string"
+ }
+ }
+ }
+}
+
+sink {
+ Sls {
+ endpoint = "cn-hangzhou.log.aliyuncs.com"
+ project = "project1"
+ logstore = "logstore1"
+ access_key_id = "xxxxxxxxxxxxxxxxxxxxxxxx"
+ access_key_secret = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
+ source = "seatunnel-streaming"
+ topic = "fake-source"
+ }
+}
+```
+
## 变更日志
<ChangeLog />
diff --git a/docs/zh/connectors/sink/Socket.md
b/docs/zh/connectors/sink/Socket.md
index 133ad93060..ac6af3e326 100644
--- a/docs/zh/connectors/sink/Socket.md
+++ b/docs/zh/connectors/sink/Socket.md
@@ -17,22 +17,28 @@ import ChangeLog from '../changelog/connector-socket.md';
## 描述
-用于向 Socket Server 发送数据,支持流模式和批模式。每条 SeaTunnel 数据会被序列化为一行 JSON。
+用于向 Socket Server 发送数据,支持流模式和批模式。每条 SeaTunnel 数据会被 `JsonSerializationSchema`
+序列化为一个 JSON 对象,并写入配置的 TCP 端口。**连接器不会追加任何分隔符**——既不会追加换行符,也不会
+在记录之间追加任何其它分隔符。因此多条记录会作为一条无分隔、连续的 TCP 字节流直接拼接在一起传输
+(例如 `{"a":1}{"a":2}{"a":3}`)。输出明确*不是*按行分隔的 JSON,因此对端需要自行处理分帧:
+使用支持连续读取多个 JSON 值的流式解析器(例如 Jackson 的 `MappingIterator`),而不是按行解析的解析器。
+`nc -l` 这类工具只会原样回显拼接后的字节,适合做单条记录的快速验证,但无法自行切分多条记录。
> 例如,如果来自上游的数据是 [`age: 17, name: jared`],则发送到 Socket Server
> 的内容如下:`{"name":"jared","age":17}`
## Sink 选项
| 名称 | 类型 | 是否必传 | 默认值 |
描述 |
-|----------------|---------|----------|---------|-----------------------------------------------------------------------------------------------------------------|
+|----------------|---------|----------|---------|----------------------------------------------------------------------------------------------------------------|
| host | String | 是 | - | socket 服务器主机
|
| port | Integer | 是 | - | socket 服务器端口
|
-| max_retries | Integer | 否 | 3 | 发送失败后的最大重试次数
|
+| max_retries | Integer | 否 | 3 | 发送失败后的最大重试次数。设置为 `-1`
表示无限重试,`0` 表示失败后立即抛出异常。 |
| common-options | | 否 | - | Sink 插件通用参数,详见 [Sink
通用选项](../common-options/sink-common-options.md) |
:::tip
Socket Sink 更适合本地调试和简单集成。它会根据 `max_retries` 进行重连和重试,但不提供精确一次写入保证。
+每个 Writer 会建立一条 TCP 连接;`host`/`port` 指的是客户端要连接的 *服务端* 地址。
:::
@@ -62,6 +68,7 @@ sink {
Socket {
host = "localhost"
port = 9999
+ max_retries = 3
}
}
```
@@ -74,10 +81,10 @@ nc -l -v 9999
* 启动 SeaTunnel 任务
-* Socket 服务器控制台打印数据
+* Socket 服务器控制台打印数据。由于不会追加分隔符,多条记录在原始字节流中以拼接的 JSON 对象形式到达(下面的换行仅为便于阅读):
```text
-{"name":"jared","age":17}
+{"name":"jared","age":17}{"name":"jared","age":18}...
```
## 变更日志
diff --git a/docs/zh/connectors/source/EdgeSocket.md
b/docs/zh/connectors/source/EdgeSocket.md
index 51bea57f3e..78f6f832d0 100644
--- a/docs/zh/connectors/source/EdgeSocket.md
+++ b/docs/zh/connectors/source/EdgeSocket.md
@@ -88,6 +88,51 @@ sink {
}
```
+### 常见下游模式
+
+EdgeSocket 输出的是 `STRING`(或声明 schema 后的结构化类型)记录。生产环境最常见的组合是先
+用 `Sql` 转换把单行字符串处理一下,再写入数据库 Sink,这样每个批次最终都会落到目标表的多条
+记录中。下面这段配置和仓库自带 e2e 中的用法保持一致:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 3000
+}
+
+source {
+ EdgeSocket {
+ port = 10091
+ auth_type = "TOKEN"
+ token = "edge-token"
+ packet_mode = "RAW"
+ max_retries = 3
+ reconnect_interval_ms = 2000
+ accept_timeout_ms = 5000
+ }
+}
+
+transform {
+ Sql {
+ query = "SELECT CONCAT(value, '_transformed') AS value_text FROM
source_table"
+ }
+}
+
+sink {
+ Jdbc {
+ url =
"jdbc:mysql://mysql-e2e:3306/seatunnel?useSSL=false&allowPublicKeyRetrieval=true"
+ driver = "com.mysql.cj.jdbc.Driver"
+ username = "root"
+ password = "mysqlpw"
+ query = "insert into edge_socket_sink(value_text) values (?)"
+ }
+}
+```
+
+其它数据库 Sink(例如 `Postgres`、`ClickHouse`、`Kafka` 等)都可以按同样方式接在 EdgeSocket
+后面;Source 与具体 Sink 之间并不存在强绑定关系,请按业务需要选择目标。
+
## Schema 模式
默认输出单个 STRING 字段(字段名 value),内容为 payload 原始文本。
diff --git a/docs/zh/connectors/source/Socket.md
b/docs/zh/connectors/source/Socket.md
index 7d1d3d719f..74d7adc7d1 100644
--- a/docs/zh/connectors/source/Socket.md
+++ b/docs/zh/connectors/source/Socket.md
@@ -21,7 +21,13 @@ import ChangeLog from '../changelog/connector-socket.md';
## 描述
-用于从 Socket 服务端读取按行分隔的文本数据。Socket 中收到的每一行都会成为一条 SeaTunnel 数据。
+用于从 Socket 服务端读取按行分隔的文本数据。Socket 中收到的每一行都会成为一条 `STRING` 类型的
+SeaTunnel 数据。流处理模式下连接器保持连接持续打开并按行处理;批处理模式下读取器只执行一次 `read`,
+将这次读取中已按 `\n` 切分得到的完整行(以及最后一行末尾不完整的部分作为一行)发送出去后即结束——
+它既不会等待对端关闭连接,也没有读取超时设置。
+
+该连接器只使用单个 split(Source 并行度固定为 1)。`host`/`port` 指的是 SeaTunnel 要连接的
+*服务端* 地址,对端可以是 Sink、Transform,也可以通过 `nc -l` 等工具手动提供。
## 数据类型映射
@@ -41,7 +47,8 @@ Socket Source 会把每一行输入读取为字符串。
:::tip
-Socket Source 更适合本地调试和简单文本流读取。它不会保存 Socket 服务端的读取位点,如果需要可重放或精确一次读取,请使用 Kafka
等具备位点管理能力的 Source。
+Socket Source 更适合本地调试和简单文本流读取。它不会保存 Socket 服务端的读取位点,如果需要可重放或精确一次读取,请使用 Kafka
等具备位点管理能力的 Source。每行都会作为一条数据
+处理;空行不会被跳过,而是会产生一个负载为空字符串的行。
:::
@@ -101,6 +108,31 @@ spark
[spark]
```
+### 流处理模式
+
+流处理模式下,源端会保持连接持续打开,持续读取新行。建议配合可以缓冲或 checkpoint 的下游 Sink:
+
+```bash
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ Socket {
+ host = "localhost"
+ port = 9999
+ }
+}
+
+sink {
+ Console {
+ parallelism = 1
+ }
+}
+```
+
## 变更日志
<ChangeLog />
diff --git a/docs/zh/connectors/source/Web3j.md
b/docs/zh/connectors/source/Web3j.md
index 707181529d..bc90c7adaf 100644
--- a/docs/zh/connectors/source/Web3j.md
+++ b/docs/zh/connectors/source/Web3j.md
@@ -26,6 +26,9 @@ Web3j 源连接器用于通过 Web3 服务端点读取区块链数据。目前
批处理模式下,source 输出一行后结束。流处理模式下,它会持续轮询服务端点,并输出观察到的最新区块号。
+该连接器只使用单个分片,不支持并行度。每一条数据对应一次 `eth_blockNumber` HTTP 调用,
+实际轮询节奏由所配置的 Provider 响应速度决定。
+
## 源选项
| 参数名 | 类型 | 必须 | 默认值 | 描述 |
@@ -44,8 +47,20 @@ Web3j 源连接器用于通过 Web3 服务端点读取区块链数据。目前
{"blockNumber":19525949,"timestamp":"2024-03-27T13:28:45.605Z"}
```
+## 注意事项
+
+- `url` 必须指向兼容 JSON-RPC 的 Web3 Provider,例如 Infura、Alchemy 或者自建的以太坊节点。
+ 推荐使用 HTTPS;连接器不再做额外的鉴权,如果 Provider 需要 API Key,直接把 Key 写在 URL
+ 里即可。
+- 连接器只暴露包含 `value` 字段的固定行结构。如需进一步处理 `blockNumber` 或 `timestamp`,
+ 请在下游使用 SQL Transform 或 JSON Path。
+- 流处理模式下,连接器会保持 HTTP 连接持续打开,并按轮询节奏把观察到的最新区块号写入下游;
+ 只有当下游 Sink 需要 checkpoint 时才建议配置 `checkpoint.interval`。
+
## 示例
+批处理模式下,Source 输出一行数据后即结束:
+
```hocon
env {
parallelism = 1
@@ -73,6 +88,42 @@ sink {
{"value":"{\"blockNumber\":19525949,\"timestamp\":\"2024-03-27T13:28:45.605Z\"}"}
```
+流处理模式下,连接器持续轮询 Provider,每次轮询都会输出一行包含当前最新区块号的记录:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 10000
+}
+
+source {
+ Web3j {
+ url = "https://mainnet.infura.io/v3/xxxxx"
+ plugin_output = "web3j"
+ }
+}
+
+sink {
+ Assert {
+ plugin_input = "web3j"
+ rules {
+ field_rules = [
+ {
+ field_name = value
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
+```
+
## 变更日志
<ChangeLog />