Apache SeaTunnel SLS Sink 连接器:将 SeaTunnel 行数据写入阿里云日志服务的实现与配置指南
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文基于 SeaTunnel 官方文档 SLS Sink 文档,系统讲解 Sls sink 连接器的功能定位、参数配置、完整作业配置示例,并结合 connector-sls 模块源码 剖析其"逐行 JSON 序列化 + PutLogs 直写"的实现原理与语义边界,帮助你正确配置并合理评估该连接器在批/流作业中的适用场景。
一、连接器定位与能力边界
Sls sink 连接器将 SeaTunnel 行数据写入阿里云 Simple Log Service(SLS)。官方文档给出的核心描述是:
The Sls sink connector writes SeaTunnel rows to Alibaba Cloud Simple Log Service (SLS). Each SeaTunnel row is serialized as JSON and written to SLS as a log item whose content key is
content.
即:每一行 SeaTunnelRow 会被整体序列化为一个 JSON 字符串,写入 SLS 日志项的content字段,而不是将行内各列映射为多个日志键值对。这一点在 SeatunnelRowSerialization 中可以直接验证:
public List<LogItem> serializeRow(SeaTunnelRow row) { List<LogItem> logGroup = new ArrayList<LogItem>(); LogItem logItem = new LogItem(); String rowJson = new String(jsonSerializationSchema.serialize(row)); LogContent content = new LogContent("content", rowJson); logItem.PushBack(content); logGroup.add(logItem); return logGroup; }可以看到实现非常直接:复用seatunnel-format-json模块的JsonSerializationSchema把整行编码为 JSON 字节,再包装为LogContent("content", rowJson)。如果你需要按列检索 SLS 日志,需要自行在 SLS 侧配置 JSON 字段提取,而不是期望连接器逐列写入。
官方文档明确声明该连接器支持三种执行引擎:
- Spark
- Flink
- SeaTunnel Zeta
同时,在 Key Features 一栏中,exactly-once、cdc、timer flush三项均未勾选,意味着该连接器不提供精确一次提交语义,也不面向 CDC 场景设计。
二、Sink 参数详解
以下是文档中完整的 Sink Options 表格,全部参数在 SlsBaseOptions 与 SlsSinkOptions 中有对应定义:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
| endpoint | String | Yes | - | Alibaba Cloud SLS endpoint, for examplecn-hangzhou.log.aliyuncs.comor an intranet endpoint. |
| project | String | Yes | - | Alibaba Cloud SLS project。 |
| logstore | String | Yes | - | Alibaba Cloud SLS logstore。 |
| access_key_id | String | Yes | - | Alibaba Cloud AccessKey ID。 |
| access_key_secret | String | Yes | - | Alibaba Cloud AccessKey secret。 |
| source | String | No | SeaTunnel-Source | Source tag written to SLS log groups。 |
| topic | String | No | SeaTunnel-Topic | Topic tag written to SLS log groups。 |
从源码可以确认几个细节:
- 必填性由 SlsSinkFactory 的
OptionRule声明:ENDPOINT、PROJECT、LOGSTORE、ACCESS_KEY_ID、ACCESS_KEY_SECRET五项为required,SOURCE、TOPIC为optional,与文档表格一致。 source与topic的默认值SeaTunnel-Source/SeaTunnel-Topic定义在SlsSinkOptions中,它们会作为 SLS log group 的标签写入,用于在 SLS 中区分数据来源与主题。SlsSinkOptions中还存在一个log_group_size(默认 100,描述为 "Aliyun sls log group write size")的选项,SlsSinkWriter构造时会读取该值;但从当前write方法的结构看,每次写入只包含单条LogItem,该参数在现行写入路径中并未实际参与批量切分,属于源码结构中的预留项。- 连接器标识(
CONNECTOR_IDENTITY)为Sls,即作业配置中sink { Sls { ... } }的名称来源;工厂类通过@AutoService(Factory.class)注册,由 SeaTunnel 的插件发现机制加载。
依赖获取方面,文档说明可通过install-plugin.sh或从 Maven 中央仓库下载org.apache.seatunnel:connector-sls,版本标注为 Universal(通用,随主版本发布)。
三、写入流程:从 write 调用到 PutLogs
理解了配置之后,值得深入 SlsSinkWriter 看数据是如何发出的。关键调用链如下:
- 构造阶段:每个 writer 实例化一个阿里云 SLS SDK 客户端:
this.client = new Client( pluginConfig.get(SlsSinkOptions.ENDPOINT), pluginConfig.get(SlsSinkOptions.ACCESS_KEY_ID), pluginConfig.get(SlsSinkOptions.ACCESS_KEY_SECRET));同时从配置读取project、logStore、topic、source,并创建SeatunnelRowSerialization序列化器。
- 写入阶段:
write(SeaTunnelRow element)被调用时,序列化该行并立即发起网络写入:
public void write(SeaTunnelRow element) throws IOException { List<LogItem> data = this.seatunnelRowSerialization.serializeRow(element); PutLogsRequest plr = new PutLogsRequest(project, logStore, topic, source, data); try { this.client.PutLogs(plr); } catch (Throwable e) { log.error("Failed to write logs to SLS", e); throw new IOException(e); } }也就是说PutLogs请求在write调用时同步发出(阿里云 SDK 客户端内部可能带有自身的缓冲与重试),写入失败会记录错误日志并抛出IOException使作业失败。
提交阶段:
prepareCommit()直接返回Optional.empty(),注释明确写着 "nothing to do, when write function, data had sended"(数据在 write 时已经发送);SlsSinkCommitter 的commit同样不做任何事,snapshotState返回空列表。这与文档 Key Features 中未勾选 exactly-once 是一致的:该连接器不存在两阶段提交协议,checkpoint 仅用于上下游状态,不覆盖 SLS 写入本身。关闭阶段:
close()调用client.shutdown()释放 SLS 客户端资源。
SlsSink类本身(见 SlsSink)实现了SeaTunnelSink<SeaTunnelRow, SlsSinkState, SlsCommitInfo, SlsAggregatedCommitInfo>,createWriter时为每个 writer 传入空的slsStates列表——因为不存在需要恢复的写入状态。
四、完整作业配置示例
以下两个示例完整继承自官方文档,可直接作为生产配置的模板。
1. 批模式写入 SLS(Batch)
env { parallelism = 1 job.mode = "BATCH" } 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-intranet.log.aliyuncs.com" project = "project1" logstore = "logstore1" access_key_id = "xxxxxxxxxxxxxxxxxxxxxxxx" access_key_secret = "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" source = "seatunnel-demo" topic = "fake-source" } }批模式示例中使用的是内网 endpoint(cn-hangzhou-intranet.log.aliyuncs.com),适用于作业与 SLS 同区域部署、走内网链路降低延迟与流量的场景。
2. 流模式写入 SLS(Streaming)
文档说明:流模式下连接器保持 SLS producer 连接打开,每到达一行就写入一次;应配置checkpoint.interval保证下游状态可恢复,但要注意每次PutLogs调用相互独立,重试仅发生在客户端会话范围内。
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" } }五、注意事项与语义边界
官方文档 Notes 一节给出的四条注意事项,结合源码可以逐条落到实处:
- 权限要求:配置的 RAM 用户必须拥有向目标 project 和 logstore 写日志的权限,否则会直接触发
write抛出的IOException导致作业失败。 - 写入时机与语义:数据在
write被调用时立即写出,连接器不提供 exactly-once 提交语义;流模式下逐行 flush,checkpoint 仅保障下游状态,不保障 SLS 写入本身。对重复敏感的场景应在 SLS 消费侧做幂等处理。 - 字段映射模型:每行序列化为一个 JSON 对象并整体存放在日志项的
content键下,其余行字段不会被拆分成独立的日志键。 - 密钥安全:不要在日志或作业描述中打印
access_key_secret。由于该值会随客户端构造传入 SLS SDK,建议结合作业配置管理手段(环境变量、密钥托管等)注入。
六、端到端验证参考
仓库中为该连接器提供了 e2e 测试骨架:SlsIT,配套的 sink 配置 sls_sink_to_console.conf 展示了"FakeSource → Sls"的最小结构(endpoint/project/logstore/密钥均为占位符,实际运行需替换为真实凭据)。此外同目录下的sls_source_with_schema_to_console.conf、sls_source_without_schema_to_console.conf覆盖 SLS 作为 source 的场景,说明该连接器在仓库中同时提供 source 与 sink 两种角色(本文聚焦 sink)。
七、演进记录
从 SLS 连接器 Changelog 可以看到该模块的演进脉络:2.3.7 引入 Aliyun SLS 连接器基础能力,2.3.9 起补齐 sink 连接器、e2e 与文档,后续版本持续优化选项结构与 enumerator API 语义。当前 sink 行为应以仓库中源码实现为准。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考