基于源码理解 SeaTunnel AmazonSqs 源连接器:从队列读取到 Schema 解析的完整实战指南
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
<output_article>
SeaTunnel AmazonSqs 源连接器实战指南:有界读取 SQS 队列消息并完成 Schema 解析
本文以 SeaTunnel 官方文档中 AmazonSqs 源连接器的说明为主体,结合仓库内该连接器的源码实现与单元测试,系统讲解如何从 Amazon SQS 队列(含 LocalStack、ElasticMQ 等 SQS 兼容本地服务)读取消息、按
format与schema解析为 SeaTunnel 行数据,并给出完整可运行的配置示例与底层原理分析。读完本文,你将掌握该连接器的全部配置项、有界读取的行为边界、凭证认证方式,以及delete_message、ignore_parse_errors等关键开关在实际任务中的正确取舍。
连接器概述
Amazon SQS 源连接器(插件名AmazonSqs)用于从一个 Amazon SQS 队列 URL 读取消息。连接器会按照format和schema解析每条消息的消息体,然后输出为 SeaTunnel 行数据(SeaTunnelRow)。
该连接器是一个单 reader 源(从源码看,它继承自AbstractSingleSplitSource,见 AmazonSqsSource.java),处理完本次接收到的消息后任务即结束,因此天然适合有界(BATCH)读取场景,而非持续订阅式的流式消费。
每次 receive 请求最多从 SQS 拉取 10 条消息。如果队列里还有更多消息,需要再次运行任务,或者使用上游调度方式(如 Airflow、Azkaban 等周期调度)重复触发这种有界读取。
支持的引擎
Spark
Flink
SeaTunnel Zeta
主要特性
- 批处理
- 流处理
- 精确一次
- 列投影
- 并行度
- 支持用户自定义分片
从源码层面看,AmazonSqsSource实现了SupportColumnProjection接口(AmazonSqsSource.java),并且其getBoundedness()方法返回Boundedness.BOUNDED(同文件 L68-L71),与文档声明的「批处理、列投影支持,流处理、并行度不支持」完全一致。所谓「不支持并行度」,指的是该源不会对队列做分片并发消费,同一时刻只有一个 reader 执行一次 receive 调用。
源选项(Source Options)
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | 要读取的完整 SQS 队列 URL,例如https://sqs.us-east-1.amazonaws.com/123456789012/source_queue。 |
| region | String | 是 | - | SQS 队列所在的 AWS 区域,例如us-east-1。 |
| schema | Config | 是 | - | 消息体结构,包含字段名和字段类型。更多说明请参考 Schema 特性。 |
| access_key_id | String | 否 | - | AWS access key ID。和secret_access_key一起配置时使用静态凭证;两者都不配置时使用 AWS 默认凭证链。 |
| secret_access_key | String | 否 | - | AWS secret access key。和access_key_id一起配置时使用静态凭证。 |
| format | String | 否 | json | 消息体格式。支持json、text、canal_json、debezium_json。 |
| field_delimiter | String | 否 | , | 当format = text时使用的字段分隔符。 |
| ignore_parse_errors | Boolean | 否 | false | 是否跳过无法解析的消息并继续处理,而不是让本次轮询失败。 |
| delete_message | Boolean | 否 | false | 读取并成功解析消息后,是否从队列中删除该消息。 |
| message_group_id | String | 否 | - | 为兼容保留的消息分组 ID 选项,普通 SQS 读取不需要配置。 |
| debezium_record_include_schema | Boolean | 否 | true | Debezium JSON 消息是否包含 schema。仅在format = debezium_json时使用。 |
| common-options | 否 | - | 源插件通用参数,详见 Source Common Options。 |
url可以指向 AWS SQS,也可以指向兼容 SQS 的本地服务,例如http://sqs-host:4566/000000000000/source_queue(LocalStack 默认端口即 4566)。
参数定义在源码中的位置
上述选项并非文档凭空罗列,而是由连接器工厂通过OptionRule声明,并由配置类在运行时解析:
- 必填项
url、region、schema以及所有可选项,声明于 AmazonSqsSourceFactory.optionRule(),其中url和region还附加了notBlank约束,配置为空字符串会在构建阶段直接校验失败; url、region、access_key_id、secret_access_key、format、field_delimiter定义于 AmazonSqsBaseOptions.java,其中format的默认值是MessageFormat.JSON,field_delimiter未设置默认值,实际取默认逗号,(见同文件常量DEFAULT_FIELD_DELIMITER = ",");delete_message(默认false)、ignore_parse_errors(默认false)、message_group_id(无默认值)、debezium_record_include_schema(默认true)定义于 AmazonSqsSourceOptions.java;- 运行期由 AmazonSqsSourceConfig.java 统一读取并封装成配置对象,其中
schema会被转换为 typesafe Config 供反序列化器使用。
格式说明(Format)
连接器支持四种消息体格式,枚举定义见 MessageFormat.java:
json:把每条消息体按 JSON 对象解析,并要求字段能对应到schema。这是默认格式,工厂会构造JsonDeserializationSchema。text:按field_delimiter切分消息体,并按schema中字段顺序映射。工厂会构造TextDeserializationSchema;未显式配置field_delimiter时使用默认分隔符,(源码位于 AmazonSqsSourceFactory.setDeserialization())。canal_json:读取 Canal JSON 消息,详见 Canal JSON。debezium_json:读取 Debezium JSON 消息,详见 Debezium JSON。
多行输出与删除时机
- 一条
canal_json或debezium_json消息可能产生多行。例如,更新事件会产生更新前(UPDATE_BEFORE)和更新后(UPDATE_AFTER)两行。这一行为在源码中有明确实现:反序列化器对这两种格式调用deserializeMultipleRows,通过内部BufferingCollector收集多行(见 AmazonSqsDeserializer.java),并被 AmazonSqsSourceReaderTest.java 中的shouldDeserializeCanalUpdateAsTwoRows、shouldDeserializeDebeziumUpdateAsTwoRows等测试用例直接验证。 ignore_parse_errors = false会让本次轮询失败并保留无法解析的消息。设置为true时,源连接器会跳过该消息并继续处理本批次中的其他消息。测试shouldSkipFailedMessageAndContinueWhenParseErrorsIgnored验证了「跳过坏消息、继续消费后续消息」的行为。- 当
ignore_parse_errors和delete_message都为true时,跳过的消息会从 SQS 中删除。如果需要保留这些消息以便重新投递,请保持delete_message = false。 - 对于产生多行的消息,只有在所有行都成功收集后才会删除消息。收集失败时,SQS 消息会保留以便重新投递。对应的测试用例包括
shouldNotDeleteMessageWhenCollectionFails、shouldNotDeleteCanalMessageWhenSecondCollectionFails。 delete_message = true会删除已经消费的 SQS 消息。如果只是检查或复制消息,建议保留默认值false。- 该源连接器只执行一次 receive 请求,最多读取 10 条消息,然后结束这个有界任务。这一点在 AmazonSqsSourceReader.pollNext() 中体现:
maxNumberOfMessages(10)拉取消息,处理结束后调用context.signalNoMoreElement()通知引擎本次读取完成。
认证方式
连接器按以下顺序解析 AWS 凭证(对应 AmazonSqsSourceReader.open() 中的分支逻辑):
- 如果同时配置了
access_key_id和secret_access_key,则使用这对静态凭证,构造StaticCredentialsProvider+AwsBasicCredentials创建SqsClient; - 否则,回退到 AWS 默认凭证链(环境变量、实例角色等),使用
DefaultCredentialsProvider创建SqsClient。
针对 LocalStack、ElasticMQ 等 SQS 兼容本地服务进行测试时,把url指向本地端点(例如http://sqs-host:4566/...),并提供任意非空的access_key_id/secret_access_key即可。SQS 兼容的测试服务通常不会校验请求里的 SigV4 签名,因此任意一对静态凭证都会被接受。另外从源码可以看到,region对本地 SQS 服务本身没有实际意义,但 AWS SDK 的客户端构建器强制要求提供(源码注释原文:"The region is meaningless for local Sqs but required for client builder validation")。
数据读取的底层流程
将上述源码串起来,一次有界读取的完整调用链是:
- 引擎启动后,
AmazonSqsSourceFactory.createSource()解析配置并依据format构造对应的DeserializationSchema(JSON / Text / Canal / Debezium); AmazonSqsSource.createReader()返回单 split 的AmazonSqsSourceReader;- reader
open()时按认证方式创建SqsClient,url同时作为端点(endpointOverride)与队列 URL(queueUrl)使用; pollNext()执行一次receiveMessage(maxNumberOfMessages=10、waitTimeSeconds=10),对返回的每条消息调用反序列化器解析出SeaTunnelRow并output.collect;- 若开启
delete_message,用消息的receiptHandle调用deleteMessage删除已成功处理的消息; - 全部处理完后调用
signalNoMoreElement()结束任务。
任务示例
在本地兼容队列之间复制消息
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" } }该示例演示了最典型的场景:用同一套静态凭证同时配置源与目标队列,把本地队列(如 LocalStack)中的消息搬运到另一个队列。由于delete_message未配置(默认false),源队列中的消息在复制后仍会保留,这与文档「复制消息时建议保持delete_message = false」的建议一致。
读取 JSON 消息
env { parallelism = 1 job.mode = "BATCH" } source { AmazonSqs { url = "https://sqs.us-east-1.amazonaws.com/123456789012/source_queue" region = "us-east-1" access_key_id = "AKIA..." secret_access_key = "SECRET..." schema = { fields { name = string } } } } sink { Console {} }format默认即为json,因此只需给出schema定义字段结构。每条消息体必须是 JSON 对象,且字段需能对得上schema,例如{"name": "seatunnel"}。读取结果可直接输出到 Console sink 进行验证。
使用自定义分隔符读取文本消息
source { AmazonSqs { url = "https://sqs.us-east-1.amazonaws.com/123456789012/source_queue" region = "us-east-1" format = text field_delimiter = "#" delete_message = true schema = { fields { artist = string album = string release_year = int } } } } sink { Console {} }当消息体是纯文本行(例如Beyond#海阔天空#1993)时,使用format = text,字段按schema中的声明顺序依次映射,release_year会被解析为int类型。delete_message = true表示处理成功后从队列删除,适合「消费即处理完毕」的任务;若想保留原始消息做后续审计,应保持默认的false。
读取 Debezium JSON 消息
当上游系统(例如 Debezium 或其他 CDC 源)以 Debezium 信封形式发布变更事件时,把format设为debezium_json,并通过debezium_record_include_schema控制消息中是否包含 schema 字段。
source { AmazonSqs { url = "https://sqs.us-east-1.amazonaws.com/123456789012/cdc_events" region = "us-east-1" format = debezium_json debezium_record_include_schema = true schema = { fields { id = bigint name = string score = double } } } }debezium_record_include_schema默认值为true(见 AmazonSqsSourceOptions.java),工厂在构造DebeziumJsonDeserializationSchema时读取该开关(AmazonSqsSourceFactory.java)。一条 UPDATE 事件会被展开为UPDATE_BEFORE与UPDATE_AFTER两行,例如测试中的{"schema":{},"payload":{"before":...,"after":...,"op":"u"}}会产出before、after两行。
注意事项与最佳实践
- 有界读取的边界:每次任务只执行一次 receive,最多 10 条消息。若要清空整个队列,需要外部调度反复触发;若追求流式持续消费,该连接器(当前版本)并不适用,可考虑 Kinesis、Kafka 等流式源。
ignore_parse_errors与delete_message的组合语义:两者都为true时,解析失败的消息也会被删除,适合「队列中允许存在脏数据且无需重放」的场景;反之希望坏消息保留在队列中用于排查或重投递时,务必设置delete_message = false。注意从源码看,即使ignore_parse_errors = true,非 JSON 解析类的运行时异常(如不支持的格式)仍会向上抛出(AmazonSqsDeserializer.java),不会被静默吞掉。- 多行消息的删除保证:Canal/Debezium 消息只有在所有产出行都成功收集后才会被删除;任一行收集失败都会保留整条 SQS 消息用于重新投递,避免数据丢失。
- 本地联调:推荐使用 LocalStack 或 ElasticMQ 通过
url指向本地端点,配合任意非空静态凭证即可完成端到端验证,无需真实 AWS 账户;region虽对本地服务无实际意义,但必须填写(客户端构建器强制校验)。 - 凭证安全:生产环境建议优先依赖 AWS 默认凭证链(环境变量、实例角色、profile 等),避免在配置文件中明文写入密钥。
变更日志
关于该连接器的历史变更记录,请参考 connector-amazonsqs 变更日志(该文件由连接器模块的 ChangeLog 组件动态渲染)。源码、单元测试与文档主体位于 connector-amazonsqs 模块,如需深入阅读可重点查看上述引用到的工厂类、配置类、reader 与反序列化实现。 </output_article>
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考