Socket
Socket 源连接器
支持这些引擎
Spark
Flink
SeaTunnel Zeta
关键特性
描述
用于从 Socket 服务端读取按行分隔的文本数据。Socket 中收到的每一行都会成为一条 STRING 类型的
SeaTunnel 数据。流处理模式下连接器保持连接持续打开并按行处理;批处理模式下读取器只执行一次 read,
将这次读取中已按 \n 切分得到的完整行(以及最后一行末尾不完整的部分作为一行)发送出去后即结束——
它既不会等待对端关闭连接,也没有读取超时设置。
该连接器只使用单个 split(Source 并行度固定为 1)。host/port 指的是 SeaTunnel 要连接的
服务端 地址,对端可以是 Sink、Transform,也可以通过 nc -l 等工具手动提供。
数据类型映射
Socket Source 会把每一行输入读取为字符串。
| SeaTunnel 数据类型 |
|---|
| STRING |
选项
| 参数名 | 类型 | 必须 | 默认值 | 描述 |
|---|---|---|---|---|
| host | String | 是 | - | socket 服务器主机 |
| port | Integer | 是 | - | socket 服务器端口 |
| common-options | 否 | - | 源插件通用参数,请参考 源通用选项 详见。 |
提示
Socket Source 更适合本地调试和简单文本流读取。它不会保存 Socket 服务端的读取位点,如果需要可重放或精确一次读取,请使用 Kafka 等具备位点管理能力的 Source。每行都会作为一条数据 处理;空行不会被跳过,而是会产生一个负载为空字符串的行。
如何创建 Socket 数据同步作业
- 配置 SeaTunnel 配置文件
以下示例演示如何创建从 Socket 读取数据并在本地客户端上打印的数据同步作业:
# 设置要执行的任务的基本配置
env {
parallelism = 1
job.mode = "BATCH"
}
# 创建源以连接到 socket
source {
Socket {
host = "localhost"
port = 9999
}
}
# 控制台打印读取的 socket 数据
sink {
Console {
parallelism = 1
}
}
- 启动端口监听
nc -l 9999
启动 SeaTunnel 任务
Socket 源发送测试数据
~ nc -l 9999
test
hello
flink
spark
- 控制台 Sink 打印数据
[test]
[hello]
[flink]
[spark]
流处理模式
流处理模式下,源端会保持连接持续打开,持续读取新行。建议配合可以缓冲或 checkpoint 的下游 Sink:
env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 10000
}
source {
Socket {
host = "localhost"
port = 9999
}
}
sink {
Console {
parallelism = 1
}
}
变更日志
Change Log
| Change | Commit | Version |
|---|---|---|
| [improve] socket options (#9517) | https://github.com/apache/seatunnel/commit/af83a302cf | 2.3.12 |
| [Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786) | https://github.com/apache/seatunnel/commit/6b7c53d03c | 2.3.9 |
[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755) | https://github.com/apache/seatunnel/commit/8de7408100 | 2.3.4 |
| [Improve][build] Give the maven module a human readable name (#4114) | https://github.com/apache/seatunnel/commit/d7cd601051 | 2.3.1 |
| [Improve][Project] Code format with spotless plugin. (#4101) | https://github.com/apache/seatunnel/commit/a2ab166561 | 2.3.1 |
| [Feature][Connector] add get source method to all source connector (#3846) | https://github.com/apache/seatunnel/commit/417178fb84 | 2.3.1 |
| [Hotfix][OptionRule] Fix option rule about all connectors (#3592) | https://github.com/apache/seatunnel/commit/226dc6a119 | 2.3.0 |
| [Improve][Connector-V2][Socket] Unified exception for socket source & sink connector (#3511) | https://github.com/apache/seatunnel/commit/581292f210 | 2.3.0 |
| [feature][connector][socket] Add Socket Connector Option Rules (#3317) | https://github.com/apache/seatunnel/commit/b85317bcbe | 2.3.0 |
| [Improve][all] change Log to @Slf4j (#3001) | https://github.com/apache/seatunnel/commit/6016100f12 | 2.3.0-beta |
| [DEV][Api] Replace SeaTunnelContext with JobContext and remove singleton pattern (#2706) | https://github.com/apache/seatunnel/commit/cbf82f755c | 2.2.0-beta |
| [#2606]Dependency management split (#2630) | https://github.com/apache/seatunnel/commit/fc047be69b | 2.2.0-beta |
| [Feature][Connector-V2] Socket Connector Sink (#2549) | https://github.com/apache/seatunnel/commit/94f4600a4e | 2.2.0-beta |
| [api-draft][Optimize] Optimize module name (#2062) | https://github.com/apache/seatunnel/commit/f79e3112b1 | 2.2.0-beta |