SeaTunnel AmazonDynamoDB Source Connector 使用指南:基于 Scan 的批量数据读取与并行分段实现解析
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
Amazon 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,也不消费 CDC(Change 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(逻辑分段数量)为整张表创建对应数量的AmazonDynamoDBSourceSplit:
int 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 请求返回的最大条数)。随后通过getSplitOwner(assignCount % 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 Item(Map<String, AttributeValue>)通过 DefaultSeaTunnelRowDeserializer 转换为SeaTunnelRow。转换完全依据用户在schema.fields中声明的 SeaTunnel 类型逐字段进行,若字段在 Item 中缺失则对应值为null。
连接器选项详解
Source 插件的全部选项定义在 AmazonDynamoDBSourceOptions 与 AmazonDynamoDBBaseOptions 中,汇总如下:
| 名称 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| url | string | 是 | - | DynamoDB 服务端点 URL |
| region | string | 是 | - | DynamoDB 服务所在 AWS 区域 |
| access_key_id | string | 是 | - | AWS 访问密钥 ID |
| secret_access_key | string | 是 | - | AWS 访问密钥 Secret |
| table | string | 是 | - | 要扫描的 DynamoDB 表名 |
| schema | config | 是 | - | 从 DynamoDB Item 中读取的 SeaTunnel 字段定义 |
| scan_item_limit | int | 否 | 1 | 每次 Scan 请求返回的最大 Item 数 |
| parallel_scan_threads | int | 否 | 2 | 并行扫描的逻辑分段数量 |
| common-options | object | 否 | - | 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:8000"region [string]
DynamoDB 服务所在 AWS 区域,例如us-east-1。在 AmazonDynamoDBSourceReader#open() 中可以看到,region 会通过Region.of()传入DynamoDbClient构建器——即便连接 DynamoDB Local,region 也是客户端构建校验所必需的(对本地服务而言无实际意义,但不可省略)。
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 = "map<string, smallint>" c_array = "array<tinyint>" 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 options
Source 插件通用参数(如parallelism、result_table_name等),详见 Source Common Options。
数据类型映射
DynamoDB 使用自身的数据类型体系(Attribute Types),下表给出 SeaTunnel 数据类型与 DynamoDB Attribute Type 的完整映射关系(由 DefaultSeaTunnelRowDeserializer 中的转换逻辑实现):
| SeaTunnel 数据类型 | DynamoDB Attribute 类型 | 转换说明 |
|---|---|---|
| BOOLEAN | BOOL | attributeValue.bool() |
| TINYINT | N | 数字字符串解析为 Byte |
| SMALLINT | N | 数字字符串解析为 Short |
| INT | N | 数字字符串解析为 Integer |
| BIGINT | N | 数字字符串解析为 Long |
| FLOAT | N | 数字字符串解析为 Float |
| DOUBLE | N | 数字字符串解析为 Double |
| DECIMAL | N | 数字字符串解析为 BigDecimal |
| STRING | S | 字符串直接读取 |
| TIME | S | 字符串解析为LocalTime |
| DATE | S | 字符串解析为LocalDate |
| TIMESTAMP | S | 字符串解析为LocalDateTime |
| BYTES | B | 二进制数据转字节数组 |
| MAP | M | 递归转换为Map<String, Object> |
| ARRAY | L | 列表元素递归转换,元素类型取自 schema |
| NULL | NULL | 返回 null |
值得注意的实现细节(从源码结构可见):
- DynamoDB 数值类型(N)在 SDK 中一律以数字字符串形式暴露,因此 Deserializer 统一通过
Integer.parseInt()、BigDecimal()等方式完成解析; - TINYINT 对
n()缺失的场景做了兼容,会尝试从字符串(S 类型)取首个字节; - ARRAY(L 类型)会按 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 请求返回的条数,不是作业总行数;调大它可减少请求次数,但会增加单批内存占用。
- 字段缺失返回 null:Item 中不存在的 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 = "map<string, smallint>" c_array = "array<tinyint>" 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 个 Item,SDK 的分页器会自动翻页直至分段扫描完成;- schema 中声明的字段类型决定了 Item 的解析方式,务必与 DynamoDB 表中的实际数据形态(S / N / B / L / M)一致;
- Sink 侧同样连接 DynamoDB,
batch_size = 25用于控制写入批大小。
相关资源
- Source 连接器源码:connector-amazondynamodb
- 连接器变更记录:connector-amazondynamodb Changelog
- Schema 语法: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),仅供参考