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_type、partition_by、sink_columns等)都可直接使用;在此之上叠加了 Redshift 的 JDBC 连接与 COPY SQL 执行能力。
为了支持更多文件类型,S3Redshift 使用 HDFS 协议对 S3 进行内部访问,因此该连接器需要一些 Hadoop 依赖,且只支持 Hadoop 版本2.6.5+。对应的依赖在 pom.xml 中体现为connector-file-base-hadoop、connector-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流程为:
- 遍历每个事务的文件映射(
transactionMap),先把临时文件renameFile到目标路径; - 通过
convertSql(mvFileEntry.getValue())把execute_sql中的${path}占位符替换为实际文件路径(见convertSql实现StringUtils.replace(executeSql, "${path}", path)); - 调用 RedshiftJdbcClient 的
execute(sql)执行这条 COPY 语句,将 S3 文件导入 Redshift; - 文件导入成功后删除已导入的文件,并清理事务目录。
而abort阶段则直接删除事务目录,保证失败时不会在目标路径留下半成品数据。整个链路把「S3 文件可见」与「Redshift 数据可见」绑定在同一个提交事务内,从而实现不丢失、不重复。
参数详解
S3Redshift 的参数由两部分组成:Redshift 专属参数(定义于 S3RedshiftSinkOptions.java)和从 S3File 继承的文件参数。完整参数表如下:
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| jdbc_url | string | 是 | - | 连接 Redshift 数据库的 JDBC URL,例如jdbc:redshift://your-cluster.region.redshift.amazonaws.com:5439/your_database。 |
| jdbc_user | string | 是 | - | 连接 Redshift 数据库的用户名。 |
| jdbc_password | string | 是 | - | 连接 Redshift 数据库的密码。 |
| execute_sql | string | 是 | - | 数据写入 S3 之后要执行的 SQL,通常是一条 RedshiftCOPY命令(必须包含${path}占位符)。 |
| path | string | 是 | - | bucket 下的目标目录路径,连接器会通过${path}占位符把实际写入路径追加到execute_sql中。 |
| bucket | string | 是 | - | S3 文件系统的 bucket 地址,例如s3a://seatunnel-test。使用 Hadoop 读写时建议使用s3a协议。 |
| access_key | string | 否 | - | S3 文件系统的 access key。如果未配置,需要正确配置 Hadoop 凭据链。 |
| access_secret | string | 否 | - | S3 文件系统的 access secret。如果未配置,需要正确配置 Hadoop 凭据链。 |
| hadoop_s3_properties | map | 否 | - | 额外的 Hadoop S3A / Hadoop-AWS 选项,可设置fs.s3a.aws.credentials.provider等。 |
| file_name_expression | string | 否 | "${transactionId}" | 在path下追加的文件名表达式,可使用${now}或${uuid}注入时间或 UUID。is_enable_transaction = true时自动在文件名前添加${transactionId}_。 |
| file_format_type | string | 否 | "text" | 写入 S3 的文件格式,支持text、csv、parquet、orc、json。最终文件名带相应后缀(text的后缀为txt)。 |
| filename_time_format | string | 否 | "yyyy.MM.dd" | 解析file_name_expression中${now}的时间格式。 |
| field_delimiter | string | 否 | '\001' | text和csv文件的列分隔符。 |
| row_delimiter | string | 否 | "\n" | text和csv文件的行分隔符。 |
| partition_by | array | 否 | - | 按指定的上游字段对数据进行分区,分区目录由partition_dir_expression推导。 |
| partition_dir_expression | string | 否 | "${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/" | 根据partition_by字段生成分区目录的表达式。 |
| is_partition_field_write_in_file | boolean | 否 | false | 为true时,分区字段及其值写入数据文件。Hive 风格数据文件请设为false。 |
| sink_columns | array | 否 | 为空时所有字段都是 sink 列 | 需要写入文件的列,字段顺序决定文件实际写入顺序。 |
| is_enable_transaction | boolean | 否 | true | 为true时保证数据写入目标目录不丢失、不重复。目前只支持true。 |
| batch_size | int | 否 | 1000000 | 单个文件的最大行数。在 SeaTunnel Zeta 引擎中,每文件行数由batch_size与checkpoint.interval共同决定。 |
| common-options | - | 否 | - | Sink 插件通用参数,见 Sink 通用选项。 |
参数校验规则(来自源码)
从 S3RedshiftSinkFactory.java 的optionRule()可以看到连接器启动时的强校验逻辑:
- 必填参数:
bucket、jdbc_url、jdbc_user、jdbc_password、execute_sql、path(FILE_PATH)、fs.s3a.aws.credentials.provider(S3A_AWS_CREDENTIALS_PROVIDER_CLASS);其中四个 Redshift 参数还带有Conditions.notBlank约束,空白字符串无法通过校验; - 条件参数:当
fs.s3a.aws.credentials.provider为SimpleAWSCredentialsProvider时,必须提供access_key与secret_key;当file_format_type为text时必须提供field_delimiter与row_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.dir、fs.s3a.fast.upload.buffer、fs.s3a.session.token、fs.s3a.assumed.role.arn等(参考 S3File 文档中的相关说明)。
file_name_expression [string]
描述将在path中创建的文件名表达式。可以在表达式中加入变量${now}或${uuid},例如test_${uuid}_${now}。${now}表示当前时间,其格式由filename_time_format定义。
注意:如果is_enable_transaction为true,连接器会自动在文件名开头添加${transactionId}_。
file_format_type [string]
支持的写入文件类型:text、csv、parquet、orc、json。
注意,最终文件名会以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:数据行中列之间的分隔符,仅text和csv文件格式需要,默认'\001';row_delimiter:文件中行之间的分隔符,仅text和csv文件格式需要,默认"\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_size和checkpoint.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),文件会先写入该临时路径,提交时再mv到path目标目录,这是 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 变更日志。
使用前提与注意事项
- Hadoop 依赖:连接器通过 HDFS 协议访问 S3,需要 Hadoop 依赖,仅支持 Hadoop2.6.5+;
- 凭据配置:
access_key/access_secret未配置时,必须确保 Hadoop 凭据提供程序链能够正确完成 S3 身份验证(例如使用InstanceProfileCredentialsProvider等); - COPY 权限:
execute_sql中使用的 IAM_ROLE 必须同时具备读取 S3 中${path}对应对象的权限与写入 Redshift 目标表的权限; - 格式一致性:
execute_sql中 COPY 命令声明的format必须与file_format_type保持一致,否则 COPY 会因解析失败而报错; - 事务开关:
is_enable_transaction目前仅支持true,该设置同时决定了文件名会携带${transactionId}_前缀; - 集群部署:若使用自定义凭据提供程序类,请确保其 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),仅供参考