news 2026/9/17 12:34:57

Apache SeaTunnel SLS Sink 连接器:将 SeaTunnel 行数据写入阿里云日志服务的实现与配置指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache SeaTunnel SLS Sink 连接器:将 SeaTunnel 行数据写入阿里云日志服务的实现与配置指南

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 iscontent.

即:每一行 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-oncecdctimer flush三项均未勾选,意味着该连接器不提供精确一次提交语义,也不面向 CDC 场景设计。

二、Sink 参数详解

以下是文档中完整的 Sink Options 表格,全部参数在 SlsBaseOptions 与 SlsSinkOptions 中有对应定义:

NameTypeRequiredDefaultDescription
endpointStringYes-Alibaba Cloud SLS endpoint, for examplecn-hangzhou.log.aliyuncs.comor an intranet endpoint.
projectStringYes-Alibaba Cloud SLS project。
logstoreStringYes-Alibaba Cloud SLS logstore。
access_key_idStringYes-Alibaba Cloud AccessKey ID。
access_key_secretStringYes-Alibaba Cloud AccessKey secret。
sourceStringNoSeaTunnel-SourceSource tag written to SLS log groups。
topicStringNoSeaTunnel-TopicTopic tag written to SLS log groups。

从源码可以确认几个细节:

  • 必填性由 SlsSinkFactory 的OptionRule声明:ENDPOINTPROJECTLOGSTOREACCESS_KEY_IDACCESS_KEY_SECRET五项为requiredSOURCETOPICoptional,与文档表格一致。
  • sourcetopic的默认值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 看数据是如何发出的。关键调用链如下:

  1. 构造阶段:每个 writer 实例化一个阿里云 SLS SDK 客户端:
this.client = new Client( pluginConfig.get(SlsSinkOptions.ENDPOINT), pluginConfig.get(SlsSinkOptions.ACCESS_KEY_ID), pluginConfig.get(SlsSinkOptions.ACCESS_KEY_SECRET));

同时从配置读取projectlogStoretopicsource,并创建SeatunnelRowSerialization序列化器。

  1. 写入阶段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使作业失败。

  1. 提交阶段prepareCommit()直接返回Optional.empty(),注释明确写着 "nothing to do, when write function, data had sended"(数据在 write 时已经发送);SlsSinkCommitter 的commit同样不做任何事,snapshotState返回空列表。这与文档 Key Features 中未勾选 exactly-once 是一致的:该连接器不存在两阶段提交协议,checkpoint 仅用于上下游状态,不覆盖 SLS 写入本身

  2. 关闭阶段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 一节给出的四条注意事项,结合源码可以逐条落到实处:

  1. 权限要求:配置的 RAM 用户必须拥有向目标 project 和 logstore 写日志的权限,否则会直接触发write抛出的IOException导致作业失败。
  2. 写入时机与语义:数据在write被调用时立即写出,连接器不提供 exactly-once 提交语义;流模式下逐行 flush,checkpoint 仅保障下游状态,不保障 SLS 写入本身。对重复敏感的场景应在 SLS 消费侧做幂等处理。
  3. 字段映射模型:每行序列化为一个 JSON 对象并整体存放在日志项的content键下,其余行字段不会被拆分成独立的日志键。
  4. 密钥安全:不要在日志或作业描述中打印access_key_secret。由于该值会随客户端构造传入 SLS SDK,建议结合作业配置管理手段(环境变量、密钥托管等)注入。

六、端到端验证参考

仓库中为该连接器提供了 e2e 测试骨架:SlsIT,配套的 sink 配置 sls_sink_to_console.conf 展示了"FakeSource → Sls"的最小结构(endpoint/project/logstore/密钥均为占位符,实际运行需替换为真实凭据)。此外同目录下的sls_source_with_schema_to_console.confsls_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 12:33:32

MySQL Workbench 8.0导入导出实战:从备份到恢复的完整指南

我最早被MySQL Workbench的导入导出功能救场&#xff0c;是在帮人做数据库课程设计的时候。辛辛苦苦在实验室机器上建好的几十张表、一堆视图和存储过程&#xff0c;要拷回宿舍电脑继续调&#xff0c;总不能把整个数据库文件目录打包搬走&#xff0c;更不可能一张表一张表重新敲…

作者头像 李华
网站建设 2026/9/17 12:31:20

2025年ISO9001质量手册编写与内审落地指南:从过程方法到文件化信息

简介&#xff1a;这是一份面向质量管理人员、内审员及企业体系负责人的最新版ISO9000质量管理体系及质量手册文档&#xff0c;旨在帮助组织系统建立、实施并持续改进质量管理体系。资源以单个docx文件呈现&#xff0c;压缩包容量约114KB&#xff0c;内容即完整质量手册正文。手…

作者头像 李华
网站建设 2026/9/17 12:30:18

从T/CPCA 1001-2022解析到术语库:PCB工程师的评审利器

简介&#xff1a;TCPCA 1001-2022《电子电路术语》团体标准正式PDF版&#xff0c;由中国电子电路行业协会发布&#xff0c;面向电子电路设计、制造、测试与维护等环节的工程师、技术管理人员及相关专业师生。标准以统一行业语言为目标&#xff0c;系统界定了基础术语、设计术语…

作者头像 李华
网站建设 2026/9/17 12:28:32

VSCode 里 Git 分支切换与合并实战:冲突、回滚与避坑

上周三晚上我在赶一个需求&#xff0c;feature 分支上改了七八个文件&#xff0c;突然要确认主干上一段历史提交的写法&#xff0c;顺手在 VSCode 底部状态栏点了下分支名切过去。弹出的对话框问我要不要把改动 stash 起来&#xff0c;我当时没细想就点了"是"。第二天…

作者头像 李华