尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

SeaTunnel AmazonSqs Sink 连接器:把每一行数据作为一条 SQS 消息写入消息队列

发布时间:2026/9/16 14:40:04

资讯中心
01
ARTICLE

SeaTunnel AmazonSqs Sink 连接器:把每一行数据作为一条 SQS 消息写入消息队列

SeaTunnel AmazonSqs Sink 连接器:把每一行数据作为一条 SQS 消息写入消息队列
SeaTunnel AmazonSqs Sink 连接器把每一行数据作为一条 SQS 消息写入消息队列【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 仓库中 AmazonSQS 官方 Sink 文档 展开系统讲解 AmazonSqs sink 连接器的定位、全部配置项、四种消息体序列化格式、AWS 凭据解析顺序与本地兼容服务LocalStack 等的接入方式并结合 sink 写入器实现、工厂类 与 E2E 集成测试 佐证其底层行为。读完后你可以直接编写可运行的 HOCON 作业配置把 SeaTunnel 作业的输出以 JSON、文本、Canal JSON 或 Debezium JSON 消息体形式写入 Amazon SQS 或 SQS 兼容队列并理解每条数据“一行一消息”的发送语义。1. 连接器定位与能力边界Amazon SQS sink 连接器将每一条进入的 SeaTunnel 数据行SeaTunnelRow序列化后作为一条 Amazon SQS 消息的 body 发送到指定的队列 URL。序列化方式由format选项决定默认是 JSON。该连接器支持三大执行引擎SparkFlinkSeaTunnel Zeta从官方文档的 Key Features 清单看它支持的能力边界是能力是否支持batch批模式支持stream流模式支持exactly-once精确一次不支持cdc不支持多表写入不支持定时 flush不支持各能力项的含义详见 connector-v2-features。结合这一能力边界该 sink 更适合“行级转发、下游最终消费”的场景如跨队列复制、为下游系统投递变更消息而不是需要事务性提交的强一致场景。2. Sink 配置项详解完整的选项定义位于 AmazonSqsBaseOptionssink 专用选项类 AmazonSqsSinkOptions 直接继承它source 与 sink 共用同一套基础选项。NameTypeRequiredDefaultDescriptionurlStringYes-目标 SQS 队列的完整 URL例如https://sqs.us-east-1.amazonaws.com/123456789012/sink_queueregionStringYes-SQS 队列所在的 AWS region例如us-east-1access_key_idStringNo-AWS access key ID。与secret_access_key一起配置以使用静态凭据两者都不配置时走 AWS 默认凭据链secret_access_keyStringNo-AWS secret access key需与access_key_id配套配置formatStringNojson消息体格式。可选json、text、canal_json、debezium_jsonfield_delimiterStringNo,format text时使用的字段分隔符common-options-No-Sink 插件通用参数见 Sink Common Options关于url的一个重要事实它既可以指向 AWS 正式 SQS也可以指向任何 SQS 兼容的本地服务例如http://sqs-host:4566/000000000000/sink_queue。从源码看url在写入器构造时被直接用于SqsClient.builder().endpointOverride(URI.create(url))见 AmazonSqsSinkWriter 构造逻辑因此只要客户端能访问该端点即可工作。region虽然对本地兼容服务“语义上无意义”但源码中注释明确指出它是 AWS SDK 客户端构建时的必填校验项The region is meaningless for local Sqs but required for client builder validation所以对本地服务也必须填写一个合法的 region 值例如us-east-1。2.1 配置校验规则AmazonSqsSinkFactory 通过OptionRule声明了校验规则url和region为必填且必须满足notBlank非空且非纯空白access_key_id、secret_access_key、format、field_delimiter为可选。工厂标识符factoryIdentifier为AmazonSqs这也是 HOCON 配置中 sink 块要使用的名字。单测 AmazonSqsSinkFactoryTest 逐条验证了这些规则缺失url、缺失region、空字符串、纯空白值均会抛出OptionValidationException而合法的urlregion组合可以正常通过。这说明配置阶段就能拦截“只填了 url 没填 region”这类常见错误作业不会在运行期才失败。3. 消息体格式format与序列化format是一个枚举选项取值由 MessageFormat 枚举限定JSON、TEXT、CANAL_JSON、DEBEZIUM_JSON默认JSON。写入器中的createSerializationSchema方法按该枚举分发到对应的序列化器见 AmazonSqsSinkWriterjson使用JsonSerializationSchema每行序列化为一个 JSON 对象text使用TextSerializationSchema将行字段按field_delimiter拼接默认逗号,常量DEFAULT_FIELD_DELIMITER定义在 AmazonSqsBaseOptionscanal_json使用CanalJsonSerializationSchema输出 Canal JSON 变更事件详见 Canal JSONdebezium_json使用DebeziumJsonSerializationSchema输出 Debezium JSON 变更事件详见 Debezium JSON遇到枚举之外的值会抛出SeaTunnelJsonFormatException。官方文档同时明确了几个重要的行为约束配置时应注意sink只发送消息体message body不暴露 SQS 消息属性message attributes、延迟秒数delay seconds、去重 IDdeduplication ID或消息组 IDmessage group ID选项access_key_id与secret_access_key可选但使用静态 AWS 凭据时必须成对配置sink每行数据发送一条独立 SQS 消息不会把多行数据批量合并进一次 SQS 请求。4. 认证与凭据解析顺序连接器按以下顺序解析 AWS 凭据若同时配置了access_key_id和secret_access_key使用这对静态凭据否则回落到AWS 默认凭据链DefaultCredentialsProvider即环境变量、实例配置文件、EC2 实例角色instance profile等标准来源。从源码可以印证这一顺序AmazonSqsSinkWriter 构造函数 中当两个凭据选项都存在时构建StaticCredentialsProvider否则构建DefaultCredentialsProvider.create()两条分支都会设置endpointOverride来自url和region。本地测试场景若用 LocalStack 或 ElasticMQ 等 SQS 兼容服务做本地验证将url指向本地端点如http://sqs-host:4566/...并填入任意非空的access_key_id/secret_access_key即可。原因是这类 SQS 兼容测试服务通常不校验请求的 SigV4 签名任意静态凭据对都会被接受。5. 作业配置示例以下四个示例完整继承自官方文档覆盖本地队列复制、JSON 写入、自定义分隔符文本写入与 Canal JSON 写入四类典型场景。5.1 在 SQS 兼容的本地队列之间复制消息env { parallelism 1 job.mode BATCH } source { AmazonSqs { url http://sqs-host:4566/000000000000/source_queue access_key_id 1234 secret_access_key abcd region us-east-1 schema { fields { name string } } } } sink { AmazonSqs { url http://sqs-host:4566/000000000000/sink_queue access_key_id 1234 secret_access_key abcd region us-east-1 } }注意 source 与 sink 使用同一个连接器插件名AmazonSqssource 侧额外需要schema来声明读取的消息结构source 侧完整选项见同目录下的 source 文档。5.2 写入 JSON 消息默认格式env { parallelism 1 job.mode BATCH } source { FakeSource { row.num 1 schema { fields { name string } } rows [ { kind INSERT fields [test_name] } ] } } sink { AmazonSqs { url https://sqs.us-east-1.amazonaws.com/123456789012/sink_queue region us-east-1 access_key_id AKIA... secret_access_key SECRET... } }该配置使用FakeSource产生一行name test_name的数据默认format json最终队列中应收到消息体{name:test_name}。5.3 使用自定义分隔符写入文本消息source { FakeSource { schema { fields { artist string album string release_year int } } } } sink { AmazonSqs { url https://sqs.us-east-1.amazonaws.com/123456789012/sink_queue region us-east-1 format text field_delimiter | } }此时每行消息体形如artist|album|release_year字段按|拼接。5.4 写入 Canal JSON 消息将format设为canal_json每条 SeaTunnel 行会被序列化为一条 Canal JSON 变更事件。适用于下游是 Canal JSON 消费方的场景如 Canal → Kafka 桥接、Canal 兼容的 BigQuery 加载链路sink { AmazonSqs { url https://sqs.us-east-1.amazonaws.com/123456789012/sink_queue region us-east-1 format canal_json } }6. 写入链路源码走读一行数据如何变成一条 SQS 消息从源码结构看整个写入链路分三层均位于seatunnel-connectors-v2/connector-amazonsqs模块工厂层AmazonSqsSinkFactory 通过AutoService(Factory.class)注册标识符AmazonSqs。它做两件事用OptionRule声明配置校验第 2.1 节已述以及createSink时返回一个 lambda 包装的 AmazonSqsSinkSink 层AmazonSqsSink继承AbstractSimpleSink在构造时从CatalogTable提取物理行类型typeInfo并在createWriter中创建AmazonSqsSinkWriterWriter 层AmazonSqsSinkWriter 的write(SeaTunnelRow row)方法就是“一行一消息”语义的直接实现byte[] bytes serializationSchema.serialize(row); // 1. 按 format 序列化 String messageBody new String(bytes, StandardCharsets.UTF_8); SendMessageRequest sendMessageRequest SendMessageRequest.builder() .queueUrl(pluginConfig.get(URL)) // 2. 目标队列来自 url .messageBody(messageBody) // 3. 只设置消息体 .build(); sqsClient.sendMessage(sendMessageRequest); // 4. 逐行发送可以看到请求中只设置了queueUrl和messageBody两个字段——这正是文档中“只发送消息体、不暴露消息属性/延迟/去重 ID/消息组 ID”这一限制的实现根源。close()时关闭SqsClient释放 SDK 线程池等资源。7. 端到端验证LocalStack 集成测试E2E 测试 AmazonsqsIT 完整覆盖了第 5.1 节“本地队列复制”场景可以照它的方式自行验证使用 Testcontainers 启动localstack/localstack:3.7容器仅启用 SQS 服务网络别名sqs-host端口4566环境变量AWS_DEFAULT_REGIONus-east-1、AWS_ACCESS_KEY_ID1234、AWS_SECRET_ACCESS_KEYabcd——与文档示例中的sqs-host:4566、1234/abcd凭据一一对应测试先创建source_queue与sink_queue两个队列向 source 队列发送一条{name:test_name}消息执行作业配置文件/amazonsqsIT_source_to_sink.conf即 source → sink 的 AmazonSqs 复制作业断言进程退出码为 0再从 sink 队列收消息断言恰好收到 1 条且 body 与源消息完全一致。测试注释特别提醒SQS 消息被接收后即进入不可见状态不要重复调用receiveMessage。这套测试同时印证了两点其一文档中“SQS 兼容服务不校验 SigV4、任意静态凭据可用”的说法在 LocalStack 上成立其二region对本地服务只是占位测试统一填us-east-1。8. 适用前提与限制小结基于当前仓库内容使用该连接器需注意以下前提版本演进连接器标识符在 2.3.4 版本从amazonsqs更名为AmazonSqs且 aws-sdk-v2 已在后续版本统一为 2.31.30变更记录见 connector-amazonsqs changelog旧配置中的插件名需按新标识符调整持久化语义当前实现不支持 exactly-once 与 CDC 语义作业失败重跑是否产生重复消息取决于上游与下游的幂等设计需要在关键链路上自行评估消息能力上限无法设置延迟发送、去重 ID 与消息组 IDFIFO 队列的高级特性在本 sink 中无法使用本地服务 region即使是 LocalStack/ElasticMQregion也必须填写这是 SDK 客户端构建的强制校验项而非业务必需批量行为单行单消息、逐行sendMessage高吞吐场景下客户端请求量与数据行数成正比可结合env.parallelism调节并发。按照上述配置项与示例你可以把该 sink 直接嵌入 SeaTunnel 作业真实 AWS 场景下填写队列 URL、region 与静态凭据或依赖默认凭据链本地验证场景下指向 LocalStack 端点并填入任意占位凭据即可让数据以 JSON/文本/Canal JSON/Debezium JSON 四种消息体格式进入目标队列。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。