SeaTunnel Maxcompute Sink 连接器完全指南:认证、自动建表与 Upload/Upsert 写入策略
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
Maxcompute 是 SeaTunnel(seatunnel-connectors-v2/connector-maxcompute)中用于向阿里云 MaxCompute 写入数据的 Sink 连接器。本指南以官方文档 docs/zh/connectors/sink/Maxcompute.md 为主体,结合连接器源码,系统讲解其三种认证方式、全量配置参数、自动建表 DDL 模板、表/数据保存模式(Save Mode)以及 upload/upsert 两种写入会话的底层原理与选型建议。读完本文,你将能够独立完成从配置编写、认证选型到多表写入与 CDC 场景(更新/删除)落地的完整实战。
概述
Maxcompute Sink 连接器用于向 MaxCompute 表写入 SeaTunnel 管道中的数据。它基于阿里云官方 ODPS SDK(com.aliyun.odps)实现,通过 MaxCompute Tunnel 服务完成批量数据上传,具备以下能力:
- 支持三种认证方式:AccessKey(
accessId/accesskey)、STS Token 临时认证、阿里云默认凭据链免密认证(ECS RAM Role、环境变量等); - 支持追加写入、覆盖整表或分区,以及基于 DDL 模板的自动建表;
- 通过
insert_strategy参数在upload 会话与upsert 会话之间切换,从而支持 CDC(变更数据捕获)场景中的插入、更新(UPDATE_AFTER)与删除(DELETE)操作; - 支持多表写入(Multi-Table Sink)与多表复制数配置。
引擎支持
- SeaTunnel Zeta
- Spark
- Flink
主要特性
| 特性 | 支持情况 |
|---|---|
| 精确一次(Exactly Once) | 不支持 |
| 支持 CDC | 支持 |
| 支持多表写入 | 支持 |
| 定时刷新 | 不支持 |
特性定义详见 连接器 V2 特性说明。
认证方式:AccessKey、STS 与免密凭据链
连接器的账号构建逻辑集中在 MaxcomputeUtil.getAccount() 中,其判定优先级为:
- 配置了
sts_token:要求accessId与accesskey同时存在,否则抛出IllegalArgumentException,最终构造StsAccount(临时认证账号)。 - 同时配置了
accessId与accesskey:构造AliyunAccount(长期 AccessKey 认证)。 - 三者均未配置:构造
AklessAccount(new DefaultCredentialsProvider()),即回退到阿里云默认凭据链com.aliyun.credentials.provider.DefaultCredentialsProvider,按顺序读取环境变量、系统属性、CLI 配置文件、OIDC 以及 ECS RAM 角色等来源的凭证,实现免密认证。
免密认证(ECS RAM Role、环境变量等):只需将
accessId、accesskey和sts_token全部留空不填,连接器即自动使用阿里云默认凭据链(DefaultCredentialsProvider)读取凭证(包括环境变量、系统属性、CLI 配置文件、OIDC 以及 ECS RAM 角色)。该能力在源码MaxcomputeUtil.getAccount中以AklessAccount分支实现。
在拿到Account后,MaxcomputeUtil.getOdps() 会构造Odps客户端并依次设置endpoint、默认 Project 与当前 Schema(setCurrentSchema,来自可选的schema_name参数)。
配置参数详解
以下为 Sink 全部选项(定义见 MaxcomputeBaseOptions.java 与 MaxcomputeSinkOptions.java):
| 参数名 | 类型 | 必须 | 默认值 | 说明 |
|---|---|---|---|---|
| accessId | string | 否 | - | 访问 MaxCompute 的 AccessKey ID。 |
| accesskey | string | 否 | - | 访问 MaxCompute 的 AccessKey Secret。 |
| sts_token | string | 否 | - | MaxCompute 临时认证 STS Token;配置sts_token时accessId与accesskey必填。 |
| endpoint | string | 是 | - | MaxCompute 端点,以http开头。 |
| project | string | 是 | - | 在阿里云中创建的 MaxCompute 项目。 |
| table_name | string | 是 | - | 目标 MaxCompute 表名,例如fake。 |
| schema_name | string | 否 | - | MaxCompute Schema 名称;仅当表位于非默认 Schema 时需要设置。 |
| partition_spec | string | 否 | - | MaxCompute 分区表的规范,例如ds='20220101'。 |
| overwrite | boolean | 否 | false | 是否覆盖整张表或单个分区。 |
| schema_save_mode | enum | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 写入前如何处理目标表结构,例如RECREATE_SCHEMA或CREATE_SCHEMA_WHEN_NOT_EXIST。 |
| data_save_mode | enum | 否 | APPEND_DATA | 写入前如何处理已有数据,例如DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS。 |
| custom_sql | string | 否 | - | 当data_save_mode = CUSTOM_PROCESSING时执行的 SQL。 |
| save_mode_create_template | string | 否 | 见下文 | 在 sink 自动建表时使用的 DDL 模板。 |
| datetime_format | string | 否 | yyyy-MM-dd HH:mm:ss | 将LocalDateTime字段序列化为字符串时使用的格式。 |
| tunnel_endpoint | string | 否 | - | MaxCompute Tunnel 服务的自定义端点;未配置时根据区域自动推断。 |
| tunnel_name | string | 否 | - | Tunnel Quota 名称;需同时将endpoint与tunnel_endpoint配置为 VPC 端点。 |
| insert_strategy | string | 否 | upload | 插入会话类型:upload使用 upload 会话,upsert使用 upsert 会话并要求目标表存在主键。 |
| multi_table_sink_replica | int | 否 | 1 | 多表写入时每张表对应的 Sink Writer 副本数。 |
| common-options | - | 否 | - | Sink 插件通用参数,例如plugin_input。 |
注意:
datetime_format与multi_table_sink_replica分别来自 SeaTunnel API 的FormatOptions.DATETIME_FORMAT与SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA,在 MaxcomputeSinkFactory.java 中被注册为可用选项。
accessId [string]
您的 Maxcompute accessId,可从阿里云访问。
accesskey [string]
您的 Maxcompute accessKey,可从阿里云访问。
sts_token [string]
您的 MaxCompute STS Token,用于临时认证。注意:如果提供了sts_token,则必须同时提供accessId和accesskey。
endpoint [string]
您的 Maxcompute endpoint,以 http 开头,例如http://service.odps.aliyun.com/api。
project [string]
您在阿里云中创建的 Maxcompute 项目名。
table_name [string]
目标 Maxcompute 表名,例如fake。支持占位符用于多表写入场景(见下文multi_table_sink_replica)。
partition_spec [string]
Maxcompute 分区表的规范,例如ds='20220101'。当目标表是分区表时使用;同时它与overwrite、schema_save_mode等参数配合控制分区级别的覆盖或重建行为(见下文保存模式)。
schema_name [string]
MaxCompute Schema 名称(Project 与 Table 之间的命名空间)。仅当表位于 MaxCompute 项目的非默认 Schema时才需要设置。
- 默认值:不设置(使用项目默认 Schema)。
- 从源码看,该值会通过
Odps#setCurrentSchema写入 ODPS 客户端,并进一步被 Tunnel 的UploadSession/UpsertSession/DownloadSession继承(见 MaxcomputeUtil.java 中的buildDownloadSession与buildUploadSession)。
overwrite [boolean]
是否覆盖表或分区,默认值:false。
兼容性说明:在 MaxcomputeSink.java 中,当
overwrite = true时,连接器会打印告警日志(The configuration of 'overwrite' is deprecated, please use 'data_save_mode' instead.)并将data_save_mode强制置为DROP_DATA,即覆盖场景最终通过data_save_mode执行。新任务建议直接使用data_save_mode。
save_mode_create_template [string]
连接器使用模板来自动创建 MaxCompute 表,它会根据上游数据和 Schema 类型生成相应的建表语句。默认模板在 MaxcomputeSinkOptions.java 中定义,等价于:
CREATE TABLE IF NOT EXISTS `${table}` ( ${rowtype_fields} ) COMMENT '${comment}' ;默认模板可以根据实际情况修改。目前仅在多表模式下工作。
如果在模板中填入自定义字段,例如添加id字段:
CREATE TABLE IF NOT EXISTS `${table}` ( id, ${rowtype_fields} ) COMMENT '${comment}';连接器将自动从上游获取相应的类型来完成填充,并从rowtype_fields中删除id字段。此方法可用于自定义修改字段类型和属性(模板解析逻辑可参考 CreateTableParser.java,它以括号配平方式解析建表语句中的列定义,并跳过PRIMARY KEY等约束行)。
您可以使用以下占位符:
| 占位符 | 说明 |
|---|---|
| database | 用于获取上游模式中的数据库 |
| table_name | 用于获取上游模式中的表名 |
| rowtype_fields | 用于获取上游模式中的所有字段,将自动映射为 MaxCompute 的字段描述 |
| rowtype_primary_key | 用于获取上游模式中的主键(可能是列表) |
| rowtype_unique_key | 用于获取上游模式中的唯一键(可能是列表) |
| comment | 用于获取上游模式中的表注释 |
schema_save_mode [Enum]
在同步任务打开之前,为目标端现有的表结构选择不同的处理方案。可选值:
RECREATE_SCHEMA:表不存在时将创建;表已存在时删除并重建。如果设置了partition_spec,分区将被删除并重建。CREATE_SCHEMA_WHEN_NOT_EXIST(默认):表不存在时将创建;表已存在时跳过。如果设置了partition_spec,分区将被创建。ERROR_WHEN_SCHEMA_NOT_EXIST:表不存在时报错。IGNORE:忽略表的处理。
从源码看,该模式由 MaxComputeSaveModeHandler.java 继承 SeaTunnel API 的DefaultSaveModeHandler实现,并在createSchemaWhenNotExist与recreateSchema两个钩子中补充了分区创建逻辑:当配置了partition_spec时,调用 MaxComputeCatalog.createPartition() 创建对应分区(createPartition(partitionSpec, true),幂等创建)。
data_save_mode [Enum]
在同步任务打开之前,为目标端现有的数据选择不同的处理方案。可选值:
DROP_DATA:保留数据库结构并删除数据。APPEND_DATA(默认):保留数据库结构,保留数据。CUSTOM_PROCESSING:用户定义的处理(需配合custom_sql)。ERROR_WHEN_DATA_EXISTS:当存在数据时报错。
从源码看,对应MaxComputeCatalog中的truncateTable/清理逻辑:当存在partition_spec时执行deletePartition + createPartition(重建分区实现覆盖),否则执行odpsTable.truncate()。
custom_sql [String]
当data_save_mode选择CUSTOM_PROCESSING时,您应该填入custom_sql参数。此参数通常填入可以执行的 SQL,SQL 将在同步任务开始之前执行(由DefaultSaveModeHandler在任务打开阶段调用)。
datetime_format [String]
用户定义的格式字符串,用于将LocalDateTime字段转换为字符串。
当您想指定与DateTimeUtils.Formatter中的预定义值之一匹配的自定义日期时间格式时,请使用此选项(例如yyyy-MM-dd HH:mm:ss、yyyyMMddHHmmss等)。在 MaxcomputeOutputFormat.java 中,该选项被包装为FormatterContext,在将 SeaTunnel 行数据映射为 MaxCompute Record 时用于格式化日期时间字段。
示例值:
yyyy-MM-dd HH:mm:ssyyyy-MM-dd HH:mm:ss.SSSSSSyyyy.MM.dd HH:mm:ssyyyy/MM/dd HH:mm:ssyyyy/M/d HH:mmyyyy-M-d HH:mmyyyy/M/d HH:mm:ssyyyy-M-d HH:mm:ssyyyyMMddHHmmss
默认值:yyyy-MM-dd HH:mm:ss
tunnel_endpoint [String]
指定 MaxCompute Tunnel 服务的自定义端点 URL。
- 默认情况下,端点从配置的区域自动推断。
- 此选项允许您覆盖默认行为并使用自定义 Tunnel 端点。
- 通常您不需要设置
tunnel_endpoint,仅在自定义网络、调试或本地开发时才需要。
示例值:
https://dt.cn-hangzhou.maxcompute.aliyun.comhttps://dt.ap-southeast-1.maxcompute.aliyun.comhttp://maxcompute:8080
默认值:未设置(从区域自动推断)。在 MaxcomputeUtil.getTableTunnel() 中,配置了该值时会对TableTunnel执行setEndpoint。
tunnel_name [String]
tunnel_name指定 Tunnel Quota 名称,用于独占资源组。
Tunnel Quota 允许您使用专用的计算资源进行 MaxCompute Tunnel 数据传输,从而提供更好的性能和资源隔离。
重要提示:Tunnel Quota 仅在VPC(虚拟私有云)端点下生效,暂不支持公共网络访问。使用tunnel_name时,必须同时将endpoint和tunnel_endpoint配置为 VPC 端点。
如果未指定,将使用默认的 Tunnel quota。源码中对应 MaxcomputeUtil.getTableTunnel() 的tableTunnel.getConfig().setQuotaName(...)调用。
示例值:your_tunnel_quota_name
默认值:未设置(使用默认 quota)
insert_strategy [string]
插入会话类型,默认upload。写入会话的创建与数据分发集中在 MaxcomputeOutputFormat.java:
- 设置为
upload:使用upload 会话(TableTunnel.UploadSession+openBufferedWriter缓冲写入,close()时commit)。 - 设置为
upsert:使用upsert 会话(TableTunnel.UpsertSession+buildUpsertStream,close()时commit(true)),要求目标表存在主键。
注意:
- 在同时存在更新或删除操作的情况下,使用 upload 会话进行插入操作,可能会导致插入的记录比预期更晚出现在表中。
- 当表中存在主键时,建议将
insert_strategy设置为upsert,以确保一致的 upsert 行为。 UPDATE_AFTER和DELETE数据都会通过 MaxCompute upsert 会话写入,所以任务包含更新或删除数据时,目标表必须有主键。当前 Sink 不支持UPDATE_BEFORE数据(MaxcomputeOutputFormat.write() 中仅处理INSERT、UPDATE_AFTER、DELETE三种 RowKind,其余类型抛出unsupportedDataType)。
对应的 RowKind 处理逻辑如下:
| SeaTunnel RowKind | upload 会话 | upsert 会话 |
|---|---|---|
| INSERT | recordWriter.write(缓冲写入) | upsertStream.upsert |
| UPDATE_AFTER | 不支持(抛出异常) | upsertStream.upsert |
| DELETE | 不支持(抛出异常) | upsertStream.delete |
| 其他(含 UPDATE_BEFORE) | 不支持(抛出异常) | 不支持(抛出异常) |
写入结束后,MaxcomputeWriter.close() 会统一关闭会话:upload 会话执行uploadSession.commit(),upsert 会话执行upsertSession.commit(true)后关闭,保证数据提交。
multi_table_sink_replica [int]
多表写入模式下的 writer 副本数,默认值为1。
当上游数据包含多张表,并且table_name使用${table_name}这类占位符时可以配置该参数。例如table_name = "${table_name}_sink"会把上游表test_table写入目标表test_table_sink。
通用选项
Sink 插件通用参数,例如plugin_input(指定当前插件处理的数据集,适用于多 source/transform/sink 场景),请参考 Sink 通用选项 详见。
配置示例
追加写入
最简单的追加写入场景,使用 AccessKey 认证:
sink { Maxcompute { accessId="<your access id>" accesskey="<your access Key>" endpoint="<http://service.odps.aliyun.com/api>" project="<your project>" table_name="<your table name>" #partition_spec="<your partition spec>" #overwrite = false } }多表写入
上游使用FakeSource的tables_configs模拟多张表,Sink 端通过table_name = "${table_name}_sink"将上游表名映射为目标表名,并配置insert_strategy = "upsert":
source { FakeSource { tables_configs = [ { schema = { table = "test_table" fields { ID = int NAME = string AGE = int } primaryKey { name = "ID" columnNames = [ID] } } rows = [ { kind = INSERT, fields = [1, "INSERT_TEST1", 20] } { kind = INSERT, fields = [2, "INSERT_TEST2", 30] } ] }, { schema = { table = "test_table_2" fields { ID = int NAME = string AGE = int } primaryKey { name = "ID" columnNames = [ID] } } rows = [ { kind = INSERT, fields = [1, "INSERT_TEST1", 20] } ] } ] } } sink { Maxcompute { accessId = "ak" accesskey = "sk" endpoint = "http://maxcompute:8080" tunnel_endpoint = "http://maxcompute:8080" project = "mocked_mc" table_name = "${table_name}_sink" insert_strategy = "upsert" multi_table_sink_replica = 1 } }该示例中两张上游表分别写入test_table_sink与test_table_2_sink。由于配置了primaryKey且insert_strategy = "upsert",写入走 upsert 会话。
更新插入或删除数据
当上游表结构有主键,并且任务里包含更新或删除数据时,建议配置insert_strategy = "upsert":
source { FakeSource { tables_configs = [ { schema = { table = "test_table_sink" fields { ID = int NAME = string AGE = int } primaryKey { name = "ID" columnNames = [ID] } } rows = [ { kind = UPDATE_AFTER fields = [1, "UPSERT_TEST", 100] } ] } ] } } sink { Maxcompute { accessId = "ak" accesskey = "sk" endpoint = "http://maxcompute:8080" tunnel_endpoint = "http://maxcompute:8080" project = "mocked_mc" table_name = "test_table_sink" insert_strategy = "upsert" } }源码级写入链路
综合以上配置,一次 Maxcompute 写入的完整调用链如下(可在仓库对应文件中逐个核对):
- Sink 构建:MaxcomputeSinkFactory 注册全部选项并构建 MaxcomputeSink。
- 保存模式处理:
MaxcomputeSink.getSaveModeHandler()通过 SPI 发现MaxComputeCatalog,组装MaxComputeSaveModeHandler,在任务启动前按schema_save_mode/data_save_mode完成建表、重建分区、清空数据或执行custom_sql(涉及 MaxComputeSaveModeHandler.java 与 MaxComputeCatalog.java)。overwrite = true时在此处被转换为DROP_DATA。 - 认证与客户端:MaxcomputeUtil 按
sts_token → accessId/accesskey → DefaultCredentialsProvider的优先级构造Account,进而创建Odps客户端与TableTunnel(可选设置 Tunnel 端点与 Quota)。 - 写入:MaxcomputeWriter 将 SeaTunnelRow 交给 MaxcomputeOutputFormat,按 RowKind 分发到 upload 会话的
RecordWriter或 upsert 会话的UpsertStream;关闭时分别commit/commit(true)。 - 类型映射:行数据经 MaxcomputeTypeMapper 与
FormatterContext(datetime_format)转换为 MaxComputeRecord后写入。
连接器同时提供了配套的单元测试用于印证配置解析与类型转换行为,例如 MaxcomputeSourceFactoryTest.java 与 MaxcomputeUtilTest.java,可作为理解选项与底层行为的参考。
变更日志
连接器的历史变更记录请见 Maxcompute 连接器变更日志。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考