
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 ServiceSLS。官方文档给出的核心描述是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 ListLogItem serializeRow(SeaTunnelRow row) { ListLogItem logGroup new ArrayListLogItem(); 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 字段提取而不是期望连接器逐列写入。官方文档明确声明该连接器支持三种执行引擎SparkFlinkSeaTunnel Zeta同时在 Key Features 一栏中exactly-once、cdc、timer flush三项均未勾选意味着该连接器不提供精确一次提交语义也不面向 CDC 场景设计。二、Sink 参数详解以下是文档中完整的 Sink Options 表格全部参数在 SlsBaseOptions 与 SlsSinkOptions 中有对应定义NameTypeRequiredDefaultDescriptionendpointStringYes-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声明ENDPOINT、PROJECT、LOGSTORE、ACCESS_KEY_ID、ACCESS_KEY_SECRET五项为requiredSOURCE、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 { ListLogItem 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实现了SeaTunnelSinkSeaTunnelRow, SlsSinkState, SlsCommitInfo, SlsAggregatedCommitInfocreateWriter时为每个 writer 传入空的slsStates列表——因为不存在需要恢复的写入状态。四、完整作业配置示例以下两个示例完整继承自官方文档可直接作为生产配置的模板。1. 批模式写入 SLSBatchenv { 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 } }批模式示例中使用的是内网 endpointcn-hangzhou-intranet.log.aliyuncs.com适用于作业与 SLS 同区域部署、走内网链路降低延迟与流量的场景。2. 流模式写入 SLSStreaming文档说明流模式下连接器保持 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 提交语义流模式下逐行 flushcheckpoint 仅保障下游状态不保障 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),仅供参考