跳到主要内容
版本:3.0.0

S3File

S3 文件 Sink 连接器

支持的引擎​

Spark
Flink
SeaTunnel Zeta

主要特性​

  • 多模态

    使用二进制文件格式读取和写入任何格式的文件,例如视频、图片等。简而言之,任何文件都可以同步到目标位置。

  • 精确一次

    默认情况下,我们使用 2PC 提交来确保 精确一次。

  • CDC

  • 支持多表写入

  • 文件格式类型

    • text
    • csv
    • parquet
    • orc
    • json
    • excel
    • xml
    • binary
    • canal_json
    • debezium_json
    • maxwell_json
  • 定时刷新

描述​

将数据输出到 AWS S3 文件系统。

支持的数据源信息​

数据源支持的版本
S3当前版本

数据库依赖​

如果您使用 Spark/Flink,为了使用此连接器,您必须确保您的 Spark/Flink 集群已经集成了 Hadoop。测试的 Hadoop 版本为 2.x。

如果您使用 SeaTunnel引擎,当您下载并安装 SeaTunnel引擎时,它会自动集成 Hadoop jar 包。您可以在 ${SEATUNNEL_HOME}/lib 下检查 jar 包以确认这一点。 要使用此连接器,您需要将 hadoop-aws-3.1.4.jar 和 aws-java-sdk-bundle-1.12.692.jar 放在 ${SEATUNNEL_HOME}/lib 目录下。

数据类型映射​

如果写入 csv、text 文件类型,所有列都将为字符串类型。

Orc 文件类型​

SeaTunnel 数据类型Orc 数据类型
STRINGSTRING
BOOLEANBOOLEAN
TINYINTBYTE
SMALLINTSHORT
INTINT
BIGINTLONG
FLOATFLOAT
FLOATFLOAT
DOUBLEDOUBLE
DECIMALDECIMAL
BYTESBINARY
DATEDATE
TIME
TIMESTAMP
TIMESTAMP
ROWSTRUCT
NULL不支持的数据类型
ARRAYLIST
MapMap

Parquet 文件类型​

SeaTunnel 数据类型Parquet 数据类型
STRINGSTRING
BOOLEANBOOLEAN
TINYINTINT_8
SMALLINTINT_16
INTINT32
BIGINTINT64
FLOATFLOAT
FLOATFLOAT
DOUBLEDOUBLE
DECIMALDECIMAL
BYTESBINARY
DATEDATE
TIME
TIMESTAMP
TIMESTAMP_MILLIS
ROWGroupType
NULL不支持的数据类型
ARRAYLIST
MapMap

Sink 选项​

名称类型是否必填默认值描述
pathstring是-
tmp_pathstring否/tmp/seatunnel结果文件将首先写入临时路径,然后使用 mv 将临时目录提交到目标目录。需要一个 S3 目录。
bucketstring是-
fs.s3a.endpointstring是-
fs.s3a.aws.credentials.providerstring是com.amazonaws.auth.InstanceProfileCredentialsProvider透传给 Hadoop 的 S3A 凭据提供程序的全限定类名。除了两个常用值 org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider(使用静态 access_key/secret_key)和 com.amazonaws.auth.InstanceProfileCredentialsProvider(默认值)之外,任何 classpath 上可用的 S3A 凭据提供程序类都可以使用,例如基于容器的 com.amazonaws.auth.ContainerCredentialsProvider 或自定义提供程序。该类必须实现 com.amazonaws.auth.AWSCredentialsProvider 接口,并提供 Hadoop 3.1.4 支持的任一创建方式:公共 (java.net.URI, org.apache.hadoop.conf.Configuration) 构造器、公共 (org.apache.hadoop.conf.Configuration) 构造器、返回 AWSCredentialsProvider 的公共静态无参 getInstance() 工厂方法,或公共无参构造器。支持 Hadoop 风格的逗号或换行分隔的提供程序链,且每个类都会被独立校验。该提供程序 jar 必须存在于每个集群节点的运行时 classpath 中(例如放在 ${SEATUNNEL_HOME}/lib 下),而不仅仅是提交作业的节点。共享/多租户集群的运维者请注意:此选项允许作业编写者按类名加载类,因此请相应地限制作业提交权限。
access_keystring否-仅当 fs.s3a.aws.credentials.provider = org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider 时使用
secret_keystring否-仅当 fs.s3a.aws.credentials.provider = org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider 时使用
custom_filenameboolean否false是否需要自定义文件名
file_name_expressionstring否"${transactionId}"仅当 custom_filename 为 true 时使用
filename_time_formatstring否"yyyy.MM.dd"仅当 custom_filename 为 true 时使用
file_format_typestring否"csv"
field_delimiterstring否'\001'仅当 file_format 为 text 时使用
row_delimiterstring否"\n"仅当 file_format 为 text、csv、json 时使用
have_partitionboolean否false是否需要处理分区。
partition_byarray否-仅当 have_partition 为 true 时使用
partition_dir_expressionstring否"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/"仅当 have_partition 为 true 时使用
is_partition_field_write_in_fileboolean否false仅当 have_partition 为 true 时使用
sink_columnsarray否当此参数为空时,所有字段均为 sink 列
is_enable_transactionboolean否true
batch_sizeint否1000000
compress_codecstring否none
common-optionsobject否-
max_rows_in_memoryint否-仅当 file_format 为 excel 时使用
sheet_namestring否Sheet${Random number}仅当 file_format 为 excel 时使用
csv_string_quote_modeenum否MINIMAL仅当 file_format 为 csv 时使用
xml_root_tagstring否RECORDS仅当 file_format 为 xml 时使用,指定 XML 文件中根元素的标签名称。
xml_row_tagstring否RECORD仅当 file_format 为 xml 时使用,指定 XML 文件中数据行的标签名称。
xml_use_attr_formatboolean否-仅当 file_format 为 xml 时使用,指定是否使用标签属性格式处理数据。
single_file_modeboolean否false每个并行度只会输出一个文件。当此参数开启时,batch_size 将不会生效。输出文件名不会有文件块后缀。
create_empty_file_when_no_databoolean否false当上游没有数据同步时,仍然会生成相应的数据文件。
parquet_avro_write_timestamp_as_int96boolean否false仅当 file_format 为 parquet 时使用
parquet_avro_write_fixed_as_int96array否-仅当 file_format 为 parquet 时使用
hadoop_s3_propertiesmap否如果您需要添加其他选项,可以在此处添加,并参考此链接
schema_evolution_enabledboolean否false开启 Schema 演变支持,适用于 CDC 管道。为 true 时,来自上游的 ADD/DROP/RENAME/MODIFY 列事件无需重启作业即可应用到 Sink。不支持 binary 格式。
schema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST在开启同步任务之前,对目标路径进行不同的处理
data_save_modeEnum否APPEND_DATA在开启同步任务之前,对目标路径中的数据文件进行不同的处理
enable_header_writeboolean否false仅当 file_format_type 为 text,csv 时使用。
false: 不写入表头, true: 写入表头。
encodingstring否"UTF-8"仅当 file_format_type 为 json,text,csv,xml 时使用。
merge_update_eventboolean否false仅当file_format_type为canal_json、debezium_json、maxwell_json.

path [string]​

存储数据文件的路径,支持变量替换。例如:path=/test/${database_name}/${schema_name}/${table_name}

hadoop_s3_properties [map]​

如果您需要添加其他选项,可以在此处添加,并参考此链接

hadoop_s3_properties {
"fs.s3a.buffer.dir" = "/data/st_test/s3a"
"fs.s3a.fast.upload.buffer" = "disk"
}

custom_filename [boolean]​

是否自定义文件名

file_name_expression [string]​

仅当 custom_filename 为 true 时使用

file_name_expression 描述了将创建到 path 中的文件表达式。我们可以在 file_name_expression 中添加变量 ${now} 或 ${uuid},例如 test_${uuid}_${now}, ${now} 表示当前时间,其格式可以通过指定选项 filename_time_format 来定义。

请注意,如果 is_enable_transaction 为 true,我们会在文件头部自动添加 ${transactionId}_。

filename_time_format [string]​

仅当 custom_filename 为 true 时使用

当 file_name_expression 参数中的格式为 xxxx-${now} 时,filename_time_format 可以指定路径的时间格式,默认值为 yyyy.MM.dd。常用的时间格式如下:

符号描述
y年
M月
d日
H小时 (0-23)
m分钟
s秒

file_format_type [string]​

我们支持以下文件类型:

text csv parquet orc json excel xml binary canal_json debezium_json maxwell_json

请注意,最终文件名将以文件格式类型的后缀结尾,文本文件的后缀为 txt。

field_delimiter [string]​

行数据中列之间的分隔符。仅在 text 文件格式中需要。

row_delimiter [string]​

文件中行之间的分隔符。仅在 text、csv、json 文件格式中需要。

have_partition [boolean]​

是否需要处理分区。

partition_by [array]​

仅当 have_partition 为 true 时使用。

根据选定的字段对数据进行分区。

partition_dir_expression [string]​

仅当 have_partition 为 true 时使用。

如果指定了 partition_by,我们将根据分区信息生成相应的分区目录,最终文件将放置在分区目录中。

默认的 partition_dir_expression 为 ${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/。k0 是第一个分区字段,v0 是第一个分区字段的值。

is_partition_field_write_in_file [boolean]​

仅当 have_partition 为 true 时使用。

如果 is_partition_field_write_in_file 为 true,分区字段及其值将被写入数据文件。

例如,如果您想写入 Hive 数据文件,其值应为 false。

sink_columns [array]​

哪些列需要写入文件,默认值为从 Transform 或 Source 获取的所有列。 字段的顺序决定了文件实际写入的顺序。

is_enable_transaction [boolean]​

如果 is_enable_transaction 为 true,我们将确保在将数据写入目标目录时不会丢失或重复。

请注意,如果 is_enable_transaction 为 true,我们会在文件头部自动添加 ${transactionId}_。

目前仅支持 true。

batch_size [int]​

文件中的最大行数。对于 SeaTunnel Engine,文件中的行数由 batch_size 和 checkpoint.interval 共同决定。如果 checkpoint.interval 的值足够大,sink writer 将一直写入文件,直到文件中的行数超过 batch_size。如果 checkpoint.interval 较小,sink writer 将在新的 checkpoint 触发时创建一个新文件。

compress_codec [string]​

文件的压缩编解码器,支持的详细信息如下:

  • txt: lzo none
  • json: lzo none
  • csv: lzo none
  • orc: lzo snappy lz4 zlib none
  • parquet: lzo snappy lz4 gzip brotli zstd none

提示:excel 类型不支持任何压缩格式

common options​

Sink 插件通用参数,请参考 Sink 通用选项 获取详细信息。

max_rows_in_memory [int]​

当文件格式为 Excel 时,内存中可以缓存的最大数据项数。

sheet_name [string]​

写入工作表的名称

csv_string_quote_mode [string]​

当文件格式为 CSV 时,CSV 的字符串引用模式。

  • ALL: 所有字符串字段都会被引用。
  • MINIMAL: 引用包含特殊字符的字段,如字段分隔符、引用字符或行分隔符字符串中的任何字符。
  • NONE: 从不引用字段。当数据中出现分隔符时,打印机会在其前面加上转义字符。如果未设置转义字符,格式验证将抛出异常。

xml_root_tag [string]​

指定 XML 文件中根元素的标签名称。

xml_row_tag [string]​

指定 XML 文件中数据行的标签名称。

xml_use_attr_format [boolean]​

指定是否使用标签属性格式处理数据。

parquet_avro_write_timestamp_as_int96 [boolean]​

支持将时间戳写入 Parquet INT96,仅对 parquet 文件有效。

parquet_avro_write_fixed_as_int96 [array]​

支持将 12-byte 字段写入 Parquet INT96,仅对 parquet 文件有效。

schema_save_mode [Enum]​

在开启同步任务之前,对目标路径进行不同的处理。
选项介绍:
RECREATE_SCHEMA :当路径不存在时创建。如果路径已存在,则删除路径并重新创建。
CREATE_SCHEMA_WHEN_NOT_EXIST :当路径不存在时创建,路径存在时使用路径。
ERROR_WHEN_SCHEMA_NOT_EXIST :当路径不存在时报错
IGNORE :忽略表的处理

data_save_mode [Enum]​

在开启同步任务之前,对目标路径中的数据文件进行不同的处理。 选项介绍:
DROP_DATA:使用路径但删除路径中的数据文件。 APPEND_DATA:使用路径,并在路径中添加新文件以写入数据。
ERROR_WHEN_DATA_EXISTS:当路径中存在数据文件时,将报错。

encoding [string]​

仅当 file_format_type 为 json,text,csv,xml 时使用。 写入文件的编码。此参数将由 Charset.forName(encoding) 解析。

merge_update_event [boolean]​

仅当file_format_type为canal_json、debezium_json、maxwell_json时使用. 设置成true,序列化数据时,UPDATE_AFTER 和 UPDATE_BEFORE 会合并成 UPDATE; 设置成false,序列化数据时,UPDATE_AFTER 和 UPDATE_BEFORE 不会合并;

示例​

简单示例​

此示例定义了一个 SeaTunnel 同步任务,通过 FakeSource 自动生成数据并将其发送到 S3File Sink。FakeSource 总共生成 16 行数据 (row.num=16),每行有两个字段,name (字符串类型) 和 age (int 类型)。最终的目标 s3 目录将创建一个文件,并将所有数据写入其中。 在运行此作业之前,您需要创建 s3 路径:/seatunnel/text。如果您尚未安装和部署 SeaTunnel,您需要按照 安装 SeaTunnel 中的说明安装和部署 SeaTunnel。然后按照 使用 SeaTunnel Engine 快速入门 中的说明运行此作业。

# 定义运行时环境
env {
parallelism = 1
job.mode = "BATCH"
}

source {
# 这是一个示例源插件,仅用于测试和演示功能源插件
FakeSource {
parallelism = 1
plugin_output = "fake"
row.num = 16
schema = {
fields {
c_map = "map<string, array<int>>"
c_array = "array<int>"
name = string
c_boolean = boolean
age = tinyint
c_smallint = smallint
c_int = int
c_bigint = bigint
c_float = float
c_double = double
c_decimal = "decimal(16, 1)"
c_null = "null"
c_bytes = bytes
c_date = date
c_timestamp = timestamp
}
}
}
# 如果您想了解更多关于如何配置SeaTunnel以及查看完整的源插件列表,
# 请访问 https://seatunnel.apache.org/docs/connectors/source
source {
}

transform {
# 如果您想了解更多关于如何配置SeaTunnel以及查看完整的转换插件列表,
# 请访问 https://seatunnel.apache.org/docs/transforms
}

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"
}
}
# 如果您想了解更多关于如何配置SeaTunnel以及查看完整的接收插件列表,
# 请访问 https://seatunnel.apache.org/docs/connectors/sink
}

对于文本文件格式,包含 have_partition、custom_filename、sink_columns 和 com.amazonaws.auth.InstanceProfileCredentialsProvider

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"
}
}

对于Parquet文件格式,简单配置使用 org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider

S3File {
bucket = "s3a://seatunnel-test"
tmp_path = "/tmp/seatunnel"
path="/seatunnel/parquet"
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 = "parquet"
hadoop_s3_properties {
"fs.s3a.buffer.dir" = "/data/st_test/s3a"
"fs.s3a.fast.upload.buffer" = "disk"
}
}

对于ORC文件格式,简单配置使用 org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider

S3File {
bucket = "s3a://seatunnel-test"
tmp_path = "/tmp/seatunnel"
path="/seatunnel/orc"
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"
}

多表写入和保存模式

env {
"job.name"="SeaTunnel_job"
"job.mode"=STREAMING
}
source {
MySQL-CDC {
database-names=[
"wls_t1"
]
table-names=[
"wls_t1.mysqlcdc_to_s3_t3",
"wls_t1.mysqlcdc_to_s3_t4",
"wls_t1.mysqlcdc_to_s3_t5",
"wls_t1.mysqlcdc_to_s3_t1",
"wls_t1.mysqlcdc_to_s3_t2"
]
password="xxxxxx"
username="xxxxxxxxxxxxx"
url="jdbc:mysql://localhost:3306/qa_source"
}
}

transform {
}

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"
}
}

enable_header_write [boolean]​

仅在 file_format_type 为 text 或 csv 时使用。false:不写入表头,true:写入表头。

schema_evolution_enabled [boolean]​

设置为 true 时,文件 Sink 可在运行时处理 CDC Schema 变更事件(ADD COLUMN、DROP COLUMN、RENAME COLUMN、MODIFY COLUMN 类型),无需重启作业。每次 Schema 变更时,当前输出文件会被关闭,并以新 Schema 打开一个新文件。

支持的格式: 除 binary 外的所有文件格式。将此选项与 file_format_type = binary 一起使用时,作业启动时会抛出配置校验错误。

分区约束: 当 have_partition = true 时,不允许删除 partition_by 中列出的列,违反时会立即抛出异常。分区列在 Schema 变更过程中必须保持稳定。

当 schema_evolution_enabled = false(默认值)时: 若上游 CDC Source 配置了 schema-changes.enabled = true 且 Sink 收到 AlterTableEvent,作业会立即抛出如下错误:

Received AlterTableEvent but schema_evolution_enabled=false at this sink. Either set schema_evolution_enabled=true to handle schema changes, or set schema-changes.enabled=false at the CDC source to suppress them.

使用默认 CDC Source 配置(schema-changes.enabled = false)的用户不受影响。

已知限制: Schema 变更与 Checkpoint 不是原子操作。若作业在文件轮转与 Schema 元数据更新之间的窗口期崩溃,恢复后写入的数据行可能使用变更前的 Schema。这是与其他 SeaTunnel Sink 共同存在的已知架构限制。完整的重启后 DDL 正确性支持需要配套的 CDC Source 修复(另行跟踪)。

CDC 管道中的使用示例:

S3File {
path = "/test/cdc/${table_name}"
fs.s3a.endpoint = "s3.cn-north-1.amazonaws.com.cn"
access_key = "xxxxxxxxxxxxxxxxx"
secret_key = "xxxxxxxxxxxxxxxxx"
file_format_type = "parquet"
schema_evolution_enabled = true
}

生产作业中不建议把长期有效的密钥直接写入任务文件。优先使用 IAM 类认证方式,例如 fs.s3a.aws.credentials.provider = com.amazonaws.auth.InstanceProfileCredentialsProvider,或通过 SeaTunnel 变量替换注入 access_key 和 secret_key。

使用 STS AssumeRole 写入(跨账号写入)​

向另一个 AWS 账号拥有的 bucket 写入时,先通过 sts:AssumeRole 拿到临时会话凭证,再通过 hadoop_s3_properties 与 TemporaryAWSCredentialsProvider 配合使用。

sink {
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"
schema_evolution_enabled = true
}
}

对于 AWS SSO / Profile 角色,把 provider 类换成 com.amazonaws.auth.profile.ProfileCredentialsProvider,并把 fs.s3a.profile、fs.s3a.credentialsFile 等 provider 特定键放在 hadoop_s3_properties 里。完整的 fs.s3a.* 键集合参见 Hadoop AWS 文档。

容器环境中的凭据提供程序​

在容器环境(Kubernetes、ECS、EKS、Docker)中运行 SeaTunnel 时,S3File 连接器接受任何实现 com.amazonaws.auth.AWSCredentialsProvider 接口且在 classpath 上可用的全限定 S3A 凭据提供程序类。fs.s3a.aws.credentials.provider 选项在配置解析时进行验证(当类在构建配置的节点上可解析时):类必须实现 AWS 凭据提供程序接口,且不能是抽象类。当类无法解析时(例如,提供程序 JAR 仅在 worker 节点上可用),验证将延迟到实际运行 S3A 的 worker 节点上的运行时进行。

支持的凭据提供程序​

提供程序类名典型场景
Simple AWSCredentialsorg.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider静态 access key / secret key
Instance Profilecom.amazonaws.auth.InstanceProfileCredentialsProviderEC2 实例角色(默认)
Containercom.amazonaws.auth.ContainerCredentialsProviderECS 任务角色
Default Chaincom.amazonaws.auth.DefaultAWSCredentialsProviderChain多源回退链
自定义任何 com.amazonaws.auth.AWSCredentialsProvider 实现用户自定义提供程序

Kubernetes / EKS 配置​

EC2 节点实例角色(推荐):如果您的 EKS 工作节点具有包含 S3 权限的 EC2 实例配置文件,默认的 InstanceProfileCredentialsProvider 会自动从实例元数据服务解析凭据:

S3File {
bucket = "s3a://my-bucket"
tmp_path = "/tmp/seatunnel"
fs.s3a.endpoint = "s3.amazonaws.com"
path = "/data/output"
file_format_type = "parquet"
}

通过 Kubernetes Secret 注入静态密钥(备选方案):如果实例角色不可用,从 Kubernetes Secret 注入凭据:

S3File {
bucket = "s3a://my-bucket"
tmp_path = "/tmp/seatunnel"
fs.s3a.endpoint = "s3.amazonaws.com"
fs.s3a.aws.credentials.provider = "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
access_key = "<from-k8s-secret>"
secret_key = "<from-k8s-secret>"
path = "/data/output"
file_format_type = "parquet"
}

DefaultAWSCredentialsProviderChain:对于需要灵活部署的场景,默认链按顺序尝试多个凭据来源(环境变量 → 系统属性 → profile → 容器 → 实例配置文件):

S3File {
bucket = "s3a://my-bucket"
tmp_path = "/tmp/seatunnel"
fs.s3a.endpoint = "s3.amazonaws.com"
fs.s3a.aws.credentials.provider = "com.amazonaws.auth.DefaultAWSCredentialsProviderChain"
path = "/data/output"
file_format_type = "parquet"
}

ECS 任务角色​

在 ECS 上运行时,ECS 代理会自动设置 AWS_CONTAINER_CREDENTIALS_RELATIVE_URI 环境变量:

S3File {
bucket = "s3a://my-bucket"
tmp_path = "/tmp/seatunnel"
fs.s3a.endpoint = "s3.amazonaws.com"
fs.s3a.aws.credentials.provider = "com.amazonaws.auth.ContainerCredentialsProvider"
path = "/data/output"
file_format_type = "parquet"
}

EKS IRSA​

EKS IAM Roles for Service Accounts (IRSA) 需要 WebIdentityTokenCredentialsProvider 类。该类在较新的 AWS SDK v1.x 版本(如 1.11.5xx+)中可用,但 不包含 在 SeaTunnel 捆绑的旧版 AWS SDK v1.x(1.11.271)中。推荐以下替代方案:

  1. 使用 EC2 节点实例角色 — 为 EKS 工作节点附加 IAM 角色,保持默认的 InstanceProfileCredentialsProvider。
  2. 使用 SimpleAWSCredentialsProvider,从 Kubernetes Secret 注入凭据。
  3. 在所有集群节点 的 ${SEATUNNEL_HOME}/lib 中添加包含 WebIdentityTokenCredentialsProvider 的较新 AWS SDK JAR。

通过 hadoop_s3_properties 传递额外选项​

对于 provider 特定的配置键(如 fs.s3a.session.token、fs.s3a.assumed.role.arn),使用 hadoop_s3_properties 映射:

hadoop_s3_properties {
"fs.s3a.session.token" = "<session-token>"
"fs.s3a.assumed.role.arn" = "arn:aws:iam::123456789012:role/my-role"
}

连接器将这些键直接传递给 Hadoop S3A 配置。注意:连接器始终会用选项值覆盖 fs.s3a.aws.credentials.provider 键,因此无法通过 hadoop_s3_properties 覆盖它。

故障排查​

您可能会看到 Factory initialize failed(或类似的类加载)错误:这通常意味着凭据提供程序类不在 classpath 上。请确保 provider JAR 存在于 每个 集群节点(不仅仅是提交节点)的 ${SEATUNNEL_HOME}/lib 中。

No AWS Credentials provided by ...:配置的凭据提供程序无法解析凭据。请检查:

  • SimpleAWSCredentialsProvider:验证 access_key 和 secret_key 已设置。
  • InstanceProfileCredentialsProvider:验证 EC2 实例已附加 IAM 角色。
  • ContainerCredentialsProvider:验证 AWS_CONTAINER_CREDENTIALS_RELATIVE_URI 环境变量已设置。

配置解析时的 IllegalArgumentException:类名格式错误或类未实现 com.amazonaws.auth.AWSCredentialsProvider。请验证全限定类名是否正确,以及类是否实现了所需的接口。

变更日志​

Change Log
ChangeCommitVersion
[Feature][Connector-V2] Support S3 source connectivity dry-run (#12242)https://github.com/apache/seatunnel/commit/3339132553.0.0
[Improve][Connector-V2][File] Migrate HDFS and S3 sink validation to declarative OptionRule (#11881)https://github.com/apache/seatunnel/commit/5acb546183.0.0
[Feature][shade]Refactor the seatunnel-shade module. (#9993)https://github.com/apache/seatunnel/commit/4ba2895953.0.0
[Feature][Connector-V2] Enable continuous discovery for S3 and OSS file sources (#11789)https://github.com/apache/seatunnel/commit/ec92180c83.0.0
[Improve][Connector-V2] Guard POI Excel reads by file size (#11591)https://github.com/apache/seatunnel/commit/261644ff43.0.0
[Feature][Connector-File-Base] Add optional PDF RAG metadata for file source (#11571)https://github.com/apache/seatunnel/commit/2b060343d3.0.0
[Feature][Connector-V2] Support arbitrary S3A credentials provider for S3File (#11436)https://github.com/apache/seatunnel/commit/d344469973.0.0
[Feature][Connector-V2] Add recursive_file_scan option for file connectors (#10505)https://github.com/apache/seatunnel/commit/2826ea3dc3.0.0
[Feature][Connector-V2] Add optional Markdown RAG metadata for file source (#10844)https://github.com/apache/seatunnel/commit/970cadb1a3.0.0
[Feature][Connector-V2] Enable file split for S3File source (#10450)https://github.com/apache/seatunnel/commit/37999ea5b3.0.0
[Feature][seatunnel-api] Integrate Gravitino as metadata service for non-relational connectors (#10402)https://github.com/apache/seatunnel/commit/e24b8c1403.0.0
[Feature][File] Add markdown parser #9714https://github.com/apache/seatunnel/commit/8b3c07844dev
[Improve][Connector-V2] Add customizable row delimiter support for text file processing (#9608)https://github.com/apache/seatunnel/commit/7898e62e012.3.12
[Improve][Connector-V2] Support maxcompute sink writer with timestamp field type (#9234)https://github.com/apache/seatunnel/commit/a513c495e32.3.12
[improve] update file connectors config (#9034)https://github.com/apache/seatunnel/commit/8041d59dc22.3.11
[Improve][File] Add row_delimiter options into text file sink (#9017)https://github.com/apache/seatunnel/commit/92aa855a342.3.11
Revert " [improve] update localfile connector config" (#9018)https://github.com/apache/seatunnel/commit/cdc79e13ad2.3.10
[improve] update localfile connector config (#8765)https://github.com/apache/seatunnel/commit/def369a85f2.3.10
[Fix][Connector-V2] Fixed incorrectly setting s3 key in some cases (#8885)https://github.com/apache/seatunnel/commit/cf4bab5be22.3.10
[Feature][Connector-V2] Add filename_extension parameter for read/write file (#8769)https://github.com/apache/seatunnel/commit/78b23c0ef52.3.10
[Improve] restruct connector common options (#8634)https://github.com/apache/seatunnel/commit/f3499a6eeb2.3.10
[improve] update S3File connector config option (#8615)https://github.com/apache/seatunnel/commit/80cc9fa6ff2.3.10
[Feature][Connector-V2] Support create emtpy file when no data (#8543)https://github.com/apache/seatunnel/commit/275db789182.3.10
[Feature][Connector-V2] Support single file mode in file sink (#8518)https://github.com/apache/seatunnel/commit/e893deed502.3.10
[Feature][File] Support config null format for text file read (#8109)https://github.com/apache/seatunnel/commit/2dbf02df472.3.9
[Hotfix][Zeta] Fix the dependency conflict between the guava in hadoop-aws and hive-exec (#7986)https://github.com/apache/seatunnel/commit/a7837f1f192.3.9
[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)https://github.com/apache/seatunnel/commit/6b7c53d03c2.3.9
[Improve][Connector-V2] Support read archive compress file (#7633)https://github.com/apache/seatunnel/commit/3f98cd8a162.3.8
[Improve] Refactor S3FileCatalog and it's factory (#7457)https://github.com/apache/seatunnel/commit/d928e8b1132.3.8
[Improve][Connector] Add multi-table sink option check (#7360)https://github.com/apache/seatunnel/commit/2489f6446b2.3.7
[Feature][Core] Support using upstream table placeholders in sink options and auto replacement (#7131)https://github.com/apache/seatunnel/commit/c4ca74122c2.3.6
[Improve][Files] Support write fixed/timestamp as int96 of parquet (#6971)https://github.com/apache/seatunnel/commit/1a48a9c4932.3.6
[Feature][S3 File] Make S3 File Connector support multiple table write (#6698)https://github.com/apache/seatunnel/commit/8f2049b2f12.3.6
[Improve][Connector-v2] The hive connector support multiple filesystem (#6648)https://github.com/apache/seatunnel/commit/8a4c01fe352.3.6
[bigfix][S3 File]:Change the [SCHEMA] attribute of the [S3CONF class] to be non-static to avoid being reassigned after deserialization (#6717)https://github.com/apache/seatunnel/commit/79bb70101a2.3.6
[Fix][Connector-V2] Fix connector support SPI but without no args constructor (#6551)https://github.com/apache/seatunnel/commit/5f3c9c36a52.3.5
Add support for XML file type to various file connectors such as SFTP, FTP, LocalFile, HdfsFile, and more. (#6327)https://github.com/apache/seatunnel/commit/ec533ecd9a2.3.5
[Test][E2E] Add thread leak check for connector (#5773)https://github.com/apache/seatunnel/commit/1f2f3fc5f02.3.4
[Feature][Connector]add s3file save mode function (#6131)https://github.com/apache/seatunnel/commit/81c51073bf2.3.4
[Refactor][File Connector] Put Multiple Table File API to File Base Module (#6033)https://github.com/apache/seatunnel/commit/c324d663b42.3.4
Support using multiple hadoop account (#5903)https://github.com/apache/seatunnel/commit/d69d88d1aa2.3.4
[Improve][Common] Introduce new error define rule (#5793)https://github.com/apache/seatunnel/commit/9d1b2582b22.3.4
[Improve][connector-file] unifiy option between file source/sink and update document (#5680)https://github.com/apache/seatunnel/commit/8d87cf8fc42.3.4
[Feature] Support LZO compress on File Read (#5083)https://github.com/apache/seatunnel/commit/a4a19010962.3.4
[Feature][Connector-V2][File] Support read empty directory (#5591)https://github.com/apache/seatunnel/commit/1f58f224a02.3.4
Support config column/primaryKey/constraintKey in schema (#5564)https://github.com/apache/seatunnel/commit/eac76b4e502.3.4
[Feature][File Connector]optionrule FILE_FORMAT_TYPE is text/csv ,add parameter BaseSinkConfig.ENABLE_HEADER_WRITE: #5566 (#5567)https://github.com/apache/seatunnel/commit/0e02db768d2.3.4
[Feature][Connector V2][File] Add config of 'file_filter_pattern', which used for filtering files. (#5153)https://github.com/apache/seatunnel/commit/a3c13e59eb2.3.3
[chore] delete unavailable S3 & Kafka Catalogs (#4477)https://github.com/apache/seatunnel/commit/e0aec5ecec2.3.2
[Feature][ConnectorV2]add file excel sink and source (#4164)https://github.com/apache/seatunnel/commit/e3b97ae5d22.3.2
Change file type to file_format_type in file source/sink (#4249)https://github.com/apache/seatunnel/commit/973a2fae3c2.3.1
[Chore] Upgrade guava to 27.0-jre (#4238)https://github.com/apache/seatunnel/commit/4851bee5752.3.1
Add redshift datatype convertor (#4245)https://github.com/apache/seatunnel/commit/b19011517f2.3.1
Merge branch 'dev' into merge/cdchttps://github.com/apache/seatunnel/commit/4324ee19122.3.1
[Improve][Project] Code format with spotless plugin.https://github.com/apache/seatunnel/commit/423b5830382.3.1
[improve][api] Refactoring schema parse (#4157)https://github.com/apache/seatunnel/commit/b2f573a13e2.3.1
[Improve][build] Give the maven module a human readable name (#4114)https://github.com/apache/seatunnel/commit/d7cd6010512.3.1
Add S3Catalog (#4121)https://github.com/apache/seatunnel/commit/7d7f5065472.3.1
[Improve][Project] Code format with spotless plugin. (#4101)https://github.com/apache/seatunnel/commit/a2ab1665612.3.1
[Feature][Connector-V2][File] Support compress (#3899)https://github.com/apache/seatunnel/commit/55602f6b1c2.3.1
[Feature][Connector] add get source method to all source connector (#3846)https://github.com/apache/seatunnel/commit/417178fb842.3.1
[Improve][Connector-V2][File] Improve file connector option rule and document (#3812)https://github.com/apache/seatunnel/commit/bd760776692.3.1
[Feature][Shade] Add seatunnel hadoop3 uber (#3755)https://github.com/apache/seatunnel/commit/5a024bdf8f2.3.0
[Engine][Checkpoint]Unified naming style (#3714)https://github.com/apache/seatunnel/commit/bc0bd3bec32.3.0
[Connector][File-S3]Set AK is not required (#3713)https://github.com/apache/seatunnel/commit/da3c5261722.3.0
[Connector&Engine]Set S3 AK to optional (#3688)https://github.com/apache/seatunnel/commit/4710918b022.3.0
[Connector][S3]Support s3a protocol (#3632)https://github.com/apache/seatunnel/commit/ae4cc9c1ec2.3.0
[Hotfix][OptionRule] Fix option rule about all connectors (#3592)https://github.com/apache/seatunnel/commit/226dc6a1192.3.0
[Improve][Connector-V2][File] Unified excetion for file source & sink connectors (#3525)https://github.com/apache/seatunnel/commit/031e8e263c2.3.0
[Feature][Connector-V2][File] Add option and factory for file connectors (#3375)https://github.com/apache/seatunnel/commit/db286e86312.3.0
[Improve][Connector-V2][File] Improve code structure (#3238)https://github.com/apache/seatunnel/commit/dd5c3538812.3.0
[Connector-V2][ElasticSearch] Add ElasticSearch Source/Sink Factory (#3325)https://github.com/apache/seatunnel/commit/38254e3f262.3.0
[Feature][Connector-V2][S3] Add S3 file source & sink connector (#3119)https://github.com/apache/seatunnel/commit/f27d68ca9c2.3.0-beta