数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnel 的S3RedshiftSink 连接器解决的是一个典型场景把海量数据先写入 AWS S3 对象存储再通过 Redshift 的COPY命令把 S3 上的数据文件批量加载进 Amazon Redshift 数据仓库。本文以 docs/en/connector-v2/sink/S3-Redshift.md 为主体结合仓库源码与示例配置完整讲解连接器的原理、全部配置参数、三类文件格式的完整配置模板以及基于 2PC 的 exactly-once 提交机制读者看完即可在 SeaTunnel 作业中落地一套「S3 中转 COPY 导入」的 Redshift 数据同步方案。工作原理一次写入两段接力S3Redshift并不是直接通过 JDBC 逐行写入 Redshift而是采用「接力式」设计第一阶段写 S3Sink 把上游数据按配置的文件格式text/csv/parquet/orc/json写出到 S3 指定路径下的临时目录第二阶段COPY 导入在事务提交阶段连接器把临时文件移动到目标路径然后通过 JDBC 执行用户配置的execute_sql这条 SQL 通常是 Redshift 的COPY ... FROM s3://...语句由 Redshift 服务端直接从 S3 批量加载数据收尾清理COPY 成功后删除已加载的数据文件并清理事务目录避免重复加载。正是由于这种设计S3Redshift在实现上完全基于 S3 文件 SinkS3File扩展而来文档明确指出S3File 的配置可以原样复用。为了支持多种文件类型连接器内部使用 HDFS 协议访问 S3s3a因此需要 Hadoop 相关依赖且只支持 Hadoop2.6.5版本。核心特性[x] exactly-once默认通过两阶段提交2PC机制保证数据不丢不重详见 connector-v2-features[x] 文件格式支持text、csv、parquet、orc、json五种格式与 S3File 相比S3Redshift 文档列出的格式以 Redshift COPY 支持范围为基准。运行环境与依赖Hadoop 版本要求2.6.5连接器通过hadoop-aws相关组件访问 S3需要 Redshift JDBC 驱动。从 pom.xml 可以看到连接器依赖com.amazon.redshift:redshift-jdbc42:2.1.0.9而 RedshiftJdbcClient.java 中通过Class.forName(com.amazon.redshift.jdbc42.Driver)显式加载驱动模块本身依赖connector-file-base-hadoop与connector-file-s3见 pom.xml因此 S3 侧的 s3a 访问能力与 S3File 完全一致。配置选项总览名称类型必填默认值jdbc_urlstring是-jdbc_userstring是-jdbc_passwordstring是-execute_sqlstring是-pathstring是-bucketstring是-access_keystring否-access_secretstring否-hadoop_s3_propertiesmap否-file_name_expressionstring否${transactionId}file_format_typestring否textfilename_time_formatstring否yyyy.MM.ddfield_delimiterstring否\001row_delimiterstring否\npartition_byarray否-partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/is_partition_field_write_in_fileboolean否falsesink_columnsarray否为空时写入全部字段is_enable_transactionboolean否truebatch_sizeint否1000000common-options-否-值得注意的是从 S3RedshiftSink.java 的prepare()校验逻辑看作业启动时必须存在以下键bucket、fs.s3a.aws.credentials.provider、jdbc_url、jdbc_user、jdbc_password、execute_sql缺少任一键都会抛出CONFIG_VALIDATION_FAILED异常并终止启动。也就是说即使文档表格中某些 S3 凭证参数标记为可选fs.s3a.aws.credentials.provider仍需要在配置中显式声明。Redshift 连接与导入 SQL 参数jdbc_url连接 Redshift 数据库的 JDBC URL例如jdbc:redshift://xxx.amazonaws.com.cn:5439/xxx。jdbc_user / jdbc_password连接 Redshift 数据库的用户名与密码。三者组合在 RedshiftJdbcClient.java 中以单例方式创建 JDBC 连接DriverManager.getConnection(url, user, password)并在 Sink 关闭时统一释放。execute_sql数据写入 S3 之后要执行的 SQL这是整个连接器的核心参数。典型写法COPY target_table FROM s3://yourbucket${path} IAM_ROLE arn:XXX REGION your region format as json auto;其中target_table是 Redshift 中的目标表名${path}是数据文件实际写入 S3 的路径占位符必须保留且不要手动替换——连接器会在执行 SQL 时自动替换为真实路径IAM_ROLE是具备 S3 访问权限的角色需要确保该角色有权读取对应的 S3 目录format必须与配置中file_format_type保持一致否则 Redshift COPY 会解析失败。从源码看替换动作发生在 S3RedshiftSinkAggregatedCommitter.java 的convertSql()方法中即StringUtils.replace(executeSql, ${path}, path)替换后的 SQL 通过RedshiftJdbcClient.getInstance(pluginConfig).execute(sql)执行。这里有个容易踩的坑如果你的 Redshift 集群与 S3 不在同一区域或使用了自定义 endpoint请确保REGION参数与实际 S3 区域一致。S3 存储与凭证参数pathS3 上的目标目录路径必填。数据文件最终会写入bucket path组合出的目录。bucketS3 文件系统的 bucket 地址。文档给出的示例是s3n://seatunnel-test如果使用s3a协议则应写为s3a://seatunnel-test。从 S3ConfigOptions.java 可以看出该键在代码中的定义名为bucket。access_key / access_secretS3 访问密钥对。若不设置请确保 Hadoop 凭证提供链credential provider chain可以正确完成认证。需要特别留意一处文档与源码的差异选项表格中列出的是access_secret而官方示例与源码 S3ConfigOptions.java 中实际使用的键名是secret_key配置时请以secret_key为准。hadoop_s3_properties当默认配置无法满足需求时可以通过该 map 追加任意 Hadoop-AWS 选项例如指定凭证提供类hadoop_s3_properties { fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider }从 S3ConfigOptions.java 的注释可以看到hadoop_s3_properties被定义为mapType()可承载fs.s3a.session.token等任意 s3a 配置键。关于凭证提供方式源码级补充S3ConfigOptions.java 将fs.s3a.aws.credentials.provider定义为枚举选项默认值是com.amazonaws.auth.InstanceProfileCredentialsProviderEC2 实例角色方式同时支持org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProviderAK/SK 方式。两种方式二选一InstanceProfileCredentialsProvider适合作业运行在 EC2/EMR 等 AWS 环境内通过实例角色自动获取临时凭证无需在配置里暴露密钥SimpleAWSCredentialsProvider显式使用access_key与secret_key适合在本地或自建环境运行。另外S3ConfigOptions还定义了fs.s3a.endpoint选项S3ConfigOptions.java当需要对接 S3 兼容对象存储如自建 MinIO 等时可以通过fs.s3a.endpoint指定地址。文件写入参数file_name_expression描述写入path的文件名表达式支持内置变量${now}当前时间与${uuid}例如test_${uuid}_${now}。${now}的格式由filename_time_format控制。注意当is_enable_transaction为true时连接器会自动在文件名头部追加${transactionId}_前缀此时默认值${transactionId}配合事务机制使用。file_format_type支持text、csv、parquet、orc、json。最终文件名会以对应格式后缀结尾其中 text 文件的后缀是txt。注意field_delimiter、row_delimiter仅对 text/csv 格式生效而 parquet/orc 是列式存储字段分隔符对其无意义。filename_time_format当file_name_expression中使用了xxxx-${now}时该参数控制时间部分的格式默认yyyy.MM.dd。常用时间符号符号含义y年M月d日H小时0-23m分钟s秒field_delimiter / row_delimiter字段分隔符与行分隔符仅 text/csv 格式需要。默认值分别为\001SOH 控制字符常用于规避数据内容冲突与\n。如果数据本身包含\n或|等字符请调整分隔符或在execute_sql的 COPY 语句中使用removequotes、escape等选项配合。sink_columns指定需要写入文件的列顺序即文件实际写入顺序。为空时写入Transform或Source输出的全部字段。分区参数partition_by按所选字段进行分区配合partition_dir_expression生成分区目录。partition_dir_expression指定分区目录的生成表达式默认值为${k0}${v0}/${k1}${v1}/.../${kn}${vn}/其中k0是第一个分区字段名v0是其值。例如partition_by[dt]且partition_dir_expression${k0}${v0}时会生成dt2024-01-01/形式的目录最终文件落在该分区目录内。is_partition_field_write_in_file为true时分区字段及其值也会写入数据文件内部为false时分区字段仅体现在目录结构中。若要生成 Hive 可直接识别的数据文件请设为falseHive 约定分区列不落文件内容。事务与批量参数is_enable_transaction为true时连接器通过事务机制保证写入目标目录的数据不丢失、不重复同时自动在文件名头部追加${transactionId}_。文档明确说明当前仅支持true即事务能力是强制开启的。batch_size单个文件的最大行数默认 1000000。在 SeaTunnel Engine 中文件行数由batch_size与checkpoint.interval共同决定若checkpoint.interval足够大Sink Writer 会持续写入直到行数超过batch_size才换文件若checkpoint.interval较小则每次 checkpoint 触发时都会生成新文件行数通常小于batch_size。common-optionsSink 插件通用参数例如source_table_name详见 Sink Common Options。完整配置示例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 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_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 } }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 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_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 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_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 } }以上示例均可在文档 S3-Redshift.md 中找到对应版本。注意 text 示例的execute_sql中同时使用了removequotes emptyasnull blanksasnull maxerror 100 delimiter |而field_delimiter保持默认\001时 COPY 的delimiter需要与之匹配实际使用时建议将两者显式对齐例如都设为|。源码实现剖析提交与回滚如何保证 exactly-onceS3Redshift的提交逻辑集中在 S3RedshiftSinkAggregatedCommitter.java它继承自FileSinkAggregatedCommitter。整体流程如下写临时文件Sink 先向tmp_path示例中为/tmp/seatunnel写临时文件SeaTunnel Engine 的 checkpoint 会记录事务状态提交阶段commit遍历事务目录中的每个文件——先通过hadoopFileSystemProxy.renameFile()把临时文件重命名移动到目标 S3 路径然后convertSql()用真实路径替换${path}并执行 COPY SQL成功后再删除已加载的 S3 文件、删除整个事务目录见 commit 方法失败回滚abort如果任务失败abort()直接删除对应的事务目录S3 上不会残留半成品数据见 abort 方法资源释放close()时关闭 Hadoop 文件系统代理与 Redshift JDBC 连接。这套「先 rename 到目标路径 → 执行 COPY → 删除源文件」的编排保证了只有 COPY 真正执行成功文件才从临时区移出一旦任何一步失败事务目录整体删除数据不会重复入库。这也是连接器能声明 exactly-once 的根基。另外RedshiftJdbcClient.java 还提供了checkTableExists()方法用于校验 Redshift 目标表是否存在连接失败或 SQL 执行失败时会抛出S3RedshiftJdbcConnectorException并携带明确的错误码。常见问题与注意事项${path}必须保留execute_sql中缺失${path}时COPY 语句会指向错误或不完整的 S3 路径导致导入失败。连接器只做简单的字符串替换StringUtils.replace不会帮你补路径。format与file_format_type必须一致例如file_format_type json时COPY 语句应为format as json autoparquet/orc 同理。IAM 权限IAM_ROLE或凭证需要同时具备对 S3 目标桶的s3:GetObject读取权限以及 Redshift 集群对 S3 的访问授权否则 COPY 会以AccessDenied失败。Hadoop 版本约束连接器只支持 Hadoop 2.6.5使用 SeaTunnel Engine 时需确认${SEATUNNEL_HOME}/lib下包含hadoop-aws与aws-java-sdk-bundle相关 jar详见 S3File 文档。密钥参数键名配置中请使用secret_key源码与示例均为该键名避免误用选项表格中的access_secret写法导致认证缺失。事务强制开启is_enable_transaction目前只支持true不要尝试关闭事务否则与 COPY 导入模型的配合会失去保障。batch_size 与 checkpoint 的关系若发现文件数量过多或单个文件过小优先检查checkpoint.interval的配置期望大文件则调大 checkpoint 间隔并让batch_size生效。Changelog2.3.0-beta 2022-10-20该版本为S3Redshift连接器的初始发布版本正式提供了上述「S3 写文件 Redshift COPY 导入」的完整能力。更详细的 COPY 语法如removequotes、emptyasnull、blanksasnull、maxerror、delimiter等选项的语义可参考 Redshift 官方文档对 COPY 命令的说明S3 凭证链的完整认证机制可参考 Hadoop-AWS 官方文档。如需了解该连接器在 SeaTunnel 整体架构中的位置可参阅 docs/en/about.md。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐使用 SeaTunnel S3Redshift Sink先写 S3 再以 COPY 命令批量导入 Amazon Redshift使用 SeaTunnel S3Redshift Sink先写 S3 再以 COPY 命令批量导入 Amazon Redshift 本文聚焦 Apache Se数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel S3Redshift Sink 连接器实战基于 S3 Redshift COPY 的海量数据导入方案SeaTunnel S3Redshift Sink 连接器实战基于 S3 Redshift COPY 的海量数据导入方案 S3Redshift 是 Sea数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel JDBC Redshift Sink 连接器实战从批量写入到 CDC 多表同步SeaTunnel JDBC Redshift Sink 连接器实战从批量写入到 CDC 多表同步 SeaTunnel 通过 connector jdbc 提数据集成ETL大数据批处理流处理变更数据捕获上一篇如何用Gemma4-26B-A4B去审查模型实现465个拒绝点零突破下一篇Fluor核心原理深度解析从应用识别到Fn键切换的完整流程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考