SeaTunnel AmazonDynamoDB Source Connector 使用指南:基于 Scan 的批量数据读取与并行分段实现解析 SeaTunnel AmazonDynamoDB Source Connector 使用指南基于 Scan 的批量数据读取与并行分段实现解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelAmazon DynamoDB Source Connector 是 SeaTunnel 提供的批量数据源插件通过 DynamoDB Scan 请求读取表内全量快照数据支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎。本文基于官方文档与仓库源码完整梳理该连接器的配置参数、Schema 定义、类型映射与并行扫描原理帮助你快速上手并理解其底层实现。概述DynamoDB 全表快照读取Amazon DynamoDB 是一款键值Key-Value与文档Document数据库它不像关系型数据库那样暴露标准化的字段类型元信息。因此SeaTunnel 的 AmazonDynamoDB Source Connector 无法自动推断完整的 SeaTunnel Schema必须由用户在配置中显式声明每一个待读取字段及其类型。该连接器通过DynamoDB Scan 请求读取表中的当前数据快照它只读取当前时刻的表数据不会读取 DynamoDB Streams也不消费 CDCChange Data Capture变更事件它是批量batch数据源作业执行完一次全量扫描后即结束它支持并行扫描Parallel Scan将整张表按逻辑分段segment切分由多个 Reader 并发读取不同分段。连接器的完整实现位于仓库 seatunnel-connectors-v2/connector-amazondynamodb 目录下官方变更记录见 connector-amazondynamodb Changelog。支持的引擎引擎支持情况Spark✅Flink✅SeaTunnel Zeta✅三种引擎共享同一套 Connector 实现与配置语义本文示例可直接在 SeaTunnel Zeta 引擎下运行。关键特性特性支持说明批量模式batch✅一次全表 Scan 后结束流模式stream❌不读取 DynamoDB Streams / CDC精确一次exactly-once❌—列投影column projection❌通过 schema 显式选取字段并行度parallelism✅支持多并行度并发扫描用户自定义分片user-defined split❌分段数量由parallel_scan_threads决定特性定义的通用说明可参考 Connector V2 Features。在源码层面AmazonDynamoDBSource 同时实现了SupportParallelism与SupportColumnProjection两个接口其中getBoundedness()返回Boundedness.BOUNDED从实现上印证了它是有界的批量 Source。工作原理Scan、分段与并行理解该连接器的关键在于它的读路径设计整个流程由三个核心类协作完成。1. 分段发现SplitEnumerator 按线程数切分表作业启动时AmazonDynamoDBSourceSplitEnumerator 的discoverySplits()会根据parallel_scan_threads逻辑分段数量为整张表创建对应数量的AmazonDynamoDBSourceSplitint totalSegments amazonDynamoDBConfig.parallelScanThreads; int itemLimit amazonDynamoDBConfig.scanItemLimit; for (int i 0; i totalSegments; i) { AmazonDynamoDBSourceSplit split new AmazonDynamoDBSourceSplit(i, totalSegments, itemLimit); allSplit.add(split); }每个 Split 携带三个关键信息splitId分段编号从 0 开始、totalSegments总分段数和itemLimit单次 Scan 请求返回的最大条数。随后通过getSplitOwnerassignCount % readerCount把分段轮询分配给各个并行 Reader实现负载均衡。2. 实际扫描Reader 基于 Segment 构造 ScanRequest每个 Reader 拿到分段后在 AmazonDynamoDBSourceReader 中基于 AWS SDK v2 构造ScanRequest并利用scanPaginator()自动分页拉取全部数据ScanRequest scanRequest ScanRequest.builder() .tableName(amazondynamodbConfig.getTable()) .limit(split.getItemCount()) .segment(split.getSplitId()) .totalSegments(split.getTotalSegments()) .build(); scan dynamoDbClient.scanPaginator(scanRequest); do { scan.items().forEach(item - { output.collect(seaTunnelRowDeserializer.deserialize(item)); }); } while (scan.iterator().hasNext() !noMoreSplit);这里segment与totalSegments正是 DynamoDB Parallel Scan 的原生参数意味着多个 Reader 可以同时对不同分段发起 Scan互不干扰。当所有分段读取完毕且收到noMoreSplit事件后Reader 会调用context.signalNoMoreElement()通知引擎数据读取结束。3. 数据转换Deserializer 完成 DynamoDB → SeaTunnel 类型映射每条 DynamoDB ItemMapString, AttributeValue通过 DefaultSeaTunnelRowDeserializer 转换为SeaTunnelRow。转换完全依据用户在schema.fields中声明的 SeaTunnel 类型逐字段进行若字段在 Item 中缺失则对应值为null。连接器选项详解Source 插件的全部选项定义在 AmazonDynamoDBSourceOptions 与 AmazonDynamoDBBaseOptions 中汇总如下名称类型必填默认值说明urlstring是-DynamoDB 服务端点 URLregionstring是-DynamoDB 服务所在 AWS 区域access_key_idstring是-AWS 访问密钥 IDsecret_access_keystring是-AWS 访问密钥 Secrettablestring是-要扫描的 DynamoDB 表名schemaconfig是-从 DynamoDB Item 中读取的 SeaTunnel 字段定义scan_item_limitint否1每次 Scan 请求返回的最大 Item 数parallel_scan_threadsint否2并行扫描的逻辑分段数量common-optionsobject否-Source 插件通用参数从 AmazonDynamoDBSourceFactory 的OptionRule可以看到url、region、access_key_id、secret_access_key、table、schema六项为必填scan_item_limit与parallel_scan_threads为可选。上述默认值均来自源码中的Option定义scan_item_limit默认 1、parallel_scan_threads默认 2。url [string]DynamoDB 服务端点 URL。连接线上服务时使用 AWS 官方端点例如url https://dynamodb.us-east-1.amazonaws.com本地联调使用 DynamoDB Local 时配置本地端点即可url http://127.0.0.1:8000region [string]DynamoDB 服务所在 AWS 区域例如us-east-1。在 AmazonDynamoDBSourceReader#open() 中可以看到region 会通过Region.of()传入DynamoDbClient构建器——即便连接 DynamoDB Localregion 也是客户端构建校验所必需的对本地服务而言无实际意义但不可省略。access_key_id [string] / secret_access_key [string]连接 DynamoDB 所用的 AWS 访问凭据。Reader 中通过StaticCredentialsProvider与AwsBasicCredentials将二者显式注入客户端。该连接器必须显式提供这两个参数使用 DynamoDB Local 时填入本地服务认可的任何占位值即可例如dummy-key/dummy-secret。table [string]要扫描的 DynamoDB 表名将作为ScanRequest.tableName传入。schema [config]定义从 DynamoDB Item 中读取的 SeaTunnel 字段。由于 DynamoDB 不暴露完整的字段类型信息必须在此列出所有需要读取的字段未列出的字段将被忽略。配置片段在 AmazonDynamoDBConfig 中通过ConnectorCommonOptions.SCHEMA读取并转换为 Typesafe Config。示例schema { fields { id string c_map mapstring, smallint c_array arraytinyint c_string string c_boolean boolean c_int int c_bigint bigint c_float float c_double double c_decimal decimal(2, 1) c_bytes bytes c_date date c_timestamp timestamp } }完整的 Schema 语法说明请参考 Schema Feature。scan_item_limit [int]每次 DynamoDB Scan 请求返回的最大 Item 数对应ScanRequest.limit。注意它是单次请求的分页大小而不是整个作业的总行数上限。值越大需要的请求次数越少但单次读取批次占用的内存越大值越小请求越轻量但总请求次数增多网络往返开销上升。parallel_scan_threads [int]并行扫描使用的逻辑分段数量决定整张表被拆成多少个 segment 并发扫描。该值应与作业并行度env.parallelism/ source 的parallelism以及表大小对齐小表保持默认值 2 即可大表应同时调大env.parallelism、sourceparallelism与parallel_scan_threads让多个 Reader 各自扫描不同分段充分发挥 DynamoDB Parallel Scan 的吞吐能力。common optionsSource 插件通用参数如parallelism、result_table_name等详见 Source Common Options。数据类型映射DynamoDB 使用自身的数据类型体系Attribute Types下表给出 SeaTunnel 数据类型与 DynamoDB Attribute Type 的完整映射关系由 DefaultSeaTunnelRowDeserializer 中的转换逻辑实现SeaTunnel 数据类型DynamoDB Attribute 类型转换说明BOOLEANBOOLattributeValue.bool()TINYINTN数字字符串解析为 ByteSMALLINTN数字字符串解析为 ShortINTN数字字符串解析为 IntegerBIGINTN数字字符串解析为 LongFLOATN数字字符串解析为 FloatDOUBLEN数字字符串解析为 DoubleDECIMALN数字字符串解析为 BigDecimalSTRINGS字符串直接读取TIMES字符串解析为LocalTimeDATES字符串解析为LocalDateTIMESTAMPS字符串解析为LocalDateTimeBYTESB二进制数据转字节数组MAPM递归转换为MapString, ObjectARRAYL列表元素递归转换元素类型取自 schemaNULLNULL返回 null值得注意的实现细节从源码结构可见DynamoDB 数值类型N在 SDK 中一律以数字字符串形式暴露因此 Deserializer 统一通过Integer.parseInt()、BigDecimal()等方式完成解析TINYINT 对n()缺失的场景做了兼容会尝试从字符串S 类型取首个字节ARRAYL 类型会按 schema 中声明的元素类型创建同类型数组并递归转换同时兼容了 DynamoDB 的SS字符串集合、NS数字集合、BS二进制集合三种集合形态DATE / TIME / TIMESTAMP 依赖 ISO 格式字符串解析LocalDate.parse等因此写入方需保证时间字段以可解析的文本形式存储。使用注意事项读取的是快照而非变更Source 使用 Scan 请求读取当前表数据不消费 DynamoDB Streams也不会感知作业运行期间的新增/修改数据。凭据必须显式配置access_key_id与secret_access_key是必填项DynamoDB Local 场景下使用本地服务可接受的任意占位值。并行度需联动调整parallel_scan_threads控制 Scan 分段数量对大数据量表应连同env.parallelism与 sourceparallelism一起调大否则分段可能无法被充分并行消费。scan_item_limit 是分页大小它限制的是单次 Scan 请求返回的条数不是作业总行数调大它可减少请求次数但会增加单批内存占用。字段缺失返回 nullItem 中不存在的 schema 字段在转换时按 null 处理见 Deserializer 中item.get(fieldNames[i])的取值方式。完整任务示例以下示例从本地 DynamoDB Local 的source_table表读取数据并写入sink_table表演示 Source 与 Sink 的完整配置env { parallelism 2 job.mode BATCH } source { AmazonDynamoDB { url http://127.0.0.1:8000 region us-east-1 access_key_id dummy-key secret_access_key dummy-secret table source_table parallelism 2 scan_item_limit 2 parallel_scan_threads 4 schema { fields { id string c_map mapstring, smallint c_array arraytinyint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_decimal decimal(2, 1) c_bytes bytes c_date date c_timestamp timestamp } } } } sink { AmazonDynamoDB { url http://127.0.0.1:8000 region us-east-1 access_key_id dummy-key secret_access_key dummy-secret table sink_table batch_size 25 } }示例要点解读env.parallelism 2与 source 的parallelism 2保持了一致结合parallel_scan_threads 4四个逻辑分段会被轮询分配给两个并行 Reader每个 Reader 处理两个分段scan_item_limit 2表示每次 Scan 请求最多返回 2 个 ItemSDK 的分页器会自动翻页直至分段扫描完成schema 中声明的字段类型决定了 Item 的解析方式务必与 DynamoDB 表中的实际数据形态S / N / B / L / M一致Sink 侧同样连接 DynamoDBbatch_size 25用于控制写入批大小。相关资源Source 连接器源码connector-amazondynamodb连接器变更记录connector-amazondynamodb ChangelogSchema 语法Schema Feature通用选项Source Common Options连接器特性定义Connector V2 Features【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考