SeaTunnel S3File 连接器完全指南:从版本演进史到源码级配置实践
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
S3File 是 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 管道中的数据写出到 S3(Sink)。从官方文档与源码工厂类可以看出其能力定位:
- 支持引擎: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.1 | file_type更名为file_format_type;支持压缩(compress);新增S3Catalog;重构 schema 解析;升级 guava 至 27.0-jre | 格式语义统一、压缩读、元数据目录能力上线 |
| 2.3.2 | 新增 excel sink 与 source;删除不可用的 S3 & Kafka Catalogs | Excel 格式读写 |
| 2.3.3 | 新增file_filter_pattern用于按正则过滤文件 | 精细化文件选择 |
| 2.3.4 | 支持 LZO 压缩读取;支持读取空目录;支持配置 schema 中的 column/primaryKey/constraintKey;ENABLE_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 写为 int96;sink 选项中支持上游表占位符并自动替换;修复 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_delimiter(csv/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与schema:
field_delimiter = "#" schema { fields { name = string age = int gender = string } }Source 端核心参数
Source 的完整参数表见 docs/en/connectors/source/S3File.md 的 Options 一节,以下是必填项与高频项:
| 参数 | 必填 | 默认值 | 说明 |
|---|---|---|---|
path | 是 | - | 需要读取的 S3 路径,可含子路径 |
file_format_type | 是 | - | textcsvparquetorcjsonexcelxmlbinarymarkdownpdf |
bucket | 是 | - | S3 bucket 地址,如s3n://seatunnel-test;使用s3a协议时为s3a://seatunnel-test |
fs.s3a.endpoint | 是 | - | S3A endpoint |
fs.s3a.aws.credentials.provider | 是 | com.amazonaws.auth.InstanceProfileCredentialsProvider | S3A 凭据提供类全限定名 |
access_key/secret_key | 条件 | - | 仅当使用SimpleAWSCredentialsProvider时必填 |
read_columns | 否 | - | 列投影;text/json/csv 需配合schema |
field_delimiter | 否 | text 为\001,csv 为, | 字段分隔符,与 Hive 默认分隔符一致 |
row_delimiter | 否 | \n | 行分隔符(text 格式) |
parse_partition_from_path | 否 | true | 是否从路径解析分区键值 |
file_filter_pattern | 否 | - | 正则过滤文件 |
filename_extension | 否 | - | 按扩展名过滤,如csv.txt |
compress_codec | 否 | none | txt/json/csv 支持lzo、none;orc/parquet 自动识别 |
archive_compress_codec | 否 | none | ZIPTARTAR_GZGZ,支持 txt/json/excel/xml |
null_format | 否 | - | text 格式下定义哪些字符串代表 null,如\N |
schema | 否 | - | 上游数据 schema |
enable_file_split/file_split_size | 否 | false/ 134217728(128MB) | 大文件逻辑分片以提升并行度,仅支持 text/csv/json/parquet 且非压缩 |
其中file_filter_pattern支持按文件名或按目录路径正则匹配。官方文档示例:当path为/data/seatunnel时,.*.txt匹配report.txt,abc.*匹配以abc开头的文件,/data/seatunnel/202410\d*/.*.csv匹配第三级目录以202410开头且扩展名为.csv的文件——模式以path开头则作用于完整路径,否则仅作用于文件名。
Source 完整示例
使用SimpleAWSCredentialsProvider静态密钥认证、读取 orc 文件并输出到 Console:
env { parallelism = 1 job.mode = "BATCH" } source { S3File { path = "/seatunnel/text" fs.s3a.endpoint="s3.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.endpoint="s3.cn-north-1.amazonaws.com.cn" fs.s3a.aws.credentials.provider="com.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/name=tyrantlucifer/age=26路径下的文件时,每条记录会被自动附加两个字段:name="tyrantlucifer"、age=26。这正是 Hive 风格分区目录在 SeaTunnel 中的直接映射,也是文件连接器通用的行为(定义在文件基础模块的公共选项中)。
S3File Sink:写出数据的完整配置
写入流程与事务语义
Sink 端数据先写入tmp_path(默认/tmp/seatunnel)下的临时目录,再通过mv将临时目录提交到目标目录;is_enable_transaction(默认true)开启时采用 2PC 提交,保证数据不丢失、不重复。注意:当is_enable_transaction=true时,文件名会自动加上${transactionId}_前缀。
Sink 端核心参数
Sink 的完整参数表见 docs/en/connectors/sink/S3File.md 的 Sink Options 一节:
| 参数 | 必填 | 默认值 | 说明 |
|---|---|---|---|
path | 是 | - | 目标路径,支持变量替换如/test/${database_name}/${schema_name}/${table_name} |
bucket/fs.s3a.endpoint | 是 | - | 同 Source |
fs.s3a.aws.credentials.provider | 是 | InstanceProfileCredentialsProvider | 同 Source |
tmp_path | 否 | /tmp/seatunnel | 临时目录,先写临时再 mv 提交,需为 S3 目录 |
file_format_type | 否 | csv | text/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_json |
filename_extension | 否 | - | 覆盖默认扩展名,如.xml.jsondat |
custom_filename/file_name_expression/filename_time_format | 否 | false/${transactionId}/yyyy.MM.dd | 自定义文件名,支持${now}、${uuid}变量 |
field_delimiter/row_delimiter | 否 | \001(text) /,(csv);\n | 仅 text/csv(row_delimiter 亦支持 json) |
have_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应设为false |
sink_columns | 否 | 全部字段 | 决定写出字段及顺序 |
batch_size | 否 | 1000000 | 单文件最大行数;Zeta 引擎下由batch_size与checkpoint.interval共同决定文件切分 |
single_file_mode | 否 | false | 每个并行度只输出一个文件,batch_size失效 |
create_empty_file_when_no_data | 否 | false | 上游无数据时仍生成对应数据文件 |
compress_codec | 否 | none | text/json/csv:lzo;orc:lzo snappy lz4 zlib;parquet:lzo snappy lz4 gzip brotli zstd |
schema_save_mode | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 对目标路径的预处理(见下文 save mode) |
data_save_mode | 否 | APPEND_DATA | 对目标路径已有数据文件的处理(见下文 save mode) |
enable_header_write | 否 | false | text/csv 是否写表头 |
encoding | 否 | UTF-8 | json/text/csv/xml 编码 |
merge_update_event | 否 | false | canal_json/debezium_json/maxwell_json 时合并 UPDATE_BEFORE/UPDATE_AFTER |
save mode 机制
S3File Sink 提供两套目标目录预处理策略(对应 changelog 中 2.3.4 引入的 save mode 功能):
schema_save_mode:RECREATE_SCHEMA:路径不存在则创建,存在则删除后重建。CREATE_SCHEMA_WHEN_NOT_EXIST:不存在则创建,存在则复用。ERROR_WHEN_SCHEMA_NOT_EXIST:路径不存在直接报错。IGNORE:忽略对目录的预处理。
data_save_mode:DROP_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.endpoint="s3.cn-north-1.amazonaws.com.cn" fs.s3a.aws.credentials.provider="com.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" } } }多表写入与 CDC
S3File 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.endpoint="s3.cn-north-1.amazonaws.com.cn" fs.s3a.aws.credentials.provider="org.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_enabled=true(默认false),Sink 可在运行时处理 CDC 的 ADD/DROP/RENAME/MODIFY 列事件而无需重启作业;但 binary 格式不支持该选项,且have_partition=true时不允许删除partition_by中的列。若上游 CDC 开启了schema-changes.enabled=true而 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_name;csv 支持csv_string_quote_mode(ALL/MINIMAL/NONE);xml 支持xml_root_tag(默认RECORDS)、xml_row_tag(默认RECORD)、xml_use_attr_format。
认证与凭据:五种 Provider 与源码级解析
S3File 连接器的认证全部委托给 Hadoop S3A 的凭据提供机制,配置键为fs.s3a.aws.credentials.provider。官方文档归纳了以下受支持的 Provider:
| Provider | 类名 | 典型场景 |
|---|---|---|
| Simple AWSCredentials | org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider | 静态 access key / secret key |
| Instance Profile | com.amazonaws.auth.InstanceProfileCredentialsProvider | EC2 实例角色(默认) |
| Container | com.amazonaws.auth.ContainerCredentialsProvider | ECS 任务角色 |
| Default Chain | com.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.NativeS3FileSystem(s3n)。 access_key/secret_key仅在显式配置时写入:s3a协议写入fs.s3a.access.key/fs.s3a.secret.key,s3n协议写入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.x(1.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_mode=continuous(默认once)可使 Source 周期扫描 S3 路径,把新出现或变更的对象当作无界数据流持续处理,典型场景是"边落盘边消费"。当前该模式有两个硬性前提:file_format_type="binary"且sync_mode="update",并需将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_mode(earliest处理存量文件 /latest仅处理启动后新增)、update_strategy(distcp/strict)、compare_mode(len_mtime/checksum,checksum 仅 strict 策略有效)、update_compare_parallelism(默认 8,范围 1-64)、post_sync_action(none/delete/backup)与备份保留策略。该能力复用文件基础模块的比对逻辑,不依赖 S3 事件通知,也不产生对象删除事件或 changelog 行。
常见问题排查
官方文档针对认证类故障给出了明确的排查路径:
Factory initialize failed或类似类加载错误:凭据 Provider 类不在 classpath。确保 Provider jar 存在于每个集群节点的${SEATUNNEL_HOME}/lib,而非仅提交节点。No AWS Credentials provided by ...:按 Provider 类型检查——SimpleAWSCredentialsProvider检查access_key/secret_key;InstanceProfileCredentialsProvider检查 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.md
- Sink 完整文档:docs/en/connectors/sink/S3File.md
- 变更日志原文:docs/en/connectors/changelog/connector-file-s3.md
- 文件类对象存储 FAQ:docs/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.java),source/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),仅供参考