- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
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.defaultFS | string | 是 | - | 以hdfs://开头的 Hadoop 集群地址,例如hdfs://hadoopcluster |
| path | string | 是 | - | 目标目录路径,最终数据文件将被写入该目录 |
| tmp_path | string | 是 | /tmp/seatunnel | 结果文件首先写入该临时路径,提交时再通过mv命令把临时目录移动到目标目录,需要 HDFS 路径 |
| hdfs_site_path | string | 否 | - | hdfs-site.xml的路径,用于加载 namenode 的 HA 配置 |
| custom_filename | boolean | 否 | false | 是否需要自定义文件名 |
| file_name_expression | string | 否 | "${transactionId}" | 仅在custom_filename为 true 时使用,描述将创建到path中的文件表达式,可加入变量${now}(当前时间,格式由filename_time_format定义)或${uuid},例如test_${uuid}_${now};注意若is_enable_transaction为 true,文件头部会自动加上${transactionId}_ |
| filename_time_format | string | 否 | "yyyy.MM.dd" | 仅在custom_filename为 true 时使用,指定file_name_expression中${now}的时间格式;常用格式符号:y=年,M=月,d=月中的一天,H=一天中的小时(0-23),m=小时中的分钟,s=分钟中的秒 |
| file_format_type | string | 否 | "csv" | 支持的文件类型:text、json、csv、orc、parquet、excel;最终文件名会以对应后缀结尾,其中 text 文件的后缀是txt |
| field_delimiter | string | 否 | '\001' | 仅 text 文件格式使用,数据行中列之间的分隔符 |
| row_delimiter | string | 否 | "\n" | 仅 text 文件格式使用,文件中行之间的分隔符 |
| have_partition | boolean | 否 | false | 是否需要处理分区 |
| partition_by | array | 否 | - | 仅在have_partition为 true 时使用,根据选定的字段对数据进行分区 |
| partition_dir_expression | string | 否 | "${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/" | 仅在have_partition为 true 时使用;指定partition_by后按分区信息生成分区目录并把最终文件放入分区目录;k0是第一个分区字段,v0是其值 |
| is_partition_field_write_in_file | boolean | 否 | false | 仅当have_partition为 true 时使用;若为 true,分区字段及其值会写入数据文件(写 Hive 数据文件时应设为 false) |
| sink_columns | array | 否 | 空 | 当为空时,所有字段都是接收器列;需要写入文件的列,默认取Transform或Source输出的所有列,字段顺序决定实际写入文件时的顺序 |
| is_enable_transaction | boolean | 否 | true | 为 true 时,写入目标目录过程中保证数据不丢失、不重复;注意为 true 时文件头部会自动加上${transactionId}_;目前仅支持 true |
| batch_size | int | 否 | 1000000 | 单个文件中的最大行数;对 SeaTunnel Engine,文件行数由batch_size与checkpoint.interval共同决定:若 checkpoint 间隔足够大,writer 会一直写入直到行数超过batch_size;若 checkpoint 间隔很小,则每次新 checkpoint 触发时都会新建文件 |
| compress_codec | string | 否 | 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_path | string | 否 | /etc/krb5.conf | Kerberos 的 krb5 配置路径 |
| kerberos_principal | string | 否 | - | Kerberos 主体(principal) |
| kerberos_keytab_path | string | 否 | - | Kerberos 的 keytab 路径 |
| common-options | object | 否 | - | 接收器插件通用参数(source_table_name、parallelism),详见 接收器通用选项 |
| max_rows_in_memory | int | 否 | - | 仅当file_format为 excel 时使用,Excel 格式下可缓存在内存中的最大数据项数 |
| sheet_name | string | 否 | 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()中完成连接初始化:
- 校验
fs.defaultFS必须存在(缺失时抛出FileConnectorException,错误类型为CONFIG_VALIDATION_FAILED); - 以
fs.defaultFS构造HadoopConf; - 若配置了
hdfs_site_path,则加载到HadoopConf,用于 namenode HA 场景下的hdfs-site.xml读取; - 若配置了
remote_user,则设置远程用户; - 若配置了
krb5_path、kerberos_principal、kerberos_keytab_path,则完成 Kerberos 登录信息装配。
实际的文件系统操作(建目录、写文件、mv提交等)由connector-file-base模块中的HadoopFileSystemProxy封装执行(HadoopFileSystemProxy.java),它基于 HadoopFileSystemAPI 并统一管理UserGroupInformation。
写入流程与"精确一次"的实现原理
HdfsFile Sink 采用"临时目录 + 两阶段提交"机制保证精确一次:
- 数据首先写入
tmp_path下的事务目录(transaction directory)。从 AbstractWriteStrategy.java 的实现可见,transactionId的格式为T_{jobId}_{uuidPrefix}_{subTaskIndex}_{checkpointId},即由任务 ID、UUID 前缀、子任务索引与检查点 ID 拼接而成,保证同一任务不同文件 Sink、不同子任务、不同检查点之间不会冲突; - 检查点提交时,通过
FileSinkAggregatedCommitter(FileSinkAggregatedCommitter.java)把临时目录下的文件以mv操作移动到目标path对应的分区目录; - 若失败(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 包来确认这一点。
因此,在实际部署前建议按以下顺序自检:
- 确认引擎与 Hadoop 的集成方式(Spark/Flink 集群自带 Hadoop;SeaTunnel Engine 则检查
${SEATUNNEL_HOME}/lib下的 hadoop 相关 jar); - 确认
fs.defaultFS可被任务所在节点解析(HA 集群建议同时配置hdfs_site_path指向包含 namenode HA 配置的hdfs-site.xml); - 若集群开启了 Kerberos,确保
krb5_path、kerberos_principal、kerberos_keytab_path配置正确且 keytab 文件对运行用户可读; - 确认
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.
相关推荐
SeaTunnel HdfsFile Sink 实战指南:从基础配置到精确一次写入
SeaTunnel HdfsFile Sink 实战指南:从基础配置到精确一次写入 本篇技术指南系统讲解 Apache SeaTunnel 中 HdfsFile
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel HdfsFile Sink Connector 使用指南:将数据写入 HDFS 的完整配置与原理剖析
SeaTunnel HdfsFile Sink Connector 使用指南:将数据写入 HDFS 的完整配置与原理剖析 HdfsFile 是 SeaTunne
数据工程大数据批处理流处理先验证再激活:MAS 激活工具完成 Windows 10/11 与 Office 激活的完整操作
先验证再激活:MAS 激活工具完成 Windows 10/11 与 Office 激活的完整操作 激活页面显示“Windows 已激活”,检查脚本对 Windo
操作系统
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考