ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

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 连接器支持以下三类运行引擎SparkFlinkSeaTunnel 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名称类型是否必填默认值描述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。sourceString否SeaTunnel-Source写入 SLS log group 的 source 标记。topicString否SeaTunnel-Topic写入 SLS log group 的 topic 标记。log_group_sizeInteger否100SLS 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 实现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 Projectenv { 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),仅供参考
返回列表