ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SeaTunnel S3File 连接器完全指南:从版本演进史到源码级配置实践

2026/9/16 15:47:18 拓冰建站 浏览量
SeaTunnel S3File 连接器完全指南:从版本演进史到源码级配置实践 SeaTunnel S3File 连接器完全指南从版本演进史到源码级配置实践【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelS3File 是 SeaTunnel 中面向 AWS S3以及兼容 S3 协议的对象存储的读写入门连接器同时承担 Source 与 Sink 双重角色。本文以其在 docs/en/connectors/changelog/connector-file-s3.md 中沉淀的版本演进记录为主线串联官方 S3File Source 文档、S3File Sink 文档 以及connector-file-s3模块源码完整讲解连接器的能力矩阵、文件格式支持、认证机制、配置项、典型作业示例与底层实现原理帮助读者在 Spark / Flink / SeaTunnel Zeta 引擎上正确构建 S3 数据同步管道。连接器概览一个连接器两种角色S3File 连接器覆盖s3://、s3n://、s3a://等协议可读取 AWS S3 文件系统中的数据Source也可将 SeaTunnel 管道中的数据写出到 S3Sink。从官方文档与源码工厂类可以看出其能力定位支持引擎Spark、Flink、SeaTunnel Zeta见 docs/en/connectors/source/S3File.md 的 Support Those Engines 一节。支持批batch与流stream两种模式且是 multimodal 连接器——使用 binary 文件格式即可同步任意格式的原始文件视频、图片等。Source 端具备 exactly-once 语义一个 split 内的数据在一次 pollNext 中读完已读 split 会保存在快照中、列投影column projection、并行度控制并支持丰富的文件格式。Sink 端通过默认的 2PC 提交机制保证 exactly-once支持多表写入multiple table write与多种 CDC 事件格式。在源码层面连接器的注册与参数校验集中在两个工厂类中S3FileSourceFactory与S3FileSinkFactory位于 seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java 与 seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/sink/S3FileSinkFactory.java它们通过OptionRule声明必填项、条件项与可选参数是理解哪些参数在什么场景下必填的第一手依据。版本演进时间线S3File 连接器两年来的功能沉淀changelog 文档以变更 | 提交 | 版本三列表格的形式记录了连接器从 2.3.0-beta 到 dev 分支的全部演进。这张表本身就是一份完整的功能能力地图下面按里程碑归纳版本关键变更能力意义2.3.0-beta首次加入 S3 file source sink connector连接器诞生2.3.0支持s3a协议Set S3 AK to optional访问密钥改为可选统一文件连接器异常新增 option 与 factory接入 Hadoop S3A 文件系统支持实例角色等免 AK 认证补齐 SPI 工厂机制2.3.1file_type更名为file_format_type支持压缩compress新增S3Catalog重构 schema 解析升级 guava 至 27.0-jre格式语义统一、压缩读、元数据目录能力上线2.3.2新增 excel sink 与 source删除不可用的 S3 Kafka CatalogsExcel 格式读写2.3.3新增file_filter_pattern用于按正则过滤文件精细化文件选择2.3.4支持 LZO 压缩读取支持读取空目录支持配置 schema 中的 column/primaryKey/constraintKeyENABLE_HEADER_WRITE参数统一 Source/Sink 选项多 Hadoop 账号支持s3file save mode 功能Multiple Table File API 重构进 file-base 模块批量能力增强与选项体系统一2.3.5修复 SPI 无参构造问题为 SFTP、FTP、LocalFile、HdfsFile 等文件连接器统一支持 XML 文件类型XML 格式落地2.3.6支持多表写入multiple table write支持 parquet 将 fixed/timestamp 写为 int96sink 选项中支持上游表占位符并自动替换修复 S3CONF 反序列化后被重新赋值的 bug多表 CDC 管道基础与 parquet 兼容性2.3.7增加多表 sink 选项检查多表配置错误前置暴露2.3.8支持读取归档压缩文件archive compress file重构 S3FileCatalog 及其工厂ZIP/TAR/TAR_GZ/GZ 归档读取2.3.9支持 text 文件读取的 null 格式配置修复 hadoop-aws 中 guava 与 hive-exec 的依赖冲突metrics 与逻辑计划节点关联null 语义自定义、依赖治理2.3.10支持单文件模式single file mode无数据时创建空文件新增filename_extension读写参数重构连接器通用选项修复某些场景下 s3 key 设置错误的问题文件命名与目录布局精细化2.3.11为 text file sink 增加row_delimiter更新文件连接器配置行分隔符可配置2.3.12为 text 文件处理增加可自定义行分隔符maxcompute sink writer 支持 timestamp 字段类型行分隔符完善dev新增 markdown 解析器RAG 场景的结构化文档抽取这张时间线揭示了一条清晰的技术主线连接器先解决能不能连上 S3再逐步补齐格式多样性、压缩与归档、多表与 CDC、文件过滤与命名、空文件与单文件等工程化细节。文章后续章节将围绕这些能力展开实操与源码佐证。环境依赖运行 S3File 连接器的前置条件无论作为 Source 还是 Sink使用 S3File 连接器都需要 Hadoop S3 相关依赖若使用 Spark/Flink需确保集群已集成 Hadoop官方测试过的版本为 2.x。若使用 SeaTunnel Zeta安装包已自动集成 hadoop jar可在${SEATUNNEL_HOME}/lib下确认。无论哪种引擎都需要将hadoop-aws-3.1.4.jar与aws-java-sdk-bundle-1.12.692.jar放入${SEATUNNEL_HOME}/lib目录详见 docs/en/connectors/source/S3File.md 的 Dependency 一节。S3File Source从 S3 读取数据的完整配置支持的文件格式与数据形态Source 端支持text、csv、parquet、orc、json、excel、xml、binary、markdown、pdf共十种格式不同格式对 schema 的要求不同json必须配置schema选项告知连接器如何解析为行。支持 JSON Lines每行一个 JSON 对象换行分隔。text、excel、csv、xml必须设置schema其中 text 格式还需field_delimitercsv/xml/excel 除外。parquet、orc不要求 schema连接器自动探测上游数据的 schema。官方文档给出的 JSON 读取示例上游数据形如{code: 200, data: get success, success: true}可多行对应的 schema 配置为schema { fields { code int data string success boolean } }text/csv 的典型配置则同时声明field_delimiter与schemafield_delimiter # schema { fields { name string age int gender string } }Source 端核心参数Source 的完整参数表见 docs/en/connectors/source/S3File.md 的 Options 一节以下是必填项与高频项参数必填默认值说明path是-需要读取的 S3 路径可含子路径file_format_type是-textcsvparquetorcjsonexcelxmlbinarymarkdownpdfbucket是-S3 bucket 地址如s3n://seatunnel-test使用s3a协议时为s3a://seatunnel-testfs.s3a.endpoint是-S3A endpointfs.s3a.aws.credentials.provider是com.amazonaws.auth.InstanceProfileCredentialsProviderS3A 凭据提供类全限定名access_key/secret_key条件-仅当使用SimpleAWSCredentialsProvider时必填read_columns否-列投影text/json/csv 需配合schemafield_delimiter否text 为\001csv 为,字段分隔符与 Hive 默认分隔符一致row_delimiter否\n行分隔符text 格式parse_partition_from_path否true是否从路径解析分区键值file_filter_pattern否-正则过滤文件filename_extension否-按扩展名过滤如csv.txtcompress_codec否nonetxt/json/csv 支持lzo、noneorc/parquet 自动识别archive_compress_codec否noneZIPTARTAR_GZGZ支持 txt/json/excel/xmlnull_format否-text 格式下定义哪些字符串代表 null如\Nschema否-上游数据 schemaenable_file_split/file_split_size否false/ 134217728128MB大文件逻辑分片以提升并行度仅支持 text/csv/json/parquet 且非压缩其中file_filter_pattern支持按文件名或按目录路径正则匹配。官方文档示例当path为/data/seatunnel时.*.txt匹配report.txtabc.*匹配以abc开头的文件/data/seatunnel/202410\d*/.*.csv匹配第三级目录以202410开头且扩展名为.csv的文件——模式以path开头则作用于完整路径否则仅作用于文件名。Source 完整示例使用SimpleAWSCredentialsProvider静态密钥认证、读取 orc 文件并输出到 Consoleenv { parallelism 1 job.mode BATCH } source { S3File { path /seatunnel/text fs.s3a.endpoints3.cn-north-1.amazonaws.com.cn fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key xxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxx bucket s3a://seatunnel-test file_format_type orc } } transform { } sink { Console {} }使用实例角色免 AK读取 json 并做列投影source { S3File { path /seatunnel/json bucket s3a://seatunnel-test fs.s3a.endpoints3.cn-north-1.amazonaws.com.cn fs.s3a.aws.credentials.providercom.amazonaws.auth.InstanceProfileCredentialsProvider file_format_type json read_columns [id, name] schema { fields { id int name string age int sex int type string } } } }路径分区解析parse_partition_from_path默认true控制是否从文件路径中解析分区键与分区值。例如读取s3n://hadoop-cluster/tmp/seatunnel/parquet/nametyrantlucifer/age26路径下的文件时每条记录会被自动附加两个字段nametyrantlucifer、age26。这正是 Hive 风格分区目录在 SeaTunnel 中的直接映射也是文件连接器通用的行为定义在文件基础模块的公共选项中。S3File Sink写出数据的完整配置写入流程与事务语义Sink 端数据先写入tmp_path默认/tmp/seatunnel下的临时目录再通过mv将临时目录提交到目标目录is_enable_transaction默认true开启时采用 2PC 提交保证数据不丢失、不重复。注意当is_enable_transactiontrue时文件名会自动加上${transactionId}_前缀。Sink 端核心参数Sink 的完整参数表见 docs/en/connectors/sink/S3File.md 的 Sink Options 一节参数必填默认值说明path是-目标路径支持变量替换如/test/${database_name}/${schema_name}/${table_name}bucket/fs.s3a.endpoint是-同 Sourcefs.s3a.aws.credentials.provider是InstanceProfileCredentialsProvider同 Sourcetmp_path否/tmp/seatunnel临时目录先写临时再 mv 提交需为 S3 目录file_format_type否csvtext/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_jsonfilename_extension否-覆盖默认扩展名如.xml.jsondatcustom_filename/file_name_expression/filename_time_format否false/${transactionId}/yyyy.MM.dd自定义文件名支持${now}、${uuid}变量field_delimiter/row_delimiter否\001(text) /,(csv)\n仅 text/csvrow_delimiter 亦支持 jsonhave_partition/partition_by/partition_dir_expression/is_partition_field_write_in_file否false/ - /${k0}${v0}/${k1}${v1}/.../false分区写入Hive 数据文件场景下is_partition_field_write_in_file应设为falsesink_columns否全部字段决定写出字段及顺序batch_size否1000000单文件最大行数Zeta 引擎下由batch_size与checkpoint.interval共同决定文件切分single_file_mode否false每个并行度只输出一个文件batch_size失效create_empty_file_when_no_data否false上游无数据时仍生成对应数据文件compress_codec否nonetext/json/csv:lzoorc:lzo snappy lz4 zlibparquet:lzo snappy lz4 gzip brotli zstdschema_save_mode否CREATE_SCHEMA_WHEN_NOT_EXIST对目标路径的预处理见下文 save modedata_save_mode否APPEND_DATA对目标路径已有数据文件的处理见下文 save modeenable_header_write否falsetext/csv 是否写表头encoding否UTF-8json/text/csv/xml 编码merge_update_event否falsecanal_json/debezium_json/maxwell_json 时合并 UPDATE_BEFORE/UPDATE_AFTERsave mode 机制S3File Sink 提供两套目标目录预处理策略对应 changelog 中 2.3.4 引入的 save mode 功能schema_save_modeRECREATE_SCHEMA路径不存在则创建存在则删除后重建。CREATE_SCHEMA_WHEN_NOT_EXIST不存在则创建存在则复用。ERROR_WHEN_SCHEMA_NOT_EXIST路径不存在直接报错。IGNORE忽略对目录的预处理。data_save_modeDROP_DATA复用路径但删除其中数据文件。APPEND_DATA复用路径并追加新文件。ERROR_WHEN_DATA_EXISTS路径中存在数据文件时报错。Sink 完整示例FakeSource 生成 16 行数据写入 S3 text 文件并启用分区、自定义文件名与事务env { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age tinyint } } } } transform { } sink { S3File { bucket s3a://seatunnel-test tmp_path /tmp/seatunnel path /seatunnel/text fs.s3a.endpoints3.cn-north-1.amazonaws.com.cn fs.s3a.aws.credentials.providercom.amazonaws.auth.InstanceProfileCredentialsProvider file_format_type text field_delimiter \t row_delimiter \n have_partition true partition_by [age] partition_dir_expression ${k0}${v0} is_partition_field_write_in_file true custom_filename true file_name_expression ${transactionId}_${now} filename_time_format yyyy.MM.dd sink_columns [name,age] is_enable_transaction true hadoop_s3_properties { fs.s3a.buffer.dir /data/st_test/s3a fs.s3a.fast.upload.buffer disk } } }多表写入与 CDCS3File Sink 支持从 CDC 上游如 MySQL-CDC多表写入path中使用${table_name}占位符即可为每张表生成独立目录这正是 changelog 中支持使用上游表占位符并自动替换2.3.6的能力sink { S3File { bucket s3a://seatunnel-test tmp_path /tmp/seatunnel/${table_name} path /test/${table_name} fs.s3a.endpoints3.cn-north-1.amazonaws.com.cn fs.s3a.aws.credentials.providerorg.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key xxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxx file_format_type orc schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }配合schema_evolution_enabledtrue默认falseSink 可在运行时处理 CDC 的 ADD/DROP/RENAME/MODIFY 列事件而无需重启作业但 binary 格式不支持该选项且have_partitiontrue时不允许删除partition_by中的列。若上游 CDC 开启了schema-changes.enabledtrue而 Sink 未开启 schema evolution作业会直接抛出可操作的错误提示。文件格式的写入差异写 text/csv 时所有列按字符串处理文件扩展名由file_format_type决定text 为txt。parquet 提供parquet_avro_write_timestamp_as_int96与parquet_avro_write_fixed_as_int96两个选项对应 changelog 2.3.6用于将 timestamp/12 字节定长字段以 INT96 写出增强与旧版 parquet 生态的兼容性。excel 格式支持max_rows_in_memory、sheet_max_rows默认 1048576、sheet_namecsv 支持csv_string_quote_modeALL/MINIMAL/NONExml 支持xml_root_tag默认RECORDS、xml_row_tag默认RECORD、xml_use_attr_format。认证与凭据五种 Provider 与源码级解析S3File 连接器的认证全部委托给 Hadoop S3A 的凭据提供机制配置键为fs.s3a.aws.credentials.provider。官方文档归纳了以下受支持的 ProviderProvider类名典型场景Simple AWSCredentialsorg.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider静态 access key / secret keyInstance Profilecom.amazonaws.auth.InstanceProfileCredentialsProviderEC2 实例角色默认Containercom.amazonaws.auth.ContainerCredentialsProviderECS 任务角色Default Chaincom.amazonaws.auth.DefaultAWSCredentialsProviderChain多来源回退链环境变量 → 系统属性 → profile → 容器 → 实例角色Custom任意com.amazonaws.auth.AWSCredentialsProvider实现用户自定义 Provider在源码层面S3HadoopConf.java 中的buildWithReadOnlyConfig完成如下装配逻辑根据bucket前缀判断协议以s3a开头则使用org.apache.hadoop.fs.s3a.S3AFileSystem否则默认走org.apache.hadoop.fs.s3native.NativeS3FileSystems3n。access_key/secret_key仅在显式配置时写入s3a协议写入fs.s3a.access.key/fs.s3a.secret.keys3n协议写入fs.s3n.awsAccessKeyId/fs.s3n.awsSecretAccessKey——这正是 changelog 2.3.10 中修复某些场景下 s3 key 设置错误所针对的行为。hadoop_s3_properties中的键值对直接透传进 S3A 配置。fs.s3a.aws.credentials.provider在配置解析期即被校验checkCredentialsProviders对逗号/换行分隔的 Provider 链逐项校验类必须实现com.amazonaws.auth.AWSCredentialsProvider且非抽象类若类在当前节点不可解析则降级为警告并延迟到 Worker 运行期校验Provider jar 只需存在于实际运行 S3A 的节点如${SEATUNNEL_HOME}/lib。这一行为被 S3HadoopConfTest.java 中的十余个测试用例系统验证包括链式 Provider 透传、空段失败、抽象类失败、隔离类加载器不可绕过校验等。容器环境Kubernetes / ECS / EKS实践Kubernetes/EKS 推荐直接使用 EC2 节点实例角色保持默认InstanceProfileCredentialsProvider即可免密钥。实例角色不可用时可从 Kubernetes Secret 注入access_key/secret_key并切换为SimpleAWSCredentialsProvider。ECS 任务角色依赖AWS_CONTAINER_CREDENTIALS_RELATIVE_URI环境变量使用ContainerCredentialsProvider。EKS IRSA 需要的WebIdentityTokenCredentialsProvider不在 SeaTunnel 捆绑的旧版 AWS SDK v1.x1.11.271中官方建议改用节点实例角色、Secret 注入或向所有节点补充新版 AWS SDK jar。跨账号 / STS AssumeRole对于跨账号读取或写入可通过hadoop_s3_properties配合TemporaryAWSCredentialsProvider传递sts:AssumeRole签发的临时会话凭据source { S3File { path /cross-account/prefix bucket s3a://target-bucket fs.s3a.endpoint s3.cn-north-1.amazonaws.com.cn fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.TemporaryAWSCredentialsProvider hadoop_s3_properties { fs.s3a.access.key assumed-role-access-key fs.s3a.secret.key assumed-role-secret-key fs.s3a.session.token assumed-role-session-token } file_format_type parquet } }注意连接器总是用选项值覆盖fs.s3a.aws.credentials.provider键因此无法通过hadoop_s3_properties二次覆盖该键。连续发现将 S3 变成准实时数据源discovery_modecontinuous默认once可使 Source 周期扫描 S3 路径把新出现或变更的对象当作无界数据流持续处理典型场景是边落盘边消费。当前该模式有两个硬性前提file_format_typebinary且sync_modeupdate并需将target_path指向 Sink 相同的基准路径以便跳过未变更对象。官方示例env { parallelism 1 job.mode STREAMING } source { S3File { path /watch/source bucket s3a://seatunnel-test fs.s3a.endpoint s3.amazonaws.com fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key xxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxx file_format_type binary discovery_mode continuous scan_interval 10S start_mode earliest sync_mode update target_path /watch/target } } sink { S3File { path /watch/target tmp_path /watch/tmp bucket s3a://seatunnel-test fs.s3a.endpoint s3.amazonaws.com fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key xxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxx file_format_type binary } }配套参数包括scan_interval轮询间隔默认10S支持 ISO-8601 如PT10S、start_modeearliest处理存量文件 /latest仅处理启动后新增、update_strategydistcp/strict、compare_modelen_mtime/checksumchecksum 仅 strict 策略有效、update_compare_parallelism默认 8范围 1-64、post_sync_actionnone/delete/backup与备份保留策略。该能力复用文件基础模块的比对逻辑不依赖 S3 事件通知也不产生对象删除事件或 changelog 行。常见问题排查官方文档针对认证类故障给出了明确的排查路径Factory initialize failed或类似类加载错误凭据 Provider 类不在 classpath。确保 Provider jar 存在于每个集群节点的${SEATUNNEL_HOME}/lib而非仅提交节点。No AWS Credentials provided by ...按 Provider 类型检查——SimpleAWSCredentialsProvider检查access_key/secret_keyInstanceProfileCredentialsProvider检查 EC2 实例是否绑定 IAM 角色ContainerCredentialsProvider检查AWS_CONTAINER_CREDENTIALS_RELATIVE_URI是否设置。配置解析期出现IllegalArgumentException类名拼写错误或类未实现com.amazonaws.auth.AWSCredentialsProvider接口核对全限定类名即可。此外读取 XML 时若文件含!DOCTYPE ...声明包括仅定义内部实体的良性声明会因 XXE 加固被拒绝并报FILE_READ_FAILED错误且无配置可恢复旧行为需要在上游预处理时移除 DOCTYPE 头。延伸阅读与源码导航Source 完整文档docs/en/connectors/source/S3File.mdSink 完整文档docs/en/connectors/sink/S3File.md变更日志原文docs/en/connectors/changelog/connector-file-s3.md文件类对象存储 FAQdocs/en/connectors/file-object-storage-faq.md源码入口连接器实现位于 seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3其中config包定义选项S3FileBaseOptions.java、S3HadoopConf.javasource/sink包分别承载读写实现catalog包提供 S3FileCatalog单元测试集中在 connector-file-s3/src/test 下可作为行为契约参考。一个完整的 JDBC→S3 实战配方可参考 docs/en/getting-started/recipes/jdbc-to-s3.md。综上S3File 连接器的能力早已超出读写 S3本身它通过file_format_type统一了十种文本与结构化格式通过hadoop_s3_properties打通了任意fs.s3a.*扩展配置通过凭据 Provider 链兼容云上四种主流认证方式再叠加分区写入、单文件模式、schema evolution、连续发现等工程能力构成了 SeaTunnel 面向对象存储数据集成场景的中坚力量。读者在落地时建议先对照 changelog 时间线确认所用 SeaTunnel 版本是否包含目标能力再按本文给出的参数表与示例完成作业配置。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考