SeaTunnel Sls Sink 连接器实战指南:将数据写入阿里云日志服务 SLS
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文聚焦 Apache SeaTunnel 的Sls Sink 连接器,系统讲解如何将 SeaTunnel 处理后的数据写入阿里云日志服务(SLS)。全文覆盖连接器的功能定位、全部 Sink 配置项与源码级实现原理、批处理与流处理两种任务配置示例,以及精确一次语义、权限、序列化等关键注意事项。读完本文,你将掌握用 SeaTunnel 把任意结构化数据以 JSON 形式写入阿里云 SLS 的完整方案,并理解其底层写入链路。
概述
Sls Sink 连接器用于把 SeaTunnel 数据写入阿里云日志服务 SLS。每条 SeaTunnel 数据(SeaTunnelRow)会先被序列化为 JSON 字符串,然后作为 SLS 日志项写入,日志内容的 key 固定为content。也就是说,SLS 侧每收到一条日志,其content字段就是一条完整 JSON,原始行中的各个字段不会映射到 SLS 的其它日志 key。
从源码结构看,该连接器位于仓库的 connector-sls 模块,同时提供 Source 与 Sink 双向能力,本文聚焦 Sink 侧。连接器的插件标识(CONNECTOR_IDENTITY)为Sls,定义于 SlsBaseOptions.java。
支持的引擎
Sls Sink 连接器支持以下三类运行引擎:
- Spark
- Flink
- SeaTunnel Zeta
主要特性
| 特性 | 支持情况 |
|---|---|
| 精确一次(Exactly Once) | 不支持 |
| CDC | 不支持 |
| 定时刷新(Scheduled Refresh) | 不支持 |
关于这些特性的统一定义与语义,可参考 Connector-V2 特性说明。
支持的数据源信息
Sls 连接器对数据源版本要求为Universal(通用)。使用前需通过install-plugin.sh或 Maven 中央仓库获取依赖,Maven 坐标为org.apache.seatunnel:connector-sls。
当前仓库的 plugin-mapping.properties 中完成了插件与模块的注册映射:
seatunnel.source.Sls = connector-sls seatunnel.sink.Sls = connector-sls从连接器的 pom.xml 可以看出,其依赖了阿里云官方日志 SDKcom.aliyun.openservices:aliyun-log(版本0.6.109),以及 SeaTunnel 的connector-common、seatunnel-format-json、seatunnel-format-text等基础模块。
Sink 选项
Sls Sink 的全部配置项定义于 SlsSinkOptions.java 与父类 SlsBaseOptions.java:
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| endpoint | String | 是 | - | 阿里云 SLS 访问地址,例如cn-hangzhou.log.aliyuncs.com或内网访问地址(如cn-hangzhou-intranet.log.aliyuncs.com)。 |
| project | String | 是 | - | 阿里云 SLS Project 名称。 |
| logstore | String | 是 | - | 阿里云 SLS Logstore 名称。 |
| access_key_id | String | 是 | - | 阿里云 AccessKey ID。 |
| access_key_secret | String | 是 | - | 阿里云 AccessKey Secret。 |
| source | String | 否 | SeaTunnel-Source | 写入 SLS log group 的 source 标记。 |
| topic | String | 否 | SeaTunnel-Topic | 写入 SLS log group 的 topic 标记。 |
| log_group_size | Integer | 否 | 100 | SLS log group 写入大小(该选项在源码中定义,文档表中未列出,配置时可按需使用)。 |
源码侧的可选项规则:在 SlsSinkFactory.java 的optionRule()中明确指定了endpoint、project、logstore、access_key_id、access_key_secret五个必填项,source、topic为可选项,这与文档表格完全一致。此外 SlsSinkOptions.java 中还定义了默认值为100的log_group_size选项,用于控制 SLS log group 的写入大小。
底层写入原理(源码解析)
写入口与数据序列化
Sink 的写入口在 SlsSinkWriter.write():
- 调用
SeatunnelRowSerialization.serializeRow(element)将SeaTunnelRow序列化为LogItem列表; - 构造
PutLogsRequest(project, logStore, topic, source, data); - 调用阿里云 SDK 的
client.PutLogs(plr)立即写入 SLS; - 写入失败时记录错误日志并抛出
IOException。
在 SeatunnelRowSerialization.java 中可以看到序列化细节:连接器基于JsonSerializationSchema(来自seatunnel-format-json)将整行数据序列化为 JSON 字符串,再封装为LogContent("content", rowJson)放入LogItem。这印证了文档中「每条数据序列化为 JSON 并写入content字段」的描述。
提交与状态管理
SlsSinkWriter 的prepareCommit()返回空Optional(因为数据在write()阶段已发出)、snapshotState()返回空列表、abortPrepare()为空实现,说明该连接器是典型的「即写即发」型 Sink,不维护跨 checkpoint 的提交状态。关闭 Writer 时调用client.shutdown()释放阿里云 SDK 客户端连接。
连接器装配
SlsSink.java 实现SeaTunnelSink,getPluginName()返回Sls,createWriter()根据CatalogTable推导出的SeaTunnelRowType构造SlsSinkWriter。
注意事项
使用 Sls Sink 连接器前,请务必确认以下几点:
- 配置的 RAM 用户需要有向目标 project 和 logstore 写入日志的权限,否则
PutLogs调用会被 SLS 拒绝。 - sink 在收到数据时立即写入,不提供精确一次提交语义。流处理模式下连接器按行写入;checkpoint 只对下游状态有用,并不能保证 SLS 端的写入语义。
- 每条数据都会被序列化为 JSON,并写入 SLS 日志项
content字段,不会映射到其它日志 key。 - 不要在日志或任务说明里打印
access_key_secret,避免凭据泄露。
任务示例
写入数据到 SLS(批处理)
以下配置使用FakeSource生成 10 条测试数据,通过SlsSink 写入内网地址对应的 SLS Project:
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" } }写入数据到 SLS(流处理)
流处理模式下,连接器会保持 SLS Producer 的连接持续打开,每来一行数据就写入一条。可以配置checkpoint.interval保护下游状态,但需要清楚每条PutLogs调用互相独立,重试只在 Producer 会话内进行,不会跨重启。
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" } }变更日志
关于该连接器各版本的功能变更记录,可查看 connector-sls 变更日志。
进一步探索
- 想了解 SLS 作为数据源的用法,可阅读 Sls Source 连接器文档(源码位于 connector-sls 的 source 包)。
- 连接器的配置项合法性由 SlsFactoryTest.java 中的单测覆盖,验证了 Source 与 Sink 工厂的
optionRule()均能正常构建。 - 若需离线安装插件,可参考仓库中 plugins 目录与
config/plugin_config的插件清单机制。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考