☰
SeaTunnel HdfsFile Sink 深度实践:将数据以精确一次语义写入 HDFS 的完整配置指南
2026/10/9 1:18:39 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

HdfsFile 是 SeaTunnel 内置的 HDFS 文件接收器(Sink)插件,用于把上游 Source/Transform 输出的数据以文本、CSV、Parquet、ORC、JSON、Excel 等格式写入 Hadoop 分布式文件系统。本文以官方中文文档 HdfsFile.md 为主线,结合仓库中connector-file-hadoop、connector-file-base的实际源码与测试,完整讲解其支持的引擎、全部接收器选项、2PC 精确一次写入原理、Kerberos 认证配置以及各类典型任务示例,帮助读者在生产环境中正确、可靠地把数据落盘到 HDFS。

支持的引擎与主要特性

HdfsFile Sink 同时支持三种计算引擎,可在对应环境下直接使用:

  • Spark
  • Flink
  • SeaTunnel Zeta(SeaTunnel 自研分布式引擎)

主要特性包括:

  • 精确一次(Exactly-Once):默认通过两阶段提交(2PC)保证数据写入 HDFS 时不丢失、不重复。特性说明可参见 Connector-V2 特性文档。
  • 多种文件格式:文本(text/txt)、CSV、Parquet、ORC、JSON、Excel(xlsx)。
  • 压缩编解码器:支持 lzo 等多种压缩(具体随文件格式不同而不同,详见下文选项表)。

在源码层面,插件的入口类为 HdfsFileSink.java,其继承自connector-file-base-hadoop模块中的BaseHdfsFileSink,通过@AutoService(SeaTunnelSink.class)注册为名为HdfsFile的插件;对应的工厂类 HdfsFileSinkFactory.java 则负责声明该插件接受的参数规则(optionRule())。

支持的数据源信息

数据源支持的版本
Hdfs 文件Hadoop 2.x 和 3.x

需要说明的是,文档中的"支持的数据源信息"指该接收器可对接的 Hadoop 集群版本;若你使用 Spark/Flink 引擎,还需确保对应集群已集成 Hadoop 依赖(详见下文"运行环境提示")。

接收器选项详解

下表完整列出了 HdfsFile Sink 的全部配置项(名称、类型、是否必须、默认值及说明)。其中大部分选项定义在连接器公共模块 BaseSinkConfig.java 中,HDFS 特有参数(如fs.defaultFS)则在 BaseHdfsFileSink.java 的prepare()阶段被解析并装配进HadoopConf。

名称类型是否必须默认值描述
fs.defaultFSstring是-以hdfs://开头的 Hadoop 集群地址,例如hdfs://hadoopcluster
pathstring是-目标目录路径,最终数据文件将被写入该目录
tmp_pathstring是/tmp/seatunnel结果文件首先写入该临时路径,提交时再通过mv命令把临时目录移动到目标目录,需要 HDFS 路径
hdfs_site_pathstring否-hdfs-site.xml的路径,用于加载 namenode 的 HA 配置
custom_filenameboolean否false是否需要自定义文件名
file_name_expressionstring否"${transactionId}"仅在custom_filename为 true 时使用,描述将创建到path中的文件表达式,可加入变量${now}(当前时间,格式由filename_time_format定义)或${uuid},例如test_${uuid}_${now};注意若is_enable_transaction为 true,文件头部会自动加上${transactionId}_
filename_time_formatstring否"yyyy.MM.dd"仅在custom_filename为 true 时使用,指定file_name_expression中${now}的时间格式;常用格式符号:y=年,M=月,d=月中的一天,H=一天中的小时(0-23),m=小时中的分钟,s=分钟中的秒
file_format_typestring否"csv"支持的文件类型:text、json、csv、orc、parquet、excel;最终文件名会以对应后缀结尾,其中 text 文件的后缀是txt
field_delimiterstring否'\001'仅 text 文件格式使用,数据行中列之间的分隔符
row_delimiterstring否"\n"仅 text 文件格式使用,文件中行之间的分隔符
have_partitionboolean否false是否需要处理分区
partition_byarray否-仅在have_partition为 true 时使用,根据选定的字段对数据进行分区
partition_dir_expressionstring否"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/"仅在have_partition为 true 时使用;指定partition_by后按分区信息生成分区目录并把最终文件放入分区目录;k0是第一个分区字段,v0是其值
is_partition_field_write_in_fileboolean否false仅当have_partition为 true 时使用;若为 true,分区字段及其值会写入数据文件(写 Hive 数据文件时应设为 false)
sink_columnsarray否空当为空时,所有字段都是接收器列;需要写入文件的列,默认取Transform或Source输出的所有列,字段顺序决定实际写入文件时的顺序
is_enable_transactionboolean否true为 true 时,写入目标目录过程中保证数据不丢失、不重复;注意为 true 时文件头部会自动加上${transactionId}_;目前仅支持 true
batch_sizeint否1000000单个文件中的最大行数;对 SeaTunnel Engine,文件行数由batch_size与checkpoint.interval共同决定:若 checkpoint 间隔足够大,writer 会一直写入直到行数超过batch_size;若 checkpoint 间隔很小,则每次新 checkpoint 触发时都会新建文件
compress_codecstring否none文件压缩编解码器。按文件格式支持如下:txt:lzo、none;json:lzo、none;csv:lzo、none;orc:lzo、snappy、lz4、zlib、none;parquet:lzo、snappy、lz4、gzip、brotli、zstd、none。提示:excel 类型不支持任何压缩格式
krb5_pathstring否/etc/krb5.confKerberos 的 krb5 配置路径
kerberos_principalstring否-Kerberos 主体(principal)
kerberos_keytab_pathstring否-Kerberos 的 keytab 路径
common-optionsobject否-接收器插件通用参数(source_table_name、parallelism),详见 接收器通用选项
max_rows_in_memoryint否-仅当file_format为 excel 时使用,Excel 格式下可缓存在内存中的最大数据项数
sheet_namestring否Sheet${随机数}仅当file_format为 excel 时使用,将工作簿写入指定的表名

选项背后的源码实现

从 BaseSinkConfig.java 可以看到各选项的默认值均与上表一一对应,例如:

  • DEFAULT_FIELD_DELIMITER取自文本格式常量TextFormatConstant.SEPARATOR[0](即\u0001,对应文档中的'\001');
  • DEFAULT_ROW_DELIMITER = "\n";
  • DEFAULT_TMP_PATH = "/tmp/seatunnel";
  • DEFAULT_FILE_NAME_EXPRESSION = "${transactionId}";
  • DEFAULT_BATCH_SIZE = 1000000;
  • DEFAULT_PARTITION_DIR_EXPRESSION = "${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/";
  • KRB5_PATH默认/etc/krb5.conf。

在 HdfsFileSinkFactory.java 的optionRule()中,fs.defaultFS与path被声明为必填项,其余选项按file_format_type、custom_filename、have_partition等前置条件做了条件化声明——例如只有选择text格式时才需要field_delimiter/row_delimiter,只有custom_filename=true时才需要file_name_expression/filename_time_format,只有have_partition=true时才需要partition_by等,配置校验失败时任务会在启动阶段报错。对应的工厂测试见 HdfsFileFactoryTest.java。

文件格式枚举 FileFormat.java 定义了每种格式对应的写策略(WriteStrategy):CSV/TEXT 使用TextWriteStrategy,PARQUET 使用ParquetWriteStrategy,ORC 使用OrcWriteStrategy,JSON 使用JsonWriteStrategy,EXCEL 使用ExcelWriteStrategy,并暴露getSuffix()返回真实文件后缀(.txt、.csv、.parquet、.orc、.json、.xlsx)。压缩选项同样按格式受限:TXT_COMPRESS仅允许none/lzo,ORC_COMPRESS允许none/lzo/snappy/lz4/zlib,PARQUET_COMPRESS额外支持gzip/brotli/zstd。

HDFS 特有参数与 Kerberos 支持

HdfsFile 在 BaseHdfsFileSink.java 的prepare()中完成连接初始化:

  1. 校验fs.defaultFS必须存在(缺失时抛出FileConnectorException,错误类型为CONFIG_VALIDATION_FAILED);
  2. 以fs.defaultFS构造HadoopConf;
  3. 若配置了hdfs_site_path,则加载到HadoopConf,用于 namenode HA 场景下的hdfs-site.xml读取;
  4. 若配置了remote_user,则设置远程用户;
  5. 若配置了krb5_path、kerberos_principal、kerberos_keytab_path,则完成 Kerberos 登录信息装配。

实际的文件系统操作(建目录、写文件、mv提交等)由connector-file-base模块中的HadoopFileSystemProxy封装执行(HadoopFileSystemProxy.java),它基于 HadoopFileSystemAPI 并统一管理UserGroupInformation。

写入流程与"精确一次"的实现原理

HdfsFile Sink 采用"临时目录 + 两阶段提交"机制保证精确一次:

  1. 数据首先写入tmp_path下的事务目录(transaction directory)。从 AbstractWriteStrategy.java 的实现可见,transactionId的格式为T_{jobId}_{uuidPrefix}_{subTaskIndex}_{checkpointId},即由任务 ID、UUID 前缀、子任务索引与检查点 ID 拼接而成,保证同一任务不同文件 Sink、不同子任务、不同检查点之间不会冲突;
  2. 检查点提交时,通过FileSinkAggregatedCommitter(FileSinkAggregatedCommitter.java)把临时目录下的文件以mv操作移动到目标path对应的分区目录;
  3. 若失败(abort),则直接删除对应事务目录,实现无副作用回滚。

这就是文档所说"默认情况下使用 2PC 提交来确保精确一次"的落地方式。也是tmp_path必须设置为 HDFS 路径的原因——mv发生在同一个分布式文件系统内部。

文件名生成逻辑(AbstractWriteStrategy.java)会先按file_name_expression替换${now}、${uuid}、${transactionId}变量,再追加_${partId}与格式后缀、压缩后缀;当is_enable_transaction=true时,文件名的${transactionId}前缀正好用于保证提交过程中文件不重名。

任务示例

简单示例:FakeSource 生成数据写入 ORC 文件

此示例定义了一个 SeaTunnel 同步任务,通过 FakeSource 自动生成数据并发送到 HDFS(BATCH 模式):

# 定义运行时环境 env { parallelism = 1 job.mode = "BATCH" } source { # 这是一个示例源插件,仅用于测试和演示功能源插件 FakeSource { parallelism = 1 result_table_name = "fake" row.num = 16 schema = { fields { c_map = "map<string, smallint>" c_array = "array<int>" c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(30, 8)" c_bytes = bytes c_date = date c_timestamp = timestamp } } } } transform { # 转换插件可在此配置,更多源/转换插件参见仓库 docs/en/connector-v2/source 与 docs/en/transform-v2 目录 } sink { HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" file_format_type = "orc" } }

ORC 文件格式的简单配置

最精简的 ORC 写法只需三个选项:

HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" file_format_type = "orc" }

Text 文件格式:分区 + 自定义文件名 + 指定列

以下配置演示了 text 格式下have_partition、custom_filename、sink_columns的组合使用,数据按age字段分区,分区目录表达式为${k0}=${v0}(即age=18这样的目录),分区字段同时写入文件,文件名由事务 ID 加当前日期构成:

HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" 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 }

Parquet 文件格式:分区 + 自定义文件名 + 指定列

Parquet 的用法与 text 基本一致,仅切换file_format_type:

HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" 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" file_format_type = "parquet" sink_columns = ["name","age"] is_enable_transaction = true }

Kerberos 认证配置

在开启 Kerberos 的 Hadoop 集群中,通过hdfs_site_path加载 HA 配置,并指定 principal 与 keytab:

HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" hdfs_site_path = "/path/to/your/hdfs_site_path" kerberos_principal = "your_principal@EXAMPLE.COM" kerberos_keytab_path = "/path/to/your/keytab/file.keytab" }

压缩配置

启用 LZO 压缩:

HdfsFile { fs.defaultFS = "hdfs://hadoopcluster" path = "/tmp/hive/warehouse/test2" compress_codec = "lzo" }

注意:压缩格式必须与file_format_type匹配(例如 ORC 可选snappy/lz4/zlib,Parquet 可选gzip/brotli/zstd等,excel 不支持压缩),否则会因参数校验失败而无法启动任务。

运行环境提示

如果你使用 Spark/Flink,为了使用此连接器,必须确保你的 Spark/Flink 集群已经集成了 Hadoop。文档记录的已验证 Hadoop 版本是 2.x。 如果你使用 SeaTunnel Engine(Zeta),在下载和安装 SeaTunnel Engine 时会自动集成 Hadoop jar,可以检查${SEATUNNEL_HOME}/lib下的 jar 包来确认这一点。

因此,在实际部署前建议按以下顺序自检:

  1. 确认引擎与 Hadoop 的集成方式(Spark/Flink 集群自带 Hadoop;SeaTunnel Engine 则检查${SEATUNNEL_HOME}/lib下的 hadoop 相关 jar);
  2. 确认fs.defaultFS可被任务所在节点解析(HA 集群建议同时配置hdfs_site_path指向包含 namenode HA 配置的hdfs-site.xml);
  3. 若集群开启了 Kerberos,确保krb5_path、kerberos_principal、kerberos_keytab_path配置正确且 keytab 文件对运行用户可读;
  4. 确认tmp_path与path均为 HDFS 路径,且运行用户对该临时目录和目标目录具有写权限。

小结

HdfsFile Sink 是 SeaTunnel 连接 HDFS 的标准文件接收器:它通过fs.defaultFS+path定位目标目录,通过file_format_type选择 text/csv/parquet/orc/json/excel 等落盘格式,通过have_partition/partition_by/partition_dir_expression组织分区目录,通过custom_filename/file_name_expression自定义文件名,并通过is_enable_transaction(默认开启)结合临时目录与 2PC 提交实现精确一次写入;对安全集群还提供了完整的 Kerberos 配置入口。结合本文给出的源码链路与示例,读者可以按需组合这些参数,构建出可靠、可复用的 HDFS 落盘任务。

  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

上一篇:【亲测免费】 Apache DataFusion Python 绑定教程
下一篇:探索Prometheus Flask Exporter:监控你的Flask应用新维度

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询