
SeaTunnel HdfsFile Sink 连接器完全指南从基础配置到 HA/Kerberos/ViewFS 实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 HdfsFile Sink 连接器用于将上游数据写入 HDFS 文件系统支持 text、csv、parquet、orc、json、excel、xml、binary 八种文件格式并提供事务提交、分区写入、自定义文件名、Kerberos 认证与 schema 演进等能力。本文以官方文档 docs/en/connectors/sink/HdfsFile.md 为主体结合仓库中connector-file-hadoop与connector-file-base模块的源码系统讲解该连接器的全部配置项、底层实现原理与可运行的配置示例帮助你直接落地到实际同步任务中。支持引擎与关键特性HdfsFile Sink 同时支持三种引擎SparkFlinkSeaTunnel Zeta引擎自带 Hadoop 依赖无需额外集成关键特性一览对应 connector-v2-features多模态multimodal底层使用二进制文件格式读写任意格式文件视频、图片等可实现任意文件的同步。精确一次exactly-once默认使用两阶段提交2PC保证数据不丢失、不重复。多表写入support multiple table write一个任务可同时向多张目标表写入。文件格式text、csv、parquet、orc、json、excel、xml、binary。压缩编码lzo、canal_json、debezium_json、maxwell_json其中 canal_json / debezium_json / maxwell_json 实际是 CDC 事件的行序列化格式与merge_update_event配合使用。不支持 timer flush文件刷新时机完全由 checkpoint 与batch_size决定。支持的 HDFS 版本Hadoop 2.x 与 3.x。使用提示若在 Spark/Flink 上使用必须确保集群已集成 Hadoop文档中已验证的版本为 2.x若使用 SeaTunnel Zeta 引擎安装包lib目录下已自动包含 Hadoop jar可通过检查${SEATUNNEL_HOME}/lib下的 jar 包确认。源码结构HdfsFile Sink 的实现骨架在深入配置之前先了解该连接器在仓库中的代码组织方便后续对照连接器入口模块connector-file-hadoop约 20 个 Java 文件仅实现 HDFS 特有逻辑文件连接器通用模块connector-file-base文件读写、事务、分区、格式策略等通用实现核心类调用链如下组件文件职责插件入口HdfsFileSink.java继承BaseMultipleTableFileSink返回插件名HdfsFile插件工厂HdfsFileSinkFactory.java声明OptionRule条件化校验、构建HadoopConf配置定义HdfsFileSinkOptions.java直接继承FileBaseSinkOptions无新增项全部 Sink 选项FileBaseSinkOptions.java所有参数定义、默认值与约束约 40 个 OptionHadoop 客户端配置HadoopConf.java构造Configuration区分 hdfs / viewfs scheme写入器BaseFileSinkWriter.java事务恢复、write/prepareCommit/snapshotState提交器FileSinkAggregatedCommitter.java通过renameFile把临时文件 mv 到目标目录完成提交写策略AbstractWriteStrategy.java文件名生成、事务目录管理、分区目录计算Sink 参数完整参考表以下参数即官方文档全量整理默认值与约束与 FileBaseSinkOptions.java 中的定义一一对应。名称类型必填默认值说明fs.defaultFSstring是-Hadoop 集群地址。支持hdfs://hadoopcluster、hdfs://namenode:9000标准 HDFS、viewfs://myclusterViewFS 联邦 HDFSViewFS 配置示例见下文pathstring是-目标目录路径必填tmp_pathstring是/tmp/seatunnel结果文件先写入临时路径提交时通过mv移动到目标目录必须是 HDFS 路径hdfs_site_pathstring否-hdfs-site.xml的路径用于加载 NameNode 的 HA 配置custom_filenameboolean否false是否需要自定义文件名file_name_expressionstring否${transactionId}仅当custom_filenametrue时生效。描述生成到path的文件名表达式支持${now}与${uuid}变量如test_${uuid}_${now}${now}的格式由filename_time_format定义。注意当is_enable_transactiontrue时会自动在文件名头部追加${transactionId}_filename_time_formatstring否yyyy.MM.dd仅当custom_filenametrue时生效。指定file_name_expression中${now}的时间格式。常用占位符y年、M月、d日、H小时 0-23、m分钟、s秒file_format_typestring否csv文件类型text、csv、parquet、orc、json、excel、xml、binary。最终文件名会以文件格式后缀结尾其中 text 文件后缀为txtfilename_extensionstring否-覆盖默认文件名后缀例如.xml、.json、dat、.customtypefield_delimiterstring否text 为 \001csv 为 ,仅 text/csv 格式生效字段间分隔符。源码中默认值来自TextFormatConstant.SEPARATOR[0]\001row_delimiterstring否\n仅 text、csv、json 格式生效行间分隔符have_partitionboolean否false是否进行分区处理partition_byarray否-仅当have_partitiontrue时生效按所选字段分区partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/仅当have_partitiontrue时生效根据分区信息生成分区目录k0为第一个分区字段v0为其值is_partition_field_write_in_fileboolean否false仅当have_partitiontrue时生效。为true时分区字段及其值会写入数据文件例如要写 Hive 数据文件该值应为falsesink_columnsarray否空全部字段需要写入文件的列。为空时写入Transform或Source传来的所有字段字段顺序决定实际写入顺序is_enable_transactionboolean否true为true时保证写入目标目录的数据不丢失、不重复且自动在文件名头部追加${transactionId}_。当前仅支持truebatch_sizeint否1000000单个文件的最大行数。SeaTunnel Zeta 引擎下文件行数由batch_size与checkpoint.interval共同决定若 checkpoint 间隔足够大writer 会持续写入直到超过batch_size若 checkpoint 间隔小则每次 checkpoint 触发时新建文件compress_codecstring否none文件压缩编码。各格式支持txtlzo、nonejsonlzo、nonecsvlzo、noneorclzo、snappy、lz4、zlib、noneparquetlzo、snappy、lz4、gzip、brotli、zstd、none。excel 不支持任何压缩krb5_pathstring否/etc/krb5.confKerberos 的 krb5 配置文件路径kerberos_principalstring否-Kerberos principalkerberos_keytab_pathstring否-Kerberos keytab 路径common-optionsobject否-Sink 插件公共参数详见 Sink Common Optionsmax_rows_in_memoryint否-仅 excel 格式生效内存中可缓存的最大数据条数sheet_max_rowsint否1048576仅 excel 格式生效sheet_namestring否Sheet${随机数}仅 excel 格式生效写入的工作表名称csv_string_quote_modeenum否MINIMAL仅 csv 格式生效字符串引号模式见下文详解xml_root_tagstring否RECORDS仅 xml 格式生效XML 根元素标签名xml_row_tagstring否RECORD仅 xml 格式生效数据行标签名xml_use_attr_formatboolean否-仅 xml 格式生效是否使用标签属性格式处理数据single_file_modeboolean否false每个并行度只输出一个文件。开启后batch_size不生效输出文件名不带文件块后缀create_empty_file_when_no_databoolean否false上游无数据同步时仍生成对应的空数据文件parquet_avro_write_timestamp_as_int96boolean否false仅 parquet 格式生效是否将 timestamp 写入 Parquet INT96parquet_avro_write_fixed_as_int96array否-仅 parquet 格式生效是否将 12 字节字段写入 Parquet INT96enable_header_writeboolean否false仅 text、csv 格式生效。false 不写表头true 写表头encodingstring否UTF-8仅 json、text、csv、xml 格式生效输出文件编码remote_userstring否-HDFS 远程用户名schema_evolution_enabledboolean否false为 CDC 管道启用 schema 演进true时 ADD/DROP/RENAME/MODIFY 列事件无需重启任务即可应用到 sink。binary 格式不支持schema_save_modestring否CREATE_SCHEMA_WHEN_NOT_EXIST目录已存在时的处理方式data_save_modestring否APPEND_DATA已有数据的处理方式multi_table_sink_replicaint否1多表 Sink 任务中每张表的 sink writer 副本数merge_update_eventboolean否false仅 canal_json、debezium_json、maxwell_json 格式生效。为true时将 UPDATE_AFTER 与 UPDATE_BEFORE 事件合并为 UPDATE 事件数据参数间的条件化校验从 HdfsFileSinkFactory.java 的optionRule()可以看到连接器使用conditional实现了参数间的依赖校验配置不符合条件会在任务启动时直接报错file_format_typetext时才接受row_delimiter、field_delimiter、TXT_COMPRESSlzo/none、enable_header_writefile_format_typecsv时接受row_delimiter、TXT_COMPRESS、enable_header_writefile_format_typejson时接受row_delimiter、TXT_COMPRESSfile_format_typeorc时接受ORC_COMPRESSlzo/snappy/lz4/zlib/nonefile_format_typeparquet时接受PARQUET_COMPRESSlzo/snappy/lz4/gzip/brotli/zstd/none及两个 INT96 选项file_format_typexml时接受xml_use_attr_format、xml_root_tag、xml_row_tagcustom_filenametrue时才接受file_name_expression、filename_time_formathave_partitiontrue时才接受partition_by、partition_dir_expression、is_partition_field_write_in_filefile_format_type为 text/json/csv/xml 时才接受encoding。这也解释了压缩格式矩阵不同文件格式的压缩选项在代码中以singleChoice枚举白名单约束见 FileBaseSinkOptions.java例如 ORC 不允许 gzip/brotli/zstdparquet 不允许 zlib。schema_save_mode 与 data_save_mode目录与数据生命周期管理这两个参数控制任务启动前对目标目录的预处理策略与文件连接器的 SaveMode 机制对应源码见 sink/config/SaveMode.java。schema_save_mode目录已存在时的处理RECREATE_SCHEMA目录不存在则创建目录存在则删除后重建CREATE_SCHEMA_WHEN_NOT_EXIST目录不存在则创建存在则跳过默认ERROR_WHEN_SCHEMA_NOT_EXIST目录不存在时报错IGNORE忽略对表的处理。data_save_mode已有数据的处理DROP_DATA保留目录删除已有数据文件APPEND_DATA保留目录保留已有数据文件默认ERROR_WHEN_DATA_EXISTS已有数据文件时报错。从代码约束看data_save_mode仅允许DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS三个取值FileBaseSinkOptions.java。典型使用场景落地 Hive 数仓表时建议schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXISTdata_save_modeDROP_DATA先清空再写保证当天数据无残留增量场景用APPEND_DATA。schema_evolution_enabledCDC 管道下的在线 Schema 演进当schema_evolution_enabledtrue时文件 Sink 可在运行期处理 CDC 的 Schema 变更事件ADD COLUMN、DROP COLUMN、RENAME COLUMN、MODIFY COLUMN无需重启任务。每次 schema 变更时当前输出文件会被关闭并以更新后的 schema 打开新文件写入。支持的格式除binary外的所有文件格式。若file_format_typebinary却开启该选项任务启动时会被配置校验拦截并报错。分区约束当have_partitiontrue时不允许 DROP 掉partition_by中列出的分区列否则会快速失败fail fast。分区列必须跨 schema 变更保持稳定。当schema_evolution_enabledfalse默认时如果上游 CDC 源开启了schema-changes.enabledtrue且AlterTableEvent到达 sink任务会立即抛出可操作的错误Received AlterTableEvent but schema_evolution_enabledfalse at this sink. Either set schema_evolution_enabledtrue to handle schema changes, or set schema-changes.enabledfalse at the CDC source to suppress them.使用默认 CDC 源配置schema-changes.enabledfalse的用户完全不受影响。已知限制schema 变更与 checkpoint 不是原子的。若任务在文件轮转与 schema 元数据更新之间的窗口内崩溃恢复后写入的行可能仍使用变更前的 schema。这是 SeaTunnel 其他 Sink 共有的架构性缺口完整的重启 DDL 恢复正确性需要配合后续 CDC 源修复另行跟踪。源码层面写入器通过SupportSchemaEvolutionSinkWriter接口暴露applySchemaChange(SchemaChangeEvent)见 BaseFileSinkWriter.java将事件转交给WriteStrategy完成文件轮转与 schema 更新。格式侧还有 FileSchemaEvolutionTest.java 与 ParquetWriteStrategyEvolutionTest.java 等测试覆盖该能力。CDC 管道中的示例HdfsFile { fs.defaultFS hdfs://hadoopcluster path /tmp/seatunnel/cdc/${table_name} file_format_type parquet schema_evolution_enabled true }multi_table_sink_replica 与 merge_update_eventmulti_table_sink_replica多表 Sink 任务中每张表使用的 sink writer 副本数默认1。仅当单张表需要更高的 writer 并行度时才需要调大。merge_update_event仅canal_json、debezium_json、maxwell_json三种 CDC 序列化格式生效。为true时UPDATE_AFTER 与 UPDATE_BEFORE 事件被合并为 UPDATE 事件数据为false时二者会作为独立事件分别序列化。对应源码定义见 FileBaseSinkOptions.java。csv_string_quote_modeCSV 字符串引号模式当文件格式为 CSV 时的字符串引号策略对应 CsvStringQuoteModeALL所有字符串字段都加引号MINIMAL默认仅当字段包含特殊字符如字段分隔符、引号字符或行分隔符字符串中的任意字符时才加引号NONE永不加引号。当数据中出现分隔符时printer 会用转义字符作为前缀若未设置转义字符格式校验会抛出异常。任务示例从简单到生产级最简配置FakeSource → HDFS以下任务通过 FakeSource 自动生成 16 行多种类型的测试数据写入 HDFS 的 ORC 文件# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { # This is a example source plugin **only for test and demonstrate the feature source plugin** FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { c_map mapstring, smallint c_array arrayint 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 { # If you would like to get more information about how to configure seatunnel and see full list of transform plugins, # please go to https://seatunnel.apache.org/docs/transforms } sink { HdfsFile { fs.defaultFS hdfs://hadoopcluster path /tmp/hive/warehouse/test2 file_format_type orc } # If you would like to get more information about how to configure seatunnel and see full list of sink plugins, # please go to https://seatunnel.apache.org/docs/connectors/sink }ORC 格式最简配置HdfsFile { fs.defaultFS hdfs://hadoopcluster path /tmp/hive/warehouse/test2 file_format_type orc }text 格式 分区 自定义文件名 列裁剪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 格式 分区 自定义文件名 列裁剪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 认证最简配置HdfsFile { fs.defaultFS hdfs://hadoopcluster path /tmp/hive/warehouse/test2 hdfs_site_path /path/to/your/hdfs_site_path kerberos_principal your_principalEXAMPLE.COM kerberos_keytab_path /path/to/your/keytab/file.keytab }压缩配置HdfsFile { fs.defaultFS hdfs://hadoopcluster path /tmp/hive/warehouse/test2 compress_codec lzo }深入原理事务提交与文件名生成两阶段提交tmp_path 与 rename连接器的 exactly-once 承诺源自临时目录 原子 mv机制链路如下数据先写入tmp_path下的事务目录默认/tmp/seatunnel。每个 checkpoint 对应一个事务事务目录结构由getTransactionDir(transactionId)计算事务 ID 格式为T_{jobId}_{uuidPrefix}_{subTaskIndex}_{checkpointId}见 AbstractWriteStrategy.javauuidPrefix是 10 位随机 UUID保证跨任务不冲突。每个 checkpoint 到达时writer 调用prepareCommit()关闭当前文件返回FileCommitInfo含needMoveFiles映射临时文件 → 目标文件见 AbstractWriteStrategy.java。checkpoint 完成时FileSinkAggregatedCommitter.java 遍历mvFileEntry调用hadoopFileSystemProxy.renameFile(临时路径, 目标路径, true)将临时文件原子 mv 到目标目录即文档所述先写 tmp再用 mv 提交到目标目录。若任务从 checkpoint 恢复writer 会找出未包含在恢复状态中的孤儿事务并调用abortPrepare(transaction)删除对应事务目录只由 aggregated committer 重放已 checkpoint 的事务BaseFileSinkWriter.java。这也解释了is_enable_transactiontrue时文件名自动带${transactionId}_前缀的原因不同事务、不同并行子任务之间的文件必须可区分才能在恢复时精确取舍。文件名生成规则generateFileName(transactionId)AbstractWriteStrategy.java按以下顺序拼接若配置了filename_extension直接使用该后缀自动补.前缀否则使用文件格式后缀 压缩编码后缀例如 parquet lzo 会得到类似.parquet.lzo的结果对file_name_expression执行变量替换${transactionId}→ 当前事务 ID、${now}→ 按filename_time_format格式化的当前时间、${uuid}→ 随机 UUID若未开启single_file_mode追加_ 子任务编号partId保证同一事务内不同并行子任务的文件不冲突追加文件后缀。single_file_mode 与 binary 自定义文件名的启动校验BaseFileSinkWriter.java 的preCheckConfig还包含两条启动期校验binary 自定义文件名 并行度大于 1文件名表达式必须包含${transactionId}或${uuid}否则报错——因为并行子任务共享同一文件名时会导致互相覆盖single_file_mode 并行度大于 1文件名表达式必须包含${transactionId}否则报错。这两条保证了多并行度下输出文件名的唯一性。ViewFS联邦 HDFS配置实战ViewFS 可将多个 HDFS 集群或 namespace 统一挂载为一个逻辑命名空间适用于 HDFS Federation 场景。连接器在 HadoopConf.java 中通过检测fs.defaultFS是否以viewfs://开头自动选择org.apache.hadoop.fs.viewfs.ViewFileSystem实现并切换fs.viewfs.impl与fs.defaultFS配置。HdfsFile { fs.defaultFS viewfs://mycluster path /data/output file_format_type parquet hdfs_site_path /path/to/core-site.xml data_save_mode DROP_DATA }在core-site.xml中配置挂载表?xml version1.0 encodingUTF-8? configuration property namefs.viewfs.mounttable.mycluster.link./data/name valuehdfs://namenode1:9000/data/value /property property namefs.viewfs.mounttable.mycluster.link./logs/name valuehdfs://namenode2:9000/logs/value /property property namefs.viewfs.mounttable.mycluster.link./tmp/name valuehdfs://namenode3:9000/tmp/value /property /configuration注意hdfs_site_path指向包含挂载表配置的core-site.xml或hdfs-site.xml。在 HadoopConf.setExtraOptionsForConfiguration 中该 XML 会被加载进 HadoopConfiguration且会先 unsetfs.defaultFS、fs.{schema}.impl、fs.{schema}.impl.disable.cache三个键避免外部 XML 覆盖连接器自身构造的核心配置。写入 HA HDFS 集群启用 Kerberos写入使用 Kerberos 的 HA HDFS 集群时除 nameservice URI 外还需提供 principal/keytab。连接器复用 Hadoop 工具链的同一套认证机制因此 principal 的 HDFS 权限必须允许写入目标目录。sink { HdfsFile { fs.defaultFS hdfs://mycluster path /data/landing/events file_format_type parquet hdfs_site_path /etc/hadoop/conf/hdfs-site.xml kerberos_principal sinkEXAMPLE.COM krb5_path /etc/krb5.conf } }关键说明kerberos_principal与krb5_path会被转发给 Hadoop FileSystem 客户端见 HdfsFileSinkFactory.initHadoopConf连接器本身不执行kinit因此 keytab 必须已能在每个 worker 节点被发现通常通过KRB5CCNAME或kinitcron 实现或通过标准 Hadoop 认证工具提供给同一 JVM。HA 场景下hdfs_site_path用于加载 NameNode HA 配置nameservice → active/standby 映射。若出现集群级认证问题检查 worker 日志中的LoginException/KrbException信息——这些表明凭据问题而非连接器缺陷。remote_user参数可显式指定 HDFS 远端用户名。源码级测试验证仓库为 HdfsFile 连接器提供了较为完整的测试覆盖可作为配置正确性的参考HdfsFileSinkTest.javaSink 工厂与配置解析测试HdfsFileFactoryTest.java插件工厂行为测试通用写入策略测试text/csv/orc/parquet 等集中在 connector-file-base 的 writer 测试目录例如 CsvWriteStrategyTest.java、OrcWriteStrategyTest.java、ParquetWriteStrategyTest.java两阶段提交行为见 FileSinkAggregatedCommitterTest.javaKerberos 相关工具类见 HadoopFileSystemProxyKerberosRenewTest.java。小结与选型建议HdfsFile Sink 是 SeaTunnel 面向 HDFS 数仓落地的核心出口通过 tmp_path 事务目录与 rename 提交实现 exactly-once通过 WriteStrategy 抽象统一八种文件格式与压缩编码通过 SaveMode 管理目录与数据生命周期并针对联邦 HDFS、HA Kerberos、CDC schema 演进等生产场景提供了完整参数。上手时建议遵循以下顺序验证先跑通最简 ORC 配置确认连通性 → 按需开启分区与自定义文件名 → 接入生产集群时再启用hdfs_site_path/ Kerberos 认证与合适的schema_save_mode/data_save_mode组合避免数据残留或误删。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考