跳到主要内容
版本:Next

Socket

Socket 数据接收器

支持引擎

Spark
Flink
SeaTunnel Zeta

主要特性

描述

用于向 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 选项

名称类型是否必传默认值描述
hostString-socket 服务器主机
portInteger-socket 服务器端口
max_retriesInteger3发送失败后的最大重试次数。设置为 -1 表示无限重试,0 表示失败后立即抛出异常。
common-options-Sink 插件通用参数,详见 Sink 通用选项
提示

Socket Sink 更适合本地调试和简单集成。它会根据 max_retries 进行重连和重试,但不提供精确一次写入保证。 每个 Writer 会建立一条 TCP 连接;host/port 指的是客户端要连接的 服务端 地址。

任务示例

以下示例把 FakeSource 随机生成的数据写入 Socket Server。

env {
parallelism = 1
job.mode = "STREAMING"
}

source {
FakeSource {
plugin_output = "fake"
schema = {
fields {
name = "string"
age = "int"
}
}
}
}

sink {
Socket {
host = "localhost"
port = 9999
max_retries = 3
}
}
  • 启动端口侦听
nc -l -v 9999
  • 启动 SeaTunnel 任务

  • Socket 服务器控制台打印数据。由于不会追加分隔符,多条记录在原始字节流中以拼接的 JSON 对象形式到达(下面的换行仅为便于阅读):

{"name":"jared","age":17}{"name":"jared","age":18}...

变更日志

Change Log
ChangeCommitVersion
[improve] socket options (#9517)https://github.com/apache/seatunnel/commit/af83a302cf2.3.12
[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)https://github.com/apache/seatunnel/commit/6b7c53d03c2.3.9
[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755)https://github.com/apache/seatunnel/commit/8de74081002.3.4
[Improve][build] Give the maven module a human readable name (#4114)https://github.com/apache/seatunnel/commit/d7cd6010512.3.1
[Improve][Project] Code format with spotless plugin. (#4101)https://github.com/apache/seatunnel/commit/a2ab1665612.3.1
[Feature][Connector] add get source method to all source connector (#3846)https://github.com/apache/seatunnel/commit/417178fb842.3.1
[Hotfix][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Improve][Connector-V2][Socket] Unified exception for socket source & sink connector (#3511)https://github.com/apache/seatunnel/commit/581292f2102.3.0
[feature][connector][socket] Add Socket Connector Option Rules (#3317)https://github.com/apache/seatunnel/commit/b85317bcbe2.3.0
[Improve][all] change Log to @Slf4j (#3001)https://github.com/apache/seatunnel/commit/6016100f122.3.0-beta
[DEV][Api] Replace SeaTunnelContext with JobContext and remove singleton pattern (#2706)https://github.com/apache/seatunnel/commit/cbf82f755c2.2.0-beta
[#2606]Dependency management split (#2630)https://github.com/apache/seatunnel/commit/fc047be69b2.2.0-beta
[Feature][Connector-V2] Socket Connector Sink (#2549)https://github.com/apache/seatunnel/commit/94f4600a4e2.2.0-beta
[api-draft][Optimize] Optimize module name (#2062)https://github.com/apache/seatunnel/commit/f79e3112b12.2.0-beta