news 2026/9/17 11:34:44

SeaTunnel AmazonDynamoDB Source Connector 使用指南:基于 Scan 的批量数据读取与并行分段实现解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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/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 同时实现了SupportParallelismSupportColumnProjection两个接口,其中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 请求返回的最大条数)。随后通过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);

这里segmenttotalSegments正是 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 中,汇总如下:

名称类型必填默认值说明
urlstring-DynamoDB 服务端点 URL
regionstring-DynamoDB 服务所在 AWS 区域
access_key_idstring-AWS 访问密钥 ID
secret_access_keystring-AWS 访问密钥 Secret
tablestring-要扫描的 DynamoDB 表名
schemaconfig-从 DynamoDB Item 中读取的 SeaTunnel 字段定义
scan_item_limitint1每次 Scan 请求返回的最大 Item 数
parallel_scan_threadsint2并行扫描的逻辑分段数量
common-optionsobject-Source 插件通用参数

从 AmazonDynamoDBSourceFactory 的OptionRule可以看到,urlregionaccess_key_idsecret_access_keytableschema六项为必填,scan_item_limitparallel_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 中通过StaticCredentialsProviderAwsBasicCredentials将二者显式注入客户端。该连接器必须显式提供这两个参数;使用 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、sourceparallelismparallel_scan_threads,让多个 Reader 各自扫描不同分段,充分发挥 DynamoDB Parallel Scan 的吞吐能力。

common options

Source 插件通用参数(如parallelismresult_table_name等),详见 Source Common Options。

数据类型映射

DynamoDB 使用自身的数据类型体系(Attribute Types),下表给出 SeaTunnel 数据类型与 DynamoDB Attribute Type 的完整映射关系(由 DefaultSeaTunnelRowDeserializer 中的转换逻辑实现):

SeaTunnel 数据类型DynamoDB Attribute 类型转换说明
BOOLEANBOOLattributeValue.bool()
TINYINTN数字字符串解析为 Byte
SMALLINTN数字字符串解析为 Short
INTN数字字符串解析为 Integer
BIGINTN数字字符串解析为 Long
FLOATN数字字符串解析为 Float
DOUBLEN数字字符串解析为 Double
DECIMALN数字字符串解析为 BigDecimal
STRINGS字符串直接读取
TIMES字符串解析为LocalTime
DATES字符串解析为LocalDate
TIMESTAMPS字符串解析为LocalDateTime
BYTESB二进制数据转字节数组
MAPM递归转换为Map<String, Object>
ARRAYL列表元素递归转换,元素类型取自 schema
NULLNULL返回 null

值得注意的实现细节(从源码结构可见):

  • DynamoDB 数值类型(N)在 SDK 中一律以数字字符串形式暴露,因此 Deserializer 统一通过Integer.parseInt()BigDecimal()等方式完成解析;
  • TINYINT 对n()缺失的场景做了兼容,会尝试从字符串(S 类型)取首个字节;
  • ARRAY(L 类型)会按 schema 中声明的元素类型创建同类型数组并递归转换,同时兼容了 DynamoDB 的SS(字符串集合)、NS(数字集合)、BS(二进制集合)三种集合形态;
  • DATE / TIME / TIMESTAMP 依赖 ISO 格式字符串解析(LocalDate.parse等),因此写入方需保证时间字段以可解析的文本形式存储。

使用注意事项

  1. 读取的是快照而非变更:Source 使用 Scan 请求读取当前表数据,不消费 DynamoDB Streams,也不会感知作业运行期间的新增/修改数据。
  2. 凭据必须显式配置access_key_idsecret_access_key是必填项;DynamoDB Local 场景下使用本地服务可接受的任意占位值。
  3. 并行度需联动调整parallel_scan_threads控制 Scan 分段数量,对大数据量表应连同env.parallelism与 sourceparallelism一起调大,否则分段可能无法被充分并行消费。
  4. scan_item_limit 是分页大小:它限制的是单次 Scan 请求返回的条数,不是作业总行数;调大它可减少请求次数,但会增加单批内存占用。
  5. 字段缺失返回 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 11:31:12

gh-aw Playwright与Web搜索能力:让Agent会查资料会点网页

gh-aw Playwright与Web搜索能力&#xff1a;让Agent会查资料会点网页 【免费下载链接】gh-aw GitHub Agentic Workflows 项目地址: https://gitcode.com/GitHub_Trending/gha/gh-aw gh-aw 是一款把 AI Agent 跑在 GitHub Actions 上的开源工具&#xff08;GitHub Agenti…

作者头像 李华
网站建设 2026/9/17 11:28:18

OCA认证模拟题15解析:SQL查询过滤与连接考点复盘

简介&#xff1a;OCA认证分类模拟题15是一份面向Oracle认证助理&#xff08;OCA&#xff09;备考者的配套练习文档&#xff0c;聚焦数据库升级、数据泵&#xff08;Data Pump&#xff09;迁移、空间管理、组件兼容性等核心模块&#xff0c;适合正在系统准备OCA考试&#xff0c;…

作者头像 李华
网站建设 2026/9/17 11:22:06

进口编码器停产后怎么替代?机械电气协议参数四层匹配与三条路径

周三凌晨两点&#xff0c;产线维护群里弹出一条消息&#xff1a;三号机送料轴报编码器故障&#xff0c;换上去的备用件开始跳数。翻出原件型号去官网一查&#xff0c;停产已经五年&#xff0c;最后一批库存也被同行扫走了。这种场面在设备圈太常见了——编码器这种看着不起眼的…

作者头像 李华
网站建设 2026/9/17 11:20:52

ZeroClaw执行机制解析:WASM沙箱、DynamicExec与lced协同原理

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/17 11:20:50

DMA完成通知机制:从MSI-X、完成队列到高性能轮询

1. "我干完了"这句话&#xff0c;硬件到底是用什么方式说出来的写驱动或者调 AI 推理服务的时候&#xff0c;最容易被忽略的一环&#xff0c;恰恰是数据搬完之后那句"我干完了"。DMA 把一块 buffer 从网卡搬到内存、从 SSD 搬到主机、从主机搬到 GPU 显存&…

作者头像 李华