跳到主要内容
版本:Next

BosFile

BOS 文件 Sink 连接器

支持引擎

Spark
Flink
SeaTunnel Zeta

描述

通过 BOS HDFS SDK 将数据写入百度智能云 BOS。

提示

使用 Spark/Flink 时,需确保集群已集成 Hadoop 2.x。

使用 SeaTunnel Engine 时,Hadoop 相关 jar 已包含在 ${SEATUNNEL_HOME}/lib 中。

使用本连接器需将 bos-hdfs-sdk(>= 1.0.4-community)放入 ${SEATUNNEL_HOME}/lib。下载:bos-hdfs-sdk-1.0.4-community.jar.zip

主要特性

  • 多模态

    使用 binary 格式读写任意类型文件,可将任意文件同步到目标位置。

  • 精确一次

    默认通过 2PC 提交保证 exactly-once

  • cdc

  • 支持多表写入

  • 定时 flush

  • 文件格式

    • text
    • csv
    • parquet
    • orc
    • json
    • excel
    • xml
    • binary
    • canal_json
    • debezium_json
    • maxwell_json

配置项

名称类型必填默认值描述
pathstring-bucket 内写入目录
tmp_pathstring/tmp/seatunnel先写入临时目录,再通过 mv 提交到目标目录
bucketstring-BOS bucket,例如 bos://my-bucket
access_keystring-BOS Access Key
secret_keystring-BOS Secret Key
endpointstring-BOS Endpoint,例如 http://bj.bcebos.com
custom_filenamebooleanfalse是否自定义文件名
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 等
filename_extensionstring-自定义文件扩展名
field_delimiterstringtext 为 \001,csv 为 ,text/csv 格式使用
row_delimiterstring"\n"text/csv/json 格式使用
have_partitionbooleanfalse是否按分区写入
partition_byarray-have_partition 为 true 时使用
partition_dir_expressionstring"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/"have_partition 为 true 时使用
is_partition_field_write_in_filebooleanfalsehave_partition 为 true 时使用
sink_columnsarray为空时写入全部字段
is_enable_transactionbooleantrue开启后保证写入不丢不重,文件名自动加 ${transactionId}_ 前缀
batch_sizeint1000000单文件最大行数
compress_codecstringnone压缩格式
xml_root_tagstringRECORDSxml 格式使用
xml_row_tagstringRECORDxml 格式使用
xml_use_attr_formatboolean-xml 格式使用
single_file_modebooleanfalse每个并行度只输出一个文件
create_empty_file_when_no_databooleanfalse上游无数据时仍生成空文件
parquet_avro_write_timestamp_as_int96booleanfalseparquet 格式使用
parquet_avro_write_fixed_as_int96array-parquet 格式使用
encodingstring"UTF-8"json/text/csv/xml 格式使用

示例

分区写入 text 文件:

sink {
BosFile {
path = "/sink"
bucket = "bos://sink-bucket"
access_key = "your-access-key"
secret_key = "your-secret-key"
endpoint = "http://bj.bcebos.com"
file_format_type = "text"
have_partition = true
partition_by = ["age"]
partition_dir_expression = "${k0}=${v0}"
is_enable_transaction = true
}
}

简单 text 写入:

sink {
BosFile {
bucket = "bos://sink-bucket"
path = "/warehouse/table/"
file_format_type = "text"
access_key = "your-access-key"
secret_key = "your-secret-key"
endpoint = "http://bj.bcebos.com"
row_delimiter = "\n"
field_delimiter = ","
is_enable_transaction = true
}
}

变更日志

变更日志
变更Commit版本
[Improve][Connector-V2] 新增 Hive BOSStorage,对齐 BosFile e2e/文档至 CosFilehttps://github.com/apache/seatunnel/pull/11952dev
[Feature][Connector-V2] 新增 BosFile Source/Sink 连接器https://github.com/apache/seatunnel/pull/11952dev