跳到主要内容
版本:3.0.0

CosFile

Cos 文件接收器连接器

支持这些引擎​

Spark
Flink
SeaTunnel Zeta

描述​

将数据输出到腾讯云 COS(云对象存储)文件系统。

:::提示

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

如果你使用 SeaTunnel Engine,当你下载并安装 SeaTunnel Engine 时,它会自动集成 Hadoop jar。你可以在 ${SEATUNNEL_HOME}/lib 下检查 jar 包以确认这一点。

要使用此连接器,你需要将 hadoop-cos-{hadoop.version}-{version}.jar 和 cos_api-bundle-{version}.jar 放在 ${SEATUNNEL_HOME}/lib 目录中,下载地址:Hadoop COS release。它仅支持 Hadoop 2.6.5+ 和 hadoop-cos 8.0.2+。

:::

关键特性​

  • 多模态

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

  • 精确一次

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

  • 文件格式类型
    • text
    • csv
    • parquet
    • orc
    • json
    • excel
    • xml
    • binary
    • canal_json
    • debezium_json
    • maxwell_json

选项​

名称类型必需默认值描述
pathstring是-Sink 写入 COS 桶内的目标目录。配合 bucket,实际路径为 cosn://<bucket><path>。
tmp_pathstring否/tmp/seatunnel结果文件将首先写入tmp路径,然后使用“mv”将tmp目录提交到目标目录。需要一个COS目录.
bucketstring是-COS 文件系统的桶地址,例如 cosn://seatunnel-test-1259587829。
secret_idstring是-腾讯云 COS 的 SecretId。
secret_keystring是-腾讯云 COS 的 SecretKey。
regionstring是-COS 桶所在地域,例如 ap-chengdu。
custom_filenameboolean否false是否需要自定义文件名
file_name_expressionstring否"${transactionId}"仅在custom_filename为true时使用
filename_time_formatstring否"yyyy.MM.dd"仅在custom_filename为true时使用
file_format_typestring否"csv"文件格式类型,支持:text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json。
filename_extensionstring否-使用自定义的文件扩展名覆盖默认的文件扩展名。 例如:.xml, .json, dat, .customtype
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否当此参数为空时,所有字段都是接收列
is_enable_transactionboolean否true若为 true,写入目标目录的数据不会丢失或重复;当为 true 时,会自动在文件名前缀添加 ${transactionId}_。
batch_sizeint否1000000单个文件的最大行数。对于 SeaTunnel Engine,文件中的行数由 batch_size 和 checkpoint.interval 共同决定。
compress_codecstring否none文件的压缩编解码器。Excel 格式不支持任何压缩格式。
common-optionsobject否-Sink 插件通用参数,请参考 Sink Common Options 了解详情。
max_rows_in_memoryint否-仅在file_format为excel时使用.
sheet_max_rowsint否1048576仅在 file_format_type 为 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_row_tagstring否RECORD仅在file_format为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时使用.
encodingstring否"UTF-8"仅当file_format_type为json、text、csv、xml时使用.
merge_update_eventboolean否false仅当file_format_type为canal_json、debezium_json、maxwell_json 时使用。设置为 true 时,会将 UPDATE_AFTER 与 UPDATE_BEFORE 合并为 UPDATE 事件数据。
schema_evolution_enabledboolean否false开启 Schema 演变支持,适用于 CDC 管道。为 true 时,来自上游的 ADD/DROP/RENAME/MODIFY 列事件无需重启作业即可应用到 Sink。不支持 binary 格式。

path [string]​

目标目录路径是必需的.

bucket [string]​

cos文件系统的bucket地址,例如:cosn://seatunnel-test-1259587829

secret_id [string]​

cos文件系统的密钥id.

secret_key [string]​

cos文件系统的密钥.

region [string]​

cos文件系统的分区.

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

请注意,最终文件名将以 file_format 的后缀结尾, 文本文件的后缀为 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引擎,文件中的行数由 batch_size 和 checkpoint.interval 共同决定. 如果 checkpoint.interval 的值足够大, 接收器写入程序将在文件中写入行,直到文件中的行大于 batch_size. 如果 checkpoint.interval 较小, 则接收器写入程序将在新的检查点触发时创建一个新文件.

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

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

common options​

接收器写入插件常用参数,请参考 Sink Common Options 了解详细信息.

max_rows_in_memory [int]​

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

sheet_max_rows [int]​

仅在 file_format_type 为 excel 时使用。该选项限制每个工作表可以写入的最大行数,默认值为 1048576。

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_avro_write_fixed_as_int96 [array]​

支持从12字节字段写入Parquet INT96,仅适用于拼花地板文件.

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 不会合并;

示例​

对于具有 have_partition 、 custom_filename 和 sink_columns 的文本文件格式


CosFile {
path="/sink"
bucket = "cosn://seatunnel-test-1259587829"
secret_id = "xxxxxxxxxxxxxxxxxxx"
secret_key = "xxxxxxxxxxxxxxxxxxx"
region = "ap-chengdu"
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
}

适用于带有have_partition 和 sink_columns的parquet 文件格式`


CosFile {
path="/sink"
bucket = "cosn://seatunnel-test-1259587829"
secret_id = "xxxxxxxxxxxxxxxxxxx"
secret_key = "xxxxxxxxxxxxxxxxxxx"
region = "ap-chengdu"
have_partition = true
partition_by = ["age"]
partition_dir_expression = "${k0}=${v0}"
is_partition_field_write_in_file = true
file_format_type = "parquet"
sink_columns = ["name","age"]
}

对于orc文件格式的简单配置


CosFile {
path="/sink"
bucket = "cosn://seatunnel-test-1259587829"
secret_id = "xxxxxxxxxxxxxxxxxxx"
secret_key = "xxxxxxxxxxxxxxxxxxx"
region = "ap-chengdu"
file_format_type = "orc"
}

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 管道中的使用示例:

LocalFile {
path = "/tmp/cdc/${table_name}"
file_format_type = "parquet"
schema_evolution_enabled = true
have_partition = true
partition_by = ["updated_at_month"]
}

变更日志​

Change Log
ChangeCommitVersion
[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] 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
[Improve] File souce refactor (#10758)https://github.com/apache/seatunnel/commit/bf25292563.0.0
[Improve][Connectors-v2] File sink refactor (#10587)https://github.com/apache/seatunnel/commit/efeed28ae3.0.0