使用 SeaTunnel S3Redshift Sink:先写 S3 再以 COPY 命令批量导入 Amazon Redshift
2026/9/17 6:03:04 网站建设 项目流程

使用 SeaTunnel S3Redshift Sink:先写 S3 再以 COPY 命令批量导入 Amazon Redshift

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

本文聚焦 Apache SeaTunnel 的S3RedshiftSink 连接器,讲解其"先写对象存储、再执行 Redshift COPY 导入"的两阶段写数原理、完整参数体系、exactly-once 事务提交机制,以及 text/parquet/orc 等文件格式的实战配置。读完本文,你将能够独立编写一条把任意上游数据经 S3 中转、最终落库到 Amazon Redshift 的 SeaTunnel 作业配置,并理解其底层提交链路为何能保证数据不丢不重。

连接器定位与工作原理

S3Redshift是 SeaTunnel Connector V2 体系中的一个 Sink 插件,其插件标识为S3Redshift。它的工作方式非常特殊:并不直接通过 JDBC 逐条写入 Redshift,而是先把数据批量写入 S3 对象存储,待文件落盘提交后,再通过 Redshift 的COPY命令把 S3 上的文件数据批量装载进目标表。

这种间接写入模式充分借力了 Redshift COPY 的高吞吐批量加载能力,同时把"分布式写文件"与"数据库装载"解耦,是海量数据入仓场景下的典型架构选择。从官方文档的说明可知:

The way of S3Redshift is to write data into S3, and then use Redshift's COPY command to import data from S3 to Redshift.

与 S3File 的关系

在连接器源码中可以看到,S3RedshiftSink直接继承自文件连接器的BaseFileSink,其 Hadoop 配置由S3HadoopConf.buildWithReadOnlyConfig(pluginConfig)构建,且createAggregatedCommitter()返回了自定义的S3RedshiftSinkAggregatedCommitter。也就是说:

  • 该连接器完全基于 S3File 实现,S3File 的所有配置参数对它同样生效;
  • 为了支持更多文件类型,内部访问 S3 使用了HDFS 协议(即 Hadoop S3A 文件系统),因此依赖若干 Hadoop 组件;
  • 仅支持 Hadoop 2.6.5+ 版本,低于该版本无法正常工作。

版本演进

从连接器 Changelog 可以确认,"先写 S3 再 COPY 到 Redshift"的能力自2.3.0版本引入,后续在 2.3.1 中修复了 S3 文件复制到 Redshift 时"文件找不到"以及 S3 路径不正确等缺陷,2.3.6 起支持多表写入。因此生产环境建议使用 2.3.6 及以上的 SeaTunnel 版本。

关键特性

精确一次(exactly-once)

该连接器默认启用is_enable_transaction = true,基于2PC(两阶段提交)机制保证写入 Redshift 的数据不丢失、不重复,对应 SeaTunnel 的 exactly-once 能力说明。

其事务语义在 S3RedshiftSinkAggregatedCommitter 中体现得非常直白,commit()方法对每个聚合提交信息执行:

  1. 先把临时文件renameFile到最终路径;
  2. 用最终文件路径替换 SQL 模板中的${path}占位符并构造出COPY语句;
  3. 通过RedshiftJdbcClient.getInstance(pluginConfig).execute(sql)执行 COPY;
  4. 数据装载成功后删除已处理的 S3 文件,并删除事务目录。

abort()方法则只负责清理事务目录,保证失败时不留残留文件。这套"临时文件 → rename → COPY → 清理"的顺序,正是"数据不丢不重"的落地关键:COPY 只有在文件完整就位后才触发,一旦失败事务被中止,已写入的文件与目录都会被清理。

文件格式支持

支持以下文件格式,写入 S3 的文件格式必须与你在execute_sql的 COPY 语句中声明的format保持一致:

  • text
  • csv
  • parquet
  • orc
  • json

完整参数说明

下表汇总了 S3Redshift 的全部配置项。其中前 4 个参数(JDBC 与 SQL 相关)是该连接器在 S3File 基础上的新增参数,其余均继承自文件/S3 连接器体系。

参数名类型必填默认值说明
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}变量;当is_enable_transaction = true时自动在文件名前加${transactionId}_
file_format_typestringtext写入 S3 的文件格式,支持textcsvparquetorcjson;最终文件名会带相应后缀(text 后缀为txt
filename_time_formatstringyyyy.MM.dd解析file_name_expression${now}的时间格式,语法遵循 JavaSimpleDateFormat
field_delimiterstring\001textcsv文件的列分隔符
row_delimiterstring\ntextcsv文件的行分隔符
partition_byarray-按列出的上游字段进行分区
partition_dir_expressionstring${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/依据partition_by字段推导分区目录结构的表达式
is_partition_field_write_in_filebooleanfalsetrue时,分区字段及其取值除了体现在目录结构中外,还会写入数据文件本身;写 Hive 风格数据文件应设为false
sink_columnsarray默认上游全字段需要写入文件的列,字段顺序即实际写文件时的列顺序
is_enable_transactionbooleantruetrue时保证写入目标目录的数据不丢不重,当前仅支持true
batch_sizeint1000000单个文件的最大行数。在 SeaTunnel Zeta 引擎中,文件行数由batch_sizecheckpoint.interval共同决定
common-options--Sink 插件通用参数,详见 Sink Common Options

参数校验规则(源码佐证)

这些参数的必填/联动规则并非仅仅写在文档里,而是由 S3RedshiftSinkFactory.optionRule() 在启动阶段强制执行:

  • bucketjdbc_urljdbc_userjdbc_passwordexecute_sql为必填项,且 4 个 JDBC 相关参数均要求非空白Conditions.notBlank);
  • path必填,fs.s3a.aws.credentials.provider必填;
  • 当凭证提供器选择org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider时,才要求提供access_keysecret_key(即S3_ACCESS_KEY/S3_SECRET_KEY联动必填);
  • file_format_typetext时,field_delimiterrow_delimiter必须配置;为csvrow_delimiter必须配置。

单元测试 S3RedshiftSinkFactoryTest 验证了这些规则:合法的完整配置不抛异常,而将jdbc_url/jdbc_user/jdbc_password/execute_sql任一置为空白字符串时,均会抛出OptionValidationException。这意味着配置错误会在作业启动阶段即被拦截,而不是运行到一半才报错。

核心参数详解

jdbc_url / jdbc_user / jdbc_password

三者共同描述 Redshift 数据库连接。连接由 RedshiftJdbcClient 负责建立:它通过Class.forName("com.amazon.redshift.jdbc42.Driver")加载 Redshift JDBC 驱动(版本见 connector-s3-redshift/pom.xml,使用redshift-jdbc422.1.0.30),再以DriverManager.getConnection(url, user, password)建立单例连接。该连接仅用于执行 COPY 等装载 SQL,并不承担逐行数据写入。

execute_sql

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 中,无需也不应手动替换——提交阶段连接器会自动完成替换。在源码中,S3RedshiftSinkAggregatedCommitter.convertSql()正是通过StringUtils.replace(executeSql, "${path}", path)${path}替换为mvFileEntry.getValue()(即最终文件路径);
  • IAM_ROLE是具备 S3 访问权限的 IAM 角色,请确认该角色确有权限读取对应 bucket 与路径;
  • format必须与file_format_type中设置的文件格式一致,否则 COPY 会因格式不匹配而失败;
  • 更多 COPY 语法细节可参考 Redshift COPY 官方文档(本文不展开外部链接,请以 AWS 官方文档为准)。

path

bucket下的目标目录路径,必填。最终文件会写到bucket + path + 文件名的位置,同时该路径会以${path}形式进入execute_sql,作为 COPY 的数据源前缀。

bucket

S3 文件系统的 bucket 地址,例如s3n://seatunnel-test;若使用s3a协议则应写成s3a://seatunnel-test。由于本连接器内部通过 Hadoop S3A 文件系统访问 S3,文档与源码都明确推荐使用s3a协议

access_key / access_secret

S3 文件系统的访问凭证。若未配置,则必须确保 Hadoop 凭证提供链(credential provider chain)能够正确完成认证。从 S3FileBaseOptions 的源码可见,SeaTunnel 内置了两类常见凭证提供器常量:

  • org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider:使用静态access_key/secret_key对认证;
  • com.amazonaws.auth.InstanceProfileCredentialsProvider:从运行环境(如 EC2 实例角色)解析凭证,这也是fs.s3a.aws.credentials.provider的默认值。

hadoop_s3_properties

当需要补充其他 Hadoop S3A/Hadoop-AWS 配置时,可以统一放在这个 map 中。例如:

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

file_name_expression

描述在path下创建的文件名表达式,支持${now}${uuid}两个变量,例如test_${uuid}_${now}${now}表示当前时间,其格式由filename_time_format控制。

特别注意:当is_enable_transaction = true时,连接器会自动在文件名头部加上${transactionId}_,以区分不同事务批次产出的文件。

file_format_type 与文件名后缀

支持textcsvparquetorcjson五种格式。最终文件名会以对应后缀结尾,其中text 文件的后缀是txt

filename_time_format

file_name_expression中包含${now}时,用此参数指定时间格式,默认yyyy.MM.dd。常用时间格式符号如下:

符号含义
y年(Year)
M月(Month)
d日(Day of month)
H小时(0-23)
m分钟(Minute in hour)
s秒(Second in minute)

完整语法遵循 JavaSimpleDateFormat规范。

field_delimiter / row_delimiter

分别指定一列数据内部的列分隔符与文件内的行分隔符,仅对textcsv格式生效。默认值分别是\001(SOH 控制字符)与\n,这两个默认值在 FileBaseSinkOptions 中定义。

partition_by / partition_dir_expression

按所选字段对数据进行分区。指定partition_by后,连接器会根据分区信息生成对应分区目录,最终文件放置在分区目录内。默认partition_dir_expression${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/,其中k0是第一个分区字段名,v0是其取值。需要注意的是,分区目录结构同样会影响 COPY 语句读取的路径范围。

is_partition_field_write_in_file

true时,分区字段及其取值除了体现为目录结构,还会被写入数据文件本身。如果目标是生成 Hive 风格的数据文件,此值应设为false(分区字段只体现在目录,不重复写进文件内容)。

sink_columns

指定需要写入文件的列,默认是来自TransformSource的全部字段。列的顺序决定了文件实际写入的列顺序,因此在使用 COPY 装载时应保证与目标表列顺序或COPY的列清单对齐。

is_enable_transaction

true时保证写入目标目录的数据不丢不重,同时自动在文件名前追加${transactionId}_当前仅支持true,即该特性默认开启且无法关闭。

batch_size

单个文件的最大行数,默认1000000(一百万行)。在 SeaTunnel Engine(Zeta)中,文件内行数由batch_sizecheckpoint.interval共同决定

  • checkpoint.interval足够大,sink writer 会持续向同一文件写行,直到文件行数超过batch_size才滚动新文件;
  • checkpoint.interval较小,则每次 checkpoint 触发时 sink writer 就会创建新文件。

这一行为意味着:文件大小并不严格等于batch_size行,而是"行数阈值"与"检查点周期"二者取先到者。

配置示例

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" access_secret = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" 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" } }

该示例中 COPY 语句的delimiter '|'与文本分隔符保持一致,同时启用了removequotesemptyasnullblanksasnull等文本清洗选项,并设置了maxerror 100容忍部分坏行。

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" access_secret = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" 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" access_secret = "xxxxxxxxxxxxxxxxx" bucket = "s3a://seatunnel-test" 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" } }

parquet 与 orc 为列式存储格式,COPY 语句中通过format as PARQUET;/format as ORC;声明即可,无需再指定文本类分隔符。json 与 csv 格式的配置方式同理,只需将file_format_type与 COPY 的format保持一致。

作业运行方式

上述 sink 配置需要放置在一个完整的 SeaTunnel 作业配置文件中(Source + Transform + Sink 三段式 HOCON 结构),例如:

env { execution.parallelism = 2 job.mode = "BATCH" } source { FakeSource { schema { fields { id = bigint name = string } } row.num = 1000 } } sink { S3Redshift { # 上文的 S3Redshift 配置 } }

然后通过 SeaTunnel 命令行提交(以本仓库 seatunnel-starter 提供的启动器为例,具体脚本名以你所使用发行版的 bin 目录与 seatunnel-env.sh 为准):

bin/seatunnel.sh --config <your-job.conf> -e local

注意:该连接器需要 Hadoop 2.6.5+ 运行环境,且redshift-jdbc42驱动、Hadoop-AWS 相关依赖需随连接器一并分发到运行节点,请确认插件目录与plugin_config配置完整后再提交作业。

常见问题排查思路

  • COPY 报文件找不到:确认execute_sql中的${path}占位符未被手动替换、bucket使用s3a://协议,且 IAM 角色对bucket + paths3:GetObject/s3:ListBucket权限(该问题在 2.3.1 版本曾专门修复,见 Changelog)。
  • 格式不匹配报错:检查file_format_type与 COPY 语句format是否一致(如 parquet 写 text)。
  • 凭证认证失败:确认access_key/access_secret已配置,或fs.s3a.aws.credentials.provider指向正确的凭证提供器类;默认的InstanceProfileCredentialsProvider依赖运行环境(如 EC2 实例角色)。
  • 作业启动即失败:先核对必填参数是否齐全且非空白——工厂层的OptionRule会在启动阶段直接抛出OptionValidationException

延伸阅读

  • 基础文件 Sink 配置:docs/en/connectors/sink/S3File.md
  • Sink 插件通用参数:Sink Common Options
  • 连接器能力说明(exactly-once 等):connector-v2-features
  • 连接器实现源码:connector-s3-redshift 模块

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

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

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

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

立即咨询