FEATURED · 精选文章

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

发布时间 / 2026/9/17 6:03:05
来源 / 创域科博编辑部
栏目 / 资讯中心
使用 SeaTunnel S3Redshift Sink:先写 S3 再以 COPY 命令批量导入 Amazon Redshift 使用 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 Redshifts 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()方法对每个聚合提交信息执行先把临时文件renameFile到最终路径用最终文件路径替换 SQL 模板中的${path}占位符并构造出COPY语句通过RedshiftJdbcClient.getInstance(pluginConfig).execute(sql)执行 COPY数据装载成功后删除已处理的 S3 文件并删除事务目录。而abort()方法则只负责清理事务目录保证失败时不留残留文件。这套临时文件 → rename → COPY → 清理的顺序正是数据不丢不重的落地关键COPY 只有在文件完整就位后才触发一旦失败事务被中止已写入的文件与目录都会被清理。文件格式支持支持以下文件格式写入 S3 的文件格式必须与你在execute_sql的 COPY 语句中声明的format保持一致textcsvparquetorcjson完整参数说明下表汇总了 S3Redshift 的全部配置项。其中前 4 个参数JDBC 与 SQL 相关是该连接器在 S3File 基础上的新增参数其余均继承自文件/S3 连接器体系。参数名类型必填默认值说明jdbc_urlstring是-连接 Redshift 数据库的 JDBC URL例如jdbc:redshift://your-cluster.region.redshift.amazonaws.com:5439/your_databasejdbc_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.providerfile_name_expressionstring否${transactionId}拼接在path下的文件名表达式。支持${now}、${uuid}变量当is_enable_transaction true时自动在文件名前加${transactionId}_file_format_typestring否text写入 S3 的文件格式支持text、csv、parquet、orc、json最终文件名会带相应后缀text 后缀为txtfilename_time_formatstring否yyyy.MM.dd解析file_name_expression中${now}的时间格式语法遵循 JavaSimpleDateFormatfield_delimiterstring否\001text与csv文件的列分隔符row_delimiterstring否\ntext与csv文件的行分隔符partition_byarray否-按列出的上游字段进行分区partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/依据partition_by字段推导分区目录结构的表达式is_partition_field_write_in_fileboolean否false为true时分区字段及其取值除了体现在目录结构中外还会写入数据文件本身写 Hive 风格数据文件应设为falsesink_columnsarray否默认上游全字段需要写入文件的列字段顺序即实际写文件时的列顺序is_enable_transactionboolean否true为true时保证写入目标目录的数据不丢不重当前仅支持truebatch_sizeint否1000000单个文件的最大行数。在 SeaTunnel Zeta 引擎中文件行数由batch_size与checkpoint.interval共同决定common-options-否-Sink 插件通用参数详见 Sink Common Options参数校验规则源码佐证这些参数的必填/联动规则并非仅仅写在文档里而是由 S3RedshiftSinkFactory.optionRule() 在启动阶段强制执行bucket、jdbc_url、jdbc_user、jdbc_password、execute_sql为必填项且 4 个 JDBC 相关参数均要求非空白Conditions.notBlankpath必填fs.s3a.aws.credentials.provider必填当凭证提供器选择org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider时才要求提供access_key与secret_key即S3_ACCESS_KEY/S3_SECRET_KEY联动必填当file_format_type为text时field_delimiter与row_delimiter必须配置为csv时row_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_sqlexecute_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 官方文档为准。pathbucket下的目标目录路径必填。最终文件会写到bucket path 文件名的位置同时该路径会以${path}形式进入execute_sql作为 COPY 的数据源前缀。bucketS3 文件系统的 bucket 地址例如s3n://seatunnel-test若使用s3a协议则应写成s3a://seatunnel-test。由于本连接器内部通过 Hadoop S3A 文件系统访问 S3文档与源码都明确推荐使用s3a协议。access_key / access_secretS3 文件系统的访问凭证。若未配置则必须确保 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 与文件名后缀支持text、csv、parquet、orc、json五种格式。最终文件名会以对应后缀结尾其中text 文件的后缀是txt。filename_time_format当file_name_expression中包含${now}时用此参数指定时间格式默认yyyy.MM.dd。常用时间格式符号如下符号含义y年YearM月Monthd日Day of monthH小时0-23m分钟Minute in hours秒Second in minute完整语法遵循 JavaSimpleDateFormat规范。field_delimiter / row_delimiter分别指定一列数据内部的列分隔符与文件内的行分隔符仅对text和csv格式生效。默认值分别是\001SOH 控制字符与\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指定需要写入文件的列默认是来自Transform或Source的全部字段。列的顺序决定了文件实际写入的列顺序因此在使用 COPY 装载时应保证与目标表列顺序或COPY的列清单对齐。is_enable_transaction为true时保证写入目标目录的数据不丢不重同时自动在文件名前追加${transactionId}_。当前仅支持true即该特性默认开启且无法关闭。batch_size单个文件的最大行数默认1000000一百万行。在 SeaTunnel EngineZeta中文件内行数由batch_size与checkpoint.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_sqlCOPY 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_filetrue file_name_expression${transactionId}_${now} file_format_type text filename_time_formatyyyy.MM.dd is_enable_transactiontrue hadoop_s3_properties { fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider } }该示例中 COPY 语句的delimiter |与文本分隔符保持一致同时启用了removequotes、emptyasnull、blanksasnull等文本清洗选项并设置了maxerror 100容忍部分坏行。parquet 格式S3Redshift { jdbc_url jdbc:redshift://xxx.amazonaws.com.cn:5439/xxx jdbc_user xxx jdbc_password xxxx execute_sqlCOPY 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_filetrue file_name_expression${transactionId}_${now} file_format_type parquet filename_time_formatyyyy.MM.dd is_enable_transactiontrue 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_sqlCOPY 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_filetrue file_name_expression${transactionId}_${now} file_format_type orc filename_time_formatyyyy.MM.dd is_enable_transactiontrue 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 path有s3: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.mdSink 插件通用参数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),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻