ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

SeaTunnel S3Redshift Sink 连接器实战:基于 S3 + Redshift COPY 的海量数据导入方案 SeaTunnel S3Redshift Sink 连接器实战基于 S3 Redshift COPY 的海量数据导入方案【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelS3Redshift 是 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直接继承BaseFileSinkS3RedshiftSink.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。文件格式类型textcsvparquetorcjson定时刷新暂不支持。精确一次的实现原理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_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 的文件格式支持text、csv、parquet、orc、json。最终文件名带相应后缀text的后缀为txt。filename_time_formatstring否yyyy.MM.dd解析file_name_expression中${now}的时间格式。field_delimiterstring否\001text和csv文件的列分隔符。row_delimiterstring否\ntext和csv文件的行分隔符。partition_byarray否-按指定的上游字段对数据进行分区分区目录由partition_dir_expression推导。partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/根据partition_by字段生成分区目录的表达式。is_partition_field_write_in_fileboolean否false为true时分区字段及其值写入数据文件。Hive 风格数据文件请设为false。sink_columnsarray否为空时所有字段都是 sink 列需要写入文件的列字段顺序决定文件实际写入顺序。is_enable_transactionboolean否true为true时保证数据写入目标目录不丢失、不重复。目前只支持true。batch_sizeint否1000000单个文件的最大行数。在 SeaTunnel Zeta 引擎中每文件行数由batch_size与checkpoint.interval共同决定。common-options-否-Sink 插件通用参数见 Sink 通用选项。参数校验规则来自源码从 S3RedshiftSinkFactory.java 的optionRule()可以看到连接器启动时的强校验逻辑必填参数bucket、jdbc_url、jdbc_user、jdbc_password、execute_sql、pathFILE_PATH、fs.s3a.aws.credentials.providerS3A_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 命令文档。文件与路径相关参数继承自 S3Filepath [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_keysecret_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文件格式需要默认\001row_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 optionsSink 插件通用参数详见 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 } }示例配置要点提示tmp_path指定了临时目录示例为/tmp/seatunnel文件会先写入该临时路径提交时再mv到path目标目录这是 2PC 提交的一部分三个示例均使用SimpleAWSCredentialsProvideraccess_key/secret_key的静态凭据方式与 S3File 文档中推荐的凭据提供程序链保持兼容分区表达式partition_dir_expression${k0}${v0}结合is_partition_field_write_in_filetrue会把分区字段写入数据文件内方便按分区管理 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),仅供参考
返回列表