SeaTunnel S3Redshift Sink 连接器实战:基于 S3 + Redshift COPY 的海量数据导入方案
2026/9/18 16:10:48 网站建设 项目流程

SeaTunnel S3Redshift Sink 连接器实战:基于 S3 + Redshift COPY 的海量数据导入方案

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

S3Redshift 是 SeaTunnel 提供的一个「S3 + Redshift 双阶段」Sink 连接器:先把数据写入 AWS S3 上的临时/目标文件,再利用 Redshift 的COPY命令把 S3 文件批量装载进 Redshift 表,借助两阶段提交(2PC)实现精确一次(Exactly-Once)语义。读完本文,你将掌握 S3Redshift 的完整参数体系、execute_sql${path}占位符的底层替换机制、五种文件格式(text/csv/parquet/orc/json)的配置差异,以及三份可直接运行的 HOCON 作业示例。

连接器概述与设计思路

S3Redshift 的作用是将数据写入 S3,然后使用 Redshift 的COPY命令将数据从 S3 导入 Redshift(见 S3-Redshift.md)。

这一设计充分利用了两个系统的优势:

  • S3 充当中间缓冲层:SeaTunnel 以文件形式把数据批量落盘到 S3,避免逐条写入 Redshift 带来的性能开销;
  • Redshift COPY 负责高速装载COPY是 Redshift 官方推荐的批量导入方式,按列并行加载,性能远高于逐条INSERT

从源码结构看,S3Redshift 是基于 S3File 文件 Sink 实现的:S3RedshiftSink直接继承BaseFileSink(S3RedshiftSink.java),文件写入逻辑完全复用 S3File 的能力,因此所有 S3File 的配置项(如file_format_typepartition_bysink_columns等)都可直接使用;在此之上叠加了 Redshift 的 JDBC 连接与 COPY SQL 执行能力。

为了支持更多文件类型,S3Redshift 使用 HDFS 协议对 S3 进行内部访问,因此该连接器需要一些 Hadoop 依赖,且只支持 Hadoop 版本2.6.5+。对应的依赖在 pom.xml 中体现为connector-file-base-hadoopconnector-file-s3以及 Redshift JDBC 驱动com.amazon.redshift:redshift-jdbc42:2.1.0.30

主要特性

  • 精确一次(Exactly-Once):默认使用 2PC commit 来确保精确一次。文件先写入临时目录,提交阶段再 rename 到目标路径并触发 COPY。
  • 文件格式类型:
    • text
    • csv
    • parquet
    • orc
    • json
  • 定时刷新:暂不支持。

精确一次的实现原理

S3Redshift 的精确一次依赖 SeaTunnel 的文件 Sink 两阶段提交框架。关键实现在 S3RedshiftSinkAggregatedCommitter.java,其commit流程为:

  1. 遍历每个事务的文件映射(transactionMap),先把临时文件renameFile到目标路径;
  2. 通过convertSql(mvFileEntry.getValue())execute_sql中的${path}占位符替换为实际文件路径(见convertSql实现StringUtils.replace(executeSql, "${path}", path));
  3. 调用 RedshiftJdbcClient 的execute(sql)执行这条 COPY 语句,将 S3 文件导入 Redshift;
  4. 文件导入成功后删除已导入的文件,并清理事务目录。

abort阶段则直接删除事务目录,保证失败时不会在目标路径留下半成品数据。整个链路把「S3 文件可见」与「Redshift 数据可见」绑定在同一个提交事务内,从而实现不丢失、不重复。

参数详解

S3Redshift 的参数由两部分组成:Redshift 专属参数(定义于 S3RedshiftSinkOptions.java)和从 S3File 继承的文件参数。完整参数表如下:

名称类型是否必填默认值描述
jdbc_urlstring-连接 Redshift 数据库的 JDBC URL,例如jdbc:redshift://your-cluster.region.redshift.amazonaws.com:5439/your_database
jdbc_userstring-连接 Redshift 数据库的用户名。
jdbc_passwordstring-连接 Redshift 数据库的密码。
execute_sqlstring-数据写入 S3 之后要执行的 SQL,通常是一条 RedshiftCOPY命令(必须包含${path}占位符)。
pathstring-bucket 下的目标目录路径,连接器会通过${path}占位符把实际写入路径追加到execute_sql中。
bucketstring-S3 文件系统的 bucket 地址,例如s3a://seatunnel-test。使用 Hadoop 读写时建议使用s3a协议。
access_keystring-S3 文件系统的 access key。如果未配置,需要正确配置 Hadoop 凭据链。
access_secretstring-S3 文件系统的 access secret。如果未配置,需要正确配置 Hadoop 凭据链。
hadoop_s3_propertiesmap-额外的 Hadoop S3A / Hadoop-AWS 选项,可设置fs.s3a.aws.credentials.provider等。
file_name_expressionstring"${transactionId}"path下追加的文件名表达式,可使用${now}${uuid}注入时间或 UUID。is_enable_transaction = true时自动在文件名前添加${transactionId}_
file_format_typestring"text"写入 S3 的文件格式,支持textcsvparquetorcjson。最终文件名带相应后缀(text的后缀为txt)。
filename_time_formatstring"yyyy.MM.dd"解析file_name_expression${now}的时间格式。
field_delimiterstring'\001'textcsv文件的列分隔符。
row_delimiterstring"\n"textcsv文件的行分隔符。
partition_byarray-按指定的上游字段对数据进行分区,分区目录由partition_dir_expression推导。
partition_dir_expressionstring"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/"根据partition_by字段生成分区目录的表达式。
is_partition_field_write_in_filebooleanfalsetrue时,分区字段及其值写入数据文件。Hive 风格数据文件请设为false
sink_columnsarray为空时所有字段都是 sink 列需要写入文件的列,字段顺序决定文件实际写入顺序。
is_enable_transactionbooleantruetrue时保证数据写入目标目录不丢失、不重复。目前只支持true
batch_sizeint1000000单个文件的最大行数。在 SeaTunnel Zeta 引擎中,每文件行数由batch_sizecheckpoint.interval共同决定。
common-options--Sink 插件通用参数,见 Sink 通用选项。

参数校验规则(来自源码)

从 S3RedshiftSinkFactory.java 的optionRule()可以看到连接器启动时的强校验逻辑:

  • 必填参数bucketjdbc_urljdbc_userjdbc_passwordexecute_sqlpathFILE_PATH)、fs.s3a.aws.credentials.providerS3A_AWS_CREDENTIALS_PROVIDER_CLASS);其中四个 Redshift 参数还带有Conditions.notBlank约束,空白字符串无法通过校验;
  • 条件参数:当fs.s3a.aws.credentials.providerSimpleAWSCredentialsProvider时,必须提供access_keysecret_key;当file_format_typetext时必须提供field_delimiterrow_delimiter;为csv时必须提供row_delimiter

上述校验规则由测试 S3RedshiftSinkFactoryTest.java 覆盖验证:testBlankRedshiftOptionsRejected断言四个 Redshift 参数为空或纯空白时会抛出OptionValidationException

Redshift 专属参数

jdbc_url

连接到 Redshift 数据库的 JDBC URL。连接器内部由 RedshiftJdbcClient.java 加载驱动com.amazon.redshift.jdbc42.Driver并调用DriverManager.getConnection(url, user, password)建立连接,因此该 URL 必须是符合 Redshift JDBC 驱动规范的完整连接串。

jdbc_user / jdbc_password

连接 Redshift 数据库的用户名与密码,用于建立 JDBC 连接。由于该客户端是进程内单例(RedshiftJdbcClient.getInstance,见 RedshiftJdbcClient.java),所有事务共享同一条数据库连接。

execute_sql

数据写入 S3 后要执行的 SQL,通常是一条 RedshiftCOPY命令。示例:

COPY target_table FROM 's3://yourbucket${path}' IAM_ROLE 'arn:XXX' REGION 'your region' format as json 'auto';
  • target_table是 Redshift 中的表名;
  • ${path}是写入 S3 的文件的路径。请务必确认您的 SQL 包含此变量,且无需手动替换——连接器会在执行 SQL 时自动将其替换为真实文件路径(StringUtils.replace(executeSql, "${path}", path));
  • IAM_ROLE是有权访问 S3 的角色,请确认该角色拥有对 S3 的访问权限;
  • format是写入 S3 的文件的格式,请确认此格式与您在配置中设置的file_format_type一致。

关于 RedshiftCOPY的更多语法细节可参考官方 COPY 命令文档。

文件与路径相关参数(继承自 S3File)

path [string]

目标目录路径,必填项。连接器会把最终写入的文件路径通过${path}占位符动态拼接到execute_sql中,因此path决定了 RedshiftCOPY实际读取的 S3 前缀。

bucket [string]

S3 文件系统的 bucket 地址,例如:s3n://seatunnel-test;如果使用s3a协议,则此参数应为s3a://seatunnel-test。由于本连接器基于 Hadoop 协议访问 S3,建议统一使用s3a协议前缀。

access_key / access_secret [string]

S3 文件系统的 access key 与 access secret。如果未设置此参数,请确认凭据提供程序链可以正确进行身份验证。示例配置中常见写法为access_key+secret_key,二者在 Hadoop S3A 层面对应fs.s3a.access.key/fs.s3a.secret.key

hadoop_s3_properties [map]

如需添加额外的 Hadoop S3A / Hadoop-AWS 选项,可在此处添加,例如自定义凭据提供程序:

hadoop_s3_properties { "fs.s3a.aws.credentials.provider" = "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider" }

其他常用键还包括fs.s3a.buffer.dirfs.s3a.fast.upload.bufferfs.s3a.session.tokenfs.s3a.assumed.role.arn等(参考 S3File 文档中的相关说明)。

file_name_expression [string]

描述将在path中创建的文件名表达式。可以在表达式中加入变量${now}${uuid},例如test_${uuid}_${now}${now}表示当前时间,其格式由filename_time_format定义。

注意:如果is_enable_transactiontrue,连接器会自动在文件名开头添加${transactionId}_

file_format_type [string]

支持的写入文件类型:textcsvparquetorcjson

注意,最终文件名会以file_format_type对应的后缀结尾,其中text文件的后缀为txt。请确保execute_sql中 COPY 命令的format与此参数保持一致(如 parquet 对应format as PARQUET、orc 对应format as ORC)。

filename_time_format [string]

file_name_expression中的格式为xxxx-${now}时,用filename_time_format指定${now}的时间格式,默认值为yyyy.MM.dd。常用时间格式符号:

符号说明
y
M
d
H小时 (0-23)
m分钟
s

详细的时间格式语法遵循 JavaSimpleDateFormat规范。

field_delimiter / row_delimiter [string]

  • field_delimiter:数据行中列之间的分隔符,仅textcsv文件格式需要,默认'\001'
  • row_delimiter:文件中行之间的分隔符,仅textcsv文件格式需要,默认"\n"

partition_by [array] / partition_dir_expression [string]

基于选定字段对数据进行分区。如果指定了partition_by,连接器会根据分区信息生成相应的分区目录,并将最终文件放置在分区目录中。

默认的partition_dir_expression${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/,其中k0是第一个分区字段名,v0是第一个分区字段的值。

is_partition_field_write_in_file [boolean]

如果为true,分区字段及其值将写入数据文件;例如想写出 Hive 风格的数据文件,此值应设为false

sink_columns [array]

哪些列需要写入文件,默认值为从 Transform 或 Source 获取的所有列。字段的顺序决定了文件实际写入的顺序。

is_enable_transaction [boolean]

如果为true,连接器将确保数据在写入目标目录时不会丢失或重复。请注意,为true时会自动在文件名开头添加${transactionId}_。目前只支持true

batch_size [int]

文件中的最大行数。对于 SeaTunnel Engine,文件中的行数由batch_sizecheckpoint.interval共同决定:如果checkpoint.interval的值足够大,sink writer 会持续向文件中写入行,直到文件中的行数超过batch_size;如果checkpoint.interval较小,sink writer 会在新的 checkpoint 触发时创建一个新文件。

common options

Sink 插件通用参数,详见 Sink Common Options。

完整配置示例

示例一:text 文件格式

S3Redshift { jdbc_url = "jdbc:redshift://xxx.amazonaws.com.cn:5439/xxx" jdbc_user = "xxx" jdbc_password = "xxxx" execute_sql="COPY table_name FROM 's3://test${path}' IAM_ROLE 'arn:aws-cn:iam::xxx' REGION 'cn-north-1' removequotes emptyasnull blanksasnull maxerror 100 delimiter '|' ;" access_key = "xxxxxxxxxxxxxxxxx" secret_key = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" tmp_path = "/tmp/seatunnel" path="/seatunnel/text" row_delimiter="\n" partition_dir_expression="${k0}=${v0}" is_partition_field_write_in_file=true file_name_expression="${transactionId}_${now}" file_format_type = "text" filename_time_format="yyyy.MM.dd" is_enable_transaction=true hadoop_s3_properties { "fs.s3a.aws.credentials.provider" = "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider" } }

示例二:parquet 文件格式

S3Redshift { jdbc_url = "jdbc:redshift://xxx.amazonaws.com.cn:5439/xxx" jdbc_user = "xxx" jdbc_password = "xxxx" execute_sql="COPY table_name FROM 's3://test${path}' IAM_ROLE 'arn:aws-cn:iam::xxx' REGION 'cn-north-1' format as PARQUET;" access_key = "xxxxxxxxxxxxxxxxx" secret_key = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" tmp_path = "/tmp/seatunnel" path="/seatunnel/parquet" row_delimiter="\n" partition_dir_expression="${k0}=${v0}" is_partition_field_write_in_file=true file_name_expression="${transactionId}_${now}" file_format_type = "parquet" filename_time_format="yyyy.MM.dd" is_enable_transaction=true hadoop_s3_properties { "fs.s3a.aws.credentials.provider" = "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider" } }

示例三:orc 文件格式

S3Redshift { jdbc_url = "jdbc:redshift://xxx.amazonaws.com.cn:5439/xxx" jdbc_user = "xxx" jdbc_password = "xxxx" execute_sql="COPY table_name FROM 's3://test${path}' IAM_ROLE 'arn:aws-cn:iam::xxx' REGION 'cn-north-1' format as ORC;" access_key = "xxxxxxxxxxxxxxxxx" secret_key = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" tmp_path = "/tmp/seatunnel" path="/seatunnel/orc" row_delimiter="\n" partition_dir_expression="${k0}=${v0}" is_partition_field_write_in_file=true file_name_expression="${transactionId}_${now}" file_format_type = "orc" filename_time_format="yyyy.MM.dd" is_enable_transaction=true hadoop_s3_properties { "fs.s3a.aws.credentials.provider" = "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider" } }

示例配置要点提示

  • tmp_path指定了临时目录(示例为/tmp/seatunnel),文件会先写入该临时路径,提交时再mvpath目标目录,这是 2PC 提交的一部分;
  • 三个示例均使用SimpleAWSCredentialsProvider+access_key/secret_key的静态凭据方式,与 S3File 文档中推荐的凭据提供程序链保持兼容;
  • 分区表达式partition_dir_expression="${k0}=${v0}"结合is_partition_field_write_in_file=true,会把分区字段写入数据文件内,方便按分区管理 S3 对象并缩小 COPY 的扫描范围。

变更日志

S3Redshift 连接器的变更记录可参见 connector-s3-redshift 变更日志。

使用前提与注意事项

  1. Hadoop 依赖:连接器通过 HDFS 协议访问 S3,需要 Hadoop 依赖,仅支持 Hadoop2.6.5+
  2. 凭据配置access_key/access_secret未配置时,必须确保 Hadoop 凭据提供程序链能够正确完成 S3 身份验证(例如使用InstanceProfileCredentialsProvider等);
  3. COPY 权限execute_sql中使用的 IAM_ROLE 必须同时具备读取 S3 中${path}对应对象的权限与写入 Redshift 目标表的权限;
  4. 格式一致性execute_sql中 COPY 命令声明的format必须与file_format_type保持一致,否则 COPY 会因解析失败而报错;
  5. 事务开关is_enable_transaction目前仅支持true,该设置同时决定了文件名会携带${transactionId}_前缀;
  6. 集群部署:若使用自定义凭据提供程序类,请确保其 JAR 存在于每个集群节点的${SEATUNNEL_HOME}/lib下,而非仅提交作业的节点(参考 S3File 文档的说明)。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询