十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel LocalFile Sink 连接器完全指南:本地文件写出、事务保障与分区配置实战

SeaTunnel LocalFile Sink 连接器完全指南:本地文件写出、事务保障与分区配置实战 SeaTunnel LocalFile Sink 连接器完全指南本地文件写出、事务保障与分区配置实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelLocalFile本地文件Sink 是 SeaTunnel 中将任意上游数据写入本地文件系统的标准连接器支持 text、csv、parquet、orc、json、excel、xml、binary 及三大 CDC JSON 格式canal_json、debezium_json、maxwell_json共 11 种文件格式并内置基于 2PC 的 exactly-once 语义保障。本文以官方文档 docs/en/connectors/sink/LocalFile.md 为主体结合仓库内连接器源码与单元测试系统讲解其全部配置参数、事务与临时目录机制、分区写出、文件命名策略及面向 CDC 管线的 schema evolution 能力帮助读者在 Spark、Flink 与 SeaTunnel Zeta 三种引擎上正确、完整地配置本地文件落盘任务。概述一个连接器三种引擎全格式本地落盘LocalFile Sink 用于将 SeaTunnel 管线中的数据输出到本地文件系统。它支持以下引擎SparkFlinkSeaTunnel Zeta使用前需注意运行环境对 Hadoop 的依赖若使用Spark / Flink必须确保集群已集成 Hadoop官方测试过的 Hadoop 版本为 2.x若使用SeaTunnel EngineZeta安装包已自动集成 Hadoop jar可检查${SEATUNNEL_HOME}/lib目录下的 jar 包加以确认。从源码结构看LocalFile 连接器的实现位于 seatunnel-connectors-v2/connector-file/connector-file-local其核心类 LocalFileSink 继承自文件连接器基座模块connector-file-base的BaseMultipleTableFileSink因此天然获得多表写入、事务提交、分区写出等通用文件能力连接器自身仅需提供LocalFileHadoopConf与插件名FileSystemType.LOCAL这种“基座复用 薄实现”的设计在全部文件类连接器Local、HDFS、S3、OSS、FTP 等中保持一致。核心特性特性支持情况说明multimodal多模态✅使用二进制文件格式读写任意格式文件视频、图片等任何文件都可被同步到目标位置exactly-once✅默认通过 2PC 提交保证数据不丢不重cdc❌本身不是 CDC 源但可作为 CDC 管线的落盘端多表写入✅支持一个 sink 同时写入多张表文件格式✅text / csv / parquet / orc / json / excel / xml / binary / canal_json / debezium_json / maxwell_jsontimer flush❌不支持定时刷盘关于multimodal与exactly-once等概念的详细定义可参考 connector-v2-features.md。参数总览下表汇总了 LocalFile Sink 的全部配置项其中“Required”标记为 yes 的仅有path一项名称类型必填默认值说明pathstringyes-Sink 写入的目标目录可通过${database_name}、${table_name}、${schema_name}注入上游 CatalogTable 信息tmp_pathstringno/tmp/seatunnel结果文件先写入临时目录再通过mv将临时目录提交到目标目录custom_filenamebooleannofalse是否需要自定义文件名file_name_expressionstringno${transactionId}仅当custom_filename为 true 时生效filename_time_formatstringnoyyyy.MM.dd仅当custom_filename为 true 时生效file_format_typestringnocsv文件格式支持 text/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_jsonfilename_extensionstringno-覆盖默认文件扩展名例如.xml、.json、dat、.customtypefield_delimiterstringnotext 为 \001csv 为 ,仅 text 与 csv 格式生效row_delimiterstringno\n仅 text、csv、json 格式生效have_partitionbooleannofalse是否进行分区处理partition_byarrayno-仅当have_partition为 true 时生效partition_dir_expressionstringno${k0}${v0}/${k1}${v1}/.../${kn}${vn}/仅当have_partition为 true 时生效is_partition_field_write_in_filebooleannofalse仅当have_partition为 true 时生效sink_columnsarrayno空全部字段需要写入文件的列字段顺序决定文件实际写入顺序is_enable_transactionbooleannotrue为 true 时保证写入目标目录的数据不丢不重并自动在文件名前追加${transactionId}_前缀batch_sizeintno1000000单个文件的最大行数SeaTunnel Engine 下文件行数由batch_size与checkpoint.interval共同决定compress_codecstringnonone文件压缩编码excel 不支持任何压缩common-optionsobjectno-Sink 插件通用参数详见 sink-common-options.mdmax_rows_in_memoryintno-仅 excel 格式生效内存中缓存的最大数据条数sheet_max_rowsintno1048576仅 excel 格式生效每个 sheet 的最大行数sheet_namestringnoSheet${随机数}仅 excel 格式生效工作表名称csv_string_quote_modeenumnoMINIMAL仅 csv 格式生效字符串引用模式xml_root_tagstringnoRECORDS仅 xml 格式生效根元素标签名xml_row_tagstringnoRECORD仅 xml 格式生效数据行标签名xml_use_attr_formatbooleanno-仅 xml 格式生效是否使用标签属性格式处理数据single_file_modebooleannofalse每个并行度只输出一个文件开启后batch_size不生效输出文件名不带文件块后缀create_empty_file_when_no_databooleannofalse上游无数据时仍生成对应数据文件parquet_avro_write_timestamp_as_int96booleannofalse仅 parquet 格式生效将 timestamp 写为 Parquet INT96parquet_avro_write_fixed_as_int96arrayno-仅 parquet 格式生效将 12 字节字段写为 Parquet INT96enable_header_writebooleannofalse仅 text、csv 格式生效是否写表头encodingstringnoUTF-8仅 json、text、csv、xml 格式生效通过Charset.forName(encoding)解析schema_save_modestringnoCREATE_SCHEMA_WHEN_NOT_EXIST同步任务开始前对目标目录的处理方式data_save_modestringnoAPPEND_DATA同步任务开始前对目标目录中已有数据文件的处理方式merge_update_eventbooleannofalse仅 canal_json/debezium_json/maxwell_json 生效为 true 时将 UPDATE_AFTER 与 UPDATE_BEFORE 事件合并为 UPDATE 事件schema_evolution_enabledbooleannofalse为 CDC 管线启用 schema evolutionADD/DROP/RENAME/MODIFY 列事件无需重启作业即可应用到 sinkbinary 格式不支持这些参数在源码中的定义与默认值均可在 FileBaseSinkOptions.java 中逐一核对例如DEFAULT_TMP_PATH /tmp/seatunnel、DEFAULT_FILE_NAME_EXPRESSION ${transactionId}、DEFAULT_BATCH_SIZE 1000000等常量与上文表格完全一致。关键参数深入解读path目标目录与元数据注入path是唯一必填参数。除静态目录外可将上游 CatalogTable 的元数据注入路径占位符包括${database_name}、${table_name}、${schema_name}。典型用法LocalFile { path /tmp/hive/warehouse/${table_name} file_format_type parquet sink_columns [name,age] }这一能力对应源码变更[Feature][Core] Support using upstream table placeholders in sink options and auto replacement2.3.6 版本引入使得多表写入场景下每张表可以自动落入独立目录。tmp_path 与事务提交机制tmp_path默认/tmp/seatunnel。结果文件先写入临时路径任务提交时再通过mv把临时目录整体移动到目标目录这是 LocalFile Sink 实现 exactly-once 的关键。底层流程对应 BaseFileSink 中的createAggregatedCommitter()返回FileSinkAggregatedCommitter各并行 writer 先各自写临时文件聚合提交阶段再统一将临时目录 move 到最终路径从而避免部分写入导致的脏数据。custom_filename 与文件名表达式custom_filename true后可用file_name_expression自定义文件名支持变量${now}当前时间格式由filename_time_format指定默认yyyy.MM.dd${uuid}随机 UUID${transactionId}事务 ID。例如test_${uuid}_${now}。注意若is_enable_transaction true系统会自动在文件名最前面追加${transactionId}_前缀。针对 binary 格式有专门的约束当file_format_type binary且custom_filename true时适用于每个 sink 任务重命名单个源文件的场景若同一 sink 任务写入多个源文件应保持custom_filename false以保留源文件相对路径当存在并行 sink 子任务时必须在file_name_expression中包含${transactionId}或${uuid}避免多个子任务写出同名文件相互覆盖。filename_time_format支持的常用时间符号符号含义y年M月d日H小时0-23m分钟s秒file_format_type 与扩展名支持text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json。最终文件名以格式后缀结尾其中 text 格式的后缀是txt。可通过filename_extension覆盖默认扩展名例如写.dat、.customtype。field_delimiter 与 row_delimiterfield_delimiter一行数据中列与列之间的分隔符仅 text 和 csv 生效text 默认\001csv 默认,row_delimiter文件中行与行之间的分隔符仅 text、csv、json 生效默认\n。have_partition 分区写出have_partition true时启用分区配套参数partition_by按选定字段分区partition_dir_expression默认${k0}${v0}/${k1}${v1}/.../${kn}${vn}/k0为第一个分区字段、v0为其值最终文件落入生成的分区目录is_partition_field_write_in_file为 true 时分区字段及其值也写入数据文件若要写出 Hive 数据文件该值应设为false。sink_columns 列裁剪与排序sink_columns决定写入文件的列默认取Transform或Source输出的全部列字段排列顺序即文件实际写入顺序可用于控制列序与裁剪。is_enable_transaction 事务开关is_enable_transaction true默认时保证写入目标目录的数据不丢不重同时自动为文件名添加${transactionId}_前缀。文档注明当前仅支持true。batch_size 与 checkpoint 联动batch_size默认 1000000为单文件最大行数。在 SeaTunnel Engine 中文件行数由batch_size与checkpoint.interval共同决定若checkpoint.interval足够大writer 会持续写同一文件直到行数超过batch_size才切换新文件若checkpoint.interval较小则每次 checkpoint 触发都会新建文件。compress_codec 压缩支持矩阵格式支持的压缩编码txtlzo、nonejsonlzo、nonecsvlzo、noneorclzo、snappy、lz4、zlib、noneparquetlzo、snappy、lz4、gzip、brotli、zstd、noneexcel 类型不支持任何压缩。源码层面每种格式的合法压缩集通过独立的 Option 声明例如PARQUET_COMPRESS的单选项枚举为 NONE/LZO/SNAPPY/LZ4/GZIP/BROTLI/ZSTD与上表一致非法组合会在配置校验阶段直接报错。Excel 专属参数max_rows_in_memory内存中缓存的最大数据条数sheet_max_rows每个 sheet 的最大行数默认 1048576sheet_name写入的工作表名称默认Sheet${随机数}。csv_string_quote_mode 字符串引用模式ALL所有字符串字段均加引号MINIMAL仅对包含特殊字符如字段分隔符、引号或行分隔符中的字符的字段加引号NONE从不加引号数据中出现分隔符时以转义字符作为前缀若未设置转义字符格式校验将抛出异常。XML 专属参数xml_root_tagXML 根元素标签名默认RECORDSxml_row_tag数据行标签名默认RECORDxml_use_attr_format是否使用标签属性格式处理数据。Parquet INT96 参数parquet_avro_write_timestamp_as_int96将 timestamp 写为 Parquet INT96parquet_avro_write_fixed_as_int96将 12 字节字段写为 Parquet INT96字段列表。enable_header_write 与 encodingenable_header_write仅 text、csv 生效false 不写表头、true 写表头encoding仅 json、text、csv、xml 生效默认 UTF-8解析方式为Charset.forName(encoding)可写gbk等。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目录中存在数据文件时报错。merge_update_event仅 canal_json、debezium_json、maxwell_json 生效。为 true 时UPDATE_AFTER与UPDATE_BEFORE事件被合并为一条UPDATE事件数据为 false 时两者分别输出。配置示例ORC 文件格式最简配置LocalFile { path /tmp/hive/warehouse/test2 file_format_type orc }json / text / csv / xml 格式搭配 encodingLocalFile { path /tmp/hive/warehouse/test2 file_format_type text encoding gbk }parquet 格式搭配 sink_columns 列裁剪LocalFile { path /tmp/hive/warehouse/test2 file_format_type parquet sink_columns [name,age] }text 格式分区 自定义文件名 列裁剪完整组合LocalFile { 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 }excel 格式工作表名 内存行数 重建目录LocalFile { path/tmp/seatunnel/excel sheet_name Sheet1 max_rows_in_memory 1024 partition_dir_expression${k0}${v0} is_partition_field_write_in_filetrue file_name_expression${transactionId}_${now} file_format_typeexcel filename_time_formatyyyy.MM.dd is_enable_transactiontrue schema_save_modeRECREATE_SCHEMA data_save_modeDROP_DATA }从上游提取元数据注入路径LocalFile { path /tmp/hive/warehouse/${table_name} file_format_type parquet sink_columns [name,age] }schema_evolution_enabledCDC 管线的在线 Schema 演进schema_evolution_enabled默认 false用于 CDC 场景下的在线表结构变更处理为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 变更保持稳定关闭时默认若上游 CDC 源开启了schema-changes.enabled true且有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.enabled false的用户完全不受影响已知限制schema 变更与 checkpoint 并非原子操作。若作业在“文件轮转与 schema 元数据更新”之间的窄窗口内崩溃恢复后写入的行可能仍使用变更前的 schema。这是 SeaTunnel 其他 sink 共有的架构性缺口如需“重启后 DDL 完全正确”还需配套的 CDC 源修复单独跟踪。CDC 管线中的示例用法LocalFile { path /tmp/cdc/${table_name} file_format_type parquet schema_evolution_enabled true have_partition true partition_by [updated_at_month] }从源码看实现与校验规则插件注册与条件参数规则LocalFile Sink 通过 LocalFileSinkFactory 注册AutoService(Factory.class)插件名为FileSystemType.LOCAL对应的本地文件插件名。其optionRule()使用 SeaTunnel 的OptionRule声明了严谨的条件参数依赖必填仅pathfile_format_type text时开放row_delimiter、field_delimiter、TXT_COMPRESS、enable_header_writefile_format_type csv时开放row_delimiter、压缩与表头json 开放row_delimiter与压缩orc 仅开放 orc 压缩集parquet 开放压缩集与两个 INT96 参数xml 开放xml_use_attr_format、xml_root_tag、xml_row_tagcustom_filename true时才开放file_name_expression与filename_time_formathave_partition true时才开放partition_by、partition_dir_expression、is_partition_field_write_in_filefile_format_type为 text/json/csv/xml 时才开放encoding。这套条件规则确保了非法参数组合在配置解析阶段即被拦截而不是等到运行时才报错。2PC 提交与状态恢复BaseFileSink 实现了 SeaTunnel 的SeaTunnelSink四段式接口Writer → CommitInfo → AggregatedCommitter → 状态序列化createWriter()/restoreWriter()创建或恢复BaseFileSinkWriter恢复时传入历史FileSinkState支撑故障后不丢不重createAggregatedCommitter()返回FileSinkAggregatedCommitter在聚合阶段统一完成临时目录到目标目录的 move 提交这正是文档所述“先写 tmp_path、再 mv 到 path”的落地实现三类状态FileSinkState、FileCommitInfo、FileAggregatedCommitInfo均通过DefaultSerializer序列化随 checkpoint 持久化。此外preCheckConfig()还在作业启动阶段做了两项守护校验开启 checkpoint/流式模式时禁止single_file_mode启用create_empty_file_when_no_data时禁止同时启用分区。这两条限制在单元测试中都有对应的断言。单元测试验证的行为LocalFileTest.java 覆盖了以下关键行为可作为配置行为的佐证testSingleFileMode验证single_file_mode true时单文件仅输出一个文件并发度 1 且文件名表达式中无${transactionId}时抛出异常“Single file mode is not supported when file_name_expression not contains ${transactionId} but has parallel subtasks.”关闭 single_file_mode 后同名文件自动追加_0、_1等文件块后缀testCreateEmptyFileWhenNoData验证 text/csv/parquet 在无数据时生成空文件csv 开启enable_header_write时仅写表头而 binary 格式抛出“BinaryWriteStrategy does not support generating empty files when no data is written”testWriteFileWithCustomFileExtension验证filename_extension传txt2、.ppp均可生效自动规范化点号且 Source 端用同样扩展名可读回testCanalJsonSink/testDebeziumJsonSink/testMaxWellJsonSink验证三大 CDC JSON 格式的输出结构与merge_update_event true时 UPDATE_BEFORE/UPDATE_AFTER 被合并为单条 UPDATE 的行为。版本演进脉络从 connector-file-local.md 的变更日志可看到 LocalFile Sink 的能力演进主线2.2.0-beta 系列逐步加入 json、orc、parquet 格式与压缩支持2.3.x 系列持续增强包括 2.3.4 的多表写入Multiple Table、2.3.6 的 save mode 与上游表占位符、2.3.10 的filename_extension、单文件模式与空文件生成、2.3.11 的row_delimiter、2.3.12 的 canal_json/debezium_json/maxwell_json 三大 CDC 格式与自定义行分隔符。截至当前仓库dev 分支又新增了 markdown 解析器#9714等能力持续扩充格式边界。结语LocalFile Sink 以极简的必填参数仅path提供了覆盖 11 种格式、分区、压缩、事务、多表写入与 schema evolution 的完整文件落盘能力是 SeaTunnel 文件类连接器体系中最基础也最常用的一环。配置时建议遵循以下要点先按file_format_type选定格式再按需叠加compress_codec、encoding、sink_columns等格式专属参数需要精确控制目录布局时使用have_partitionpartition_dir_expression需要自定义文件名时开启custom_filename并注意并行子任务场景下必须携带${transactionId}/${uuid}运行 Spark/Flink 时务必确认集群已集成 Hadoop 2.x运行 SeaTunnel Engine 时则无需额外处理。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表