news 2026/9/19 0:35:16

SeaTunnel Sls Sink 连接器实战指南:将数据写入阿里云日志服务 SLS

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SeaTunnel Sls Sink 连接器实战指南:将数据写入阿里云日志服务 SLS

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-commonseatunnel-format-jsonseatunnel-format-text等基础模块。

Sink 选项

Sls Sink 的全部配置项定义于 SlsSinkOptions.java 与父类 SlsBaseOptions.java:

名称类型是否必填默认值描述
endpointString-阿里云 SLS 访问地址,例如cn-hangzhou.log.aliyuncs.com或内网访问地址(如cn-hangzhou-intranet.log.aliyuncs.com)。
projectString-阿里云 SLS Project 名称。
logstoreString-阿里云 SLS Logstore 名称。
access_key_idString-阿里云 AccessKey ID。
access_key_secretString-阿里云 AccessKey Secret。
sourceStringSeaTunnel-Source写入 SLS log group 的 source 标记。
topicStringSeaTunnel-Topic写入 SLS log group 的 topic 标记。
log_group_sizeInteger100SLS log group 写入大小(该选项在源码中定义,文档表中未列出,配置时可按需使用)。

源码侧的可选项规则:在 SlsSinkFactory.java 的optionRule()中明确指定了endpointprojectlogstoreaccess_key_idaccess_key_secret五个必填项,sourcetopic为可选项,这与文档表格完全一致。此外 SlsSinkOptions.java 中还定义了默认值为100log_group_size选项,用于控制 SLS log group 的写入大小。

底层写入原理(源码解析)

写入口与数据序列化

Sink 的写入口在 SlsSinkWriter.write():

  1. 调用SeatunnelRowSerialization.serializeRow(element)SeaTunnelRow序列化为LogItem列表;
  2. 构造PutLogsRequest(project, logStore, topic, source, data)
  3. 调用阿里云 SDK 的client.PutLogs(plr)立即写入 SLS;
  4. 写入失败时记录错误日志并抛出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 实现SeaTunnelSinkgetPluginName()返回SlscreateWriter()根据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),仅供参考

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

VSCode Python开发环境配置与调试实战指南

从第一次摸到VSCode写Python,到真正把它变成主力工具,其实中间隔着一大堆细节问题。你可能已经装好了Python和VSCode,打开编辑器,准备写第一行代码却发现没有代码提示,运行时终端全是英文报错,明明刚pip装完…

作者头像 李华
网站建设 2026/9/19 0:34:17

PCANet结合遮挡定位的人脸识别:原理、实现与调优

简介:《PCANet下的遮挡定位人脸识别算法》是一篇发表在《计算机科学与探索》上的学术论文,面向人脸识别与深度学习研究人员,聚焦自然环境下遮挡导致识别率下降的难题。论文提出将深度学习和特征点遮挡检测相结合的PCANet遮挡定位识别算法&…

作者头像 李华
网站建设 2026/9/19 0:32:35

Agent-Reach:打通工具调用、记忆与A2A的触达链路

一个 Agent 项目从 Demo 走到线上,最常见的死法不是模型不够聪明,而是它"够不着"。你在本地跑一个问答式 agent,它谈吐得体、逻辑清晰;一旦把真实工单系统、数据库、内部接口丢给它,完成率立马掉到三成以下。…

作者头像 李华
网站建设 2026/9/19 0:30:49

SSM农产品供销服务系统:从业务拆解到部署避坑全指南

干过课程设计、毕业设计,或者接手过学长留下的 SSM 老旧项目的朋友,应该都懂这种感受:项目标题写得规规矩矩,叫“SSM292的农产品供销服务系统”,乍一看平平无奇,但真正动手去跑、去改、去部署的时候&#x…

作者头像 李华
网站建设 2026/9/19 0:28:22

OpenCV中文手册不存在?手动生成可搜索本地文档

简介:本资源是一份面向计算机视觉初学者与OpenCV开发者的中文技术手册,系统梳理图像处理核心算法与API用法,助力快速掌握OpenCV 1.x/2.x经典函数体系。手册共10大章节,涵盖梯度与边缘检测(Sobel、Laplace、Canny&#…

作者头像 李华
网站建设 2026/9/19 0:28:15

智慧社区规划方案PPT编制指南:架构、场景与评审要点

简介:面向智慧社区建设方案的74页PPT,适合社区管理者、智能化集成商及方案汇报人参考。内容以“智慧、互联、共享、融合”为主线,从需求分析与总体规划切入,系统梳理顶层设计、基础系统建设、智能化系统建设,重点落至A…

作者头像 李华