十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel Maxcompute Sink 连接器完全指南:认证、自动建表与 Upload/Upsert 写入策略

SeaTunnel Maxcompute Sink 连接器完全指南:认证、自动建表与 Upload/Upsert 写入策略 SeaTunnel Maxcompute Sink 连接器完全指南认证、自动建表与 Upload/Upsert 写入策略【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelMaxcompute 是 SeaTunnelseatunnel-connectors-v2/connector-maxcompute中用于向阿里云 MaxCompute 写入数据的 Sink 连接器。本指南以官方文档 docs/zh/connectors/sink/Maxcompute.md 为主体结合连接器源码系统讲解其三种认证方式、全量配置参数、自动建表 DDL 模板、表/数据保存模式Save Mode以及 upload/upsert 两种写入会话的底层原理与选型建议。读完本文你将能够独立完成从配置编写、认证选型到多表写入与 CDC 场景更新/删除落地的完整实战。概述Maxcompute Sink 连接器用于向 MaxCompute 表写入 SeaTunnel 管道中的数据。它基于阿里云官方 ODPS SDKcom.aliyun.odps实现通过 MaxCompute Tunnel 服务完成批量数据上传具备以下能力支持三种认证方式AccessKeyaccessId/accesskey、STS Token 临时认证、阿里云默认凭据链免密认证ECS RAM Role、环境变量等支持追加写入、覆盖整表或分区以及基于 DDL 模板的自动建表通过insert_strategy参数在upload 会话与upsert 会话之间切换从而支持 CDC变更数据捕获场景中的插入、更新UPDATE_AFTER与删除DELETE操作支持多表写入Multi-Table Sink与多表复制数配置。引擎支持SeaTunnel ZetaSparkFlink主要特性特性支持情况精确一次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 与当前 SchemasetCurrentSchema来自可选的schema_name参数。配置参数详解以下为 Sink 全部选项定义见 MaxcomputeBaseOptions.java 与 MaxcomputeSinkOptions.java参数名类型必须默认值说明accessIdstring否-访问 MaxCompute 的 AccessKey ID。accesskeystring否-访问 MaxCompute 的 AccessKey Secret。sts_tokenstring否-MaxCompute 临时认证 STS Token配置sts_token时accessId与accesskey必填。endpointstring是-MaxCompute 端点以http开头。projectstring是-在阿里云中创建的 MaxCompute 项目。table_namestring是-目标 MaxCompute 表名例如fake。schema_namestring否-MaxCompute Schema 名称仅当表位于非默认 Schema 时需要设置。partition_specstring否-MaxCompute 分区表的规范例如ds20220101。overwriteboolean否false是否覆盖整张表或单个分区。schema_save_modeenum否CREATE_SCHEMA_WHEN_NOT_EXIST写入前如何处理目标表结构例如RECREATE_SCHEMA或CREATE_SCHEMA_WHEN_NOT_EXIST。data_save_modeenum否APPEND_DATA写入前如何处理已有数据例如DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS。custom_sqlstring否-当data_save_mode CUSTOM_PROCESSING时执行的 SQL。save_mode_create_templatestring否见下文在 sink 自动建表时使用的 DDL 模板。datetime_formatstring否yyyy-MM-dd HH:mm:ss将LocalDateTime字段序列化为字符串时使用的格式。tunnel_endpointstring否-MaxCompute Tunnel 服务的自定义端点未配置时根据区域自动推断。tunnel_namestring否-Tunnel Quota 名称需同时将endpoint与tunnel_endpoint配置为 VPC 端点。insert_strategystring否upload插入会话类型upload使用 upload 会话upsert使用 upsert 会话并要求目标表存在主键。multi_table_sink_replicaint否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 分区表的规范例如ds20220101。当目标表是分区表时使用同时它与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参数。此参数通常填入可以执行的 SQLSQL 将在同步任务开始之前执行由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:sstunnel_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默认值未设置使用默认 quotainsert_strategy [string]插入会话类型默认upload。写入会话的创建与数据分发集中在 MaxcomputeOutputFormat.java设置为upload使用upload 会话TableTunnel.UploadSessionopenBufferedWriter缓冲写入close()时commit。设置为upsert使用upsert 会话TableTunnel.UpsertSessionbuildUpsertStreamclose()时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 RowKindupload 会话upsert 会话INSERTrecordWriter.write缓冲写入upsertStream.upsertUPDATE_AFTER不支持抛出异常upsertStream.upsertDELETE不支持抛出异常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 { accessIdyour access id accesskeyyour access Key endpointhttp://service.odps.aliyun.com/api projectyour project table_nameyour table name #partition_specyour partition spec #overwrite false } }多表写入上游使用FakeSource的tables_configs模拟多张表Sink 端通过table_name ${table_name}_sink将上游表名映射为目标表名并配置insert_strategy upsertsource { 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 upsertsource { 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 与FormatterContextdatetime_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),仅供参考
返回列表