尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于源码理解 SeaTunnel AmazonSqs 源连接器:从队列读取到 Schema 解析的完整实战指南

基于源码理解 SeaTunnel AmazonSqs 源连接器:从队列读取到 Schema 解析的完整实战指南 基于源码理解 SeaTunnel AmazonSqs 源连接器从队列读取到 Schema 解析的完整实战指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunneloutput_articleSeaTunnel 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 等周期调度重复触发这种有界读取。支持的引擎SparkFlinkSeaTunnel Zeta主要特性批处理流处理精确一次列投影并行度支持用户自定义分片从源码层面看AmazonSqsSource实现了SupportColumnProjection接口AmazonSqsSource.java并且其getBoundedness()方法返回Boundedness.BOUNDED同文件 L68-L71与文档声明的「批处理、列投影支持流处理、并行度不支持」完全一致。所谓「不支持并行度」指的是该源不会对队列做分片并发消费同一时刻只有一个 reader 执行一次 receive 调用。源选项Source Options名称类型是否必填默认值描述urlString是-要读取的完整 SQS 队列 URL例如https://sqs.us-east-1.amazonaws.com/123456789012/source_queue。regionString是-SQS 队列所在的 AWS 区域例如us-east-1。schemaConfig是-消息体结构包含字段名和字段类型。更多说明请参考 Schema 特性。access_key_idString否-AWS access key ID。和secret_access_key一起配置时使用静态凭证两者都不配置时使用 AWS 默认凭证链。secret_access_keyString否-AWS secret access key。和access_key_id一起配置时使用静态凭证。formatString否json消息体格式。支持json、text、canal_json、debezium_json。field_delimiterString否,当format text时使用的字段分隔符。ignore_parse_errorsBoolean否false是否跳过无法解析的消息并继续处理而不是让本次轮询失败。delete_messageBoolean否false读取并成功解析消息后是否从队列中删除该消息。message_group_idString否-为兼容保留的消息分组 ID 选项普通 SQS 读取不需要配置。debezium_record_include_schemaBoolean否trueDebezium JSON 消息是否包含 schema。仅在format debezium_json时使用。common-options否-源插件通用参数详见 Source Common Options。url可以指向 AWS SQS也可以指向兼容 SQS 的本地服务例如http://sqs-host:4566/000000000000/source_queueLocalStack 默认端口即 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.JSONfield_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.javajson把每条消息体按 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则使用这对静态凭证构造StaticCredentialsProviderAwsBasicCredentials创建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构造对应的DeserializationSchemaJSON / Text / Canal / DebeziumAmazonSqsSource.createReader()返回单 split 的AmazonSqsSourceReaderreaderopen()时按认证方式创建SqsClienturl同时作为端点endpointOverride与队列 URLqueueUrl使用pollNext()执行一次receiveMessagemaxNumberOfMessages10、waitTimeSeconds10对返回的每条消息调用反序列化器解析出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),仅供参考
返回列表