ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

SeaTunnel CosFile Sink 详解:将数据落地到腾讯云 COS 的文件格式、精确一次提交与 Schema 演变

SeaTunnel CosFile Sink 详解:将数据落地到腾讯云 COS 的文件格式、精确一次提交与 Schema 演变 SeaTunnel CosFile Sink 详解将数据落地到腾讯云 COS 的文件格式、精确一次提交与 Schema 演变【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的CosFile连接器是一个文件类 Sink负责把结构化数据写入腾讯云 COSCloud Object Storage桶支持 text、csv、parquet、orc、json、excel、xml、binary 及 CDC JSON 等多种文件格式并通过 2PC 提交机制保证数据不丢不重。读完本文你可以掌握 CosFile 全部配置参数的含义与默认值、各引擎下的部署依赖要求、分区与自定义文件名的配置方式以及从源码层面理解其 Hadoop 文件系统适配和精确一次exact-once提交流程的实现原理。引擎支持与部署依赖CosFile支持以下引擎SparkFlinkSeaTunnel Zeta由于它基于 Hadoop 文件系统抽象来访问 COS各引擎下的依赖要求不同Spark / Flink必须确保集群已集成 Hadoop官方测试过的 Hadoop 版本为 2.x。SeaTunnel EngineZeta安装时会自动集成 Hadoop jar可在${SEATUNNEL_HOME}/lib下确认。此外还必须将hadoop-cos-{hadoop.version}-{version}.jar和cos_api-bundle-{version}.jar放入${SEATUNNEL_HOME}/lib目录。注意版本约束仅支持 Hadoop 2.6.5 与 hadoop-cos 8.0.2jar 包可从 hadoop-cos 的 release 页获取。从源码结构看CosFile的插件标识符由 FileSystemType 枚举定义COS(CosFile)并且 plugin-mapping.properties 中同时注册了 source 与 sink 两个方向seatunnel.source.CosFile connector-file-cos seatunnel.sink.CosFile connector-file-cos也就是说同一模块 connector-file-cos 既提供CosFileSink也提供CosFileSource本文聚焦 Sink 方向。关键特性多模态Multimodal使用二进制文件格式读取和写入任意格式的文件例如视频、图片等任何文件都可以同步到目标位置。精确一次Exactly-Once默认使用 2PC commit 确保精确一次。文件格式类型text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json。配置选项总览名称类型必需默认值描述pathstring是-Sink 写入 COS 桶内的目标目录。配合bucket实际路径为cosn://bucketpath。tmp_pathstring否/tmp/seatunnel结果文件将首先写入 tmp 路径然后使用 “mv” 将 tmp 目录提交到目标目录。需要一个 COS 目录。bucketstring是-COS 文件系统的桶地址例如cosn://seatunnel-test-1259587829。secret_idstring是-腾讯云 COS 的 SecretId。secret_keystring是-腾讯云 COS 的 SecretKey。regionstring是-COS 桶所在地域例如ap-chengdu。custom_filenameboolean否false是否需要自定义文件名。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、canal_json、debezium_json、maxwell_json。filename_extensionstring否-使用自定义的文件扩展名覆盖默认的文件扩展名例如.xml、.json、dat、.customtype。field_delimiterstring否\001仅在 file_format 为 text 时使用。row_delimiterstring否\n仅在 file_format 为text、csv、json时使用。have_partitionboolean否false是否需要处理分区。partition_byarray否-只有在 have_partition 为 true 时才使用。partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/只有在 have_partition 为 true 时才使用。is_partition_field_write_in_fileboolean否false只有在 have_partition 为 true 时才使用。sink_columnsarray否-当此参数为空时所有字段都是接收列。is_enable_transactionboolean否true若为true写入目标目录的数据不会丢失或重复当为true时会自动在文件名前缀添加${transactionId}_。batch_sizeint否1000000单个文件的最大行数。对于 SeaTunnel Engine文件中的行数由batch_size和checkpoint.interval共同决定。compress_codecstring否none文件的压缩编解码器。Excel 格式不支持任何压缩格式。common-optionsobject否-Sink 插件通用参数请参考 Sink Common Options 了解详情。max_rows_in_memoryint否-仅在 file_format 为 excel 时使用。sheet_max_rowsint否1048576仅在file_format_type为excel时使用每个工作表允许写入的最大行数。sheet_namestring否Sheet${Random number}仅在 file_format 为 excel 时使用。csv_string_quote_modeenum否MINIMAL仅在 file_format 为 csv 时使用。xml_root_tagstring否RECORDS仅在 file_format 为 xml 时使用。xml_row_tagstring否RECORD仅在 file_format 为 xml 时使用。xml_use_attr_formatboolean否-仅在 file_format 为 xml 时使用。single_file_modeboolean否false每个并行处理只会输出一个文件。启用此参数后batch_size 将不会生效。输出文件名没有文件块后缀。create_empty_file_when_no_databoolean否false当上游没有数据同步时仍然会生成相应的数据文件。parquet_avro_write_timestamp_as_int96boolean否false仅在 file_format 为 parquet 时使用。parquet_avro_write_fixed_as_int96array否-仅在 file_format 为 parquet 时使用。encodingstring否UTF-8仅当 file_format_type 为 json、text、csv、xml 时使用。merge_update_eventboolean否false仅当 file_format_type 为 canal_json、debezium_json、maxwell_json 时使用。设置为true时会将UPDATE_AFTER与UPDATE_BEFORE合并为UPDATE事件数据。schema_evolution_enabledboolean否false开启 Schema 演变支持适用于 CDC 管道。为 true 时来自上游的 ADD/DROP/RENAME/MODIFY 列事件无需重启作业即可应用到 Sink。不支持 binary 格式。path [string]目标目录路径必需。实际写入位置为cosn://bucket 去掉 cosn:// 前缀后的桶名path。bucket [string]COS 文件系统的 bucket 地址例如cosn://seatunnel-test-1259587829。注意必须带cosn://schema 前缀。secret_id / secret_key / regionsecret_id是腾讯云 COS 的密钥 IDsecret_key是对应密钥region是桶所在地域例如ap-chengdu。三者均会注入到 Hadoop 的 COSN 文件系统配置中见下文“Hadoop 文件系统适配”一节。custom_filename / file_name_expression / filename_time_formatcustom_filename为true时启用自定义文件名。file_name_expression描述将在path中创建的文件表达式可包含变量${now}或${uuid}例如test_${uuid}_${now}其中${now}表示当前时间其格式由filename_time_format定义默认yyyy.MM.dd。常用时间格式符号描述y年M月d日H时 (0-23)m分s秒请注意如果is_enable_transaction为true会自动在文件名开头添加${transactionId}_前缀。file_format_type [string]支持的文件类型textcsvparquetorcjsonexcelxmlbinarycanal_jsondebezium_jsonmaxwell_json。请注意最终文件名将以 file_format 的后缀结尾文本文件的后缀为txt。filename_extension可用来自定义覆盖该后缀。field_delimiter / row_delimiterfield_delimiter是数据行中列之间的分隔符仅text格式需要row_delimiter是行之间的分隔符仅text、csv、json格式需要。分区相关have_partition / partition_by / partition_dir_expression / is_partition_field_write_in_filehave_partition为true时启用分区处理。partition_by指定基于哪些字段分区。partition_dir_expression指定分区目录模板。默认是${k0}${v0}/${k1}${v1}/.../${kn}${vn}/其中k0是第一个分区字段v0是第一个分区字段的值指定partition_by后会据此生成分区目录最终文件放置在分区目录中。is_partition_field_write_in_file为true时分区字段及其值也会写入数据文件。如果你想写 Hive 数据文件该值应为falseHive 表从目录结构识别分区。sink_columns [array]指定哪些列需要写入文件默认是从Transform或Source获取的所有列字段的顺序决定文件实际写入的顺序。is_enable_transaction [boolean]为true时确保数据写入目标目录时不丢失、不重复并自动在文件名开头添加${transactionId}_。当前只支持true。其底层提交机制见后文“精确一次提交”一节。batch_size [int]单个文件中的最大行数。对于 SeaTunnel 引擎文件中的行数由batch_size和checkpoint.interval共同决定如果checkpoint.interval足够大写入程序会在文件中一直写行直到超过batch_size才轮转如果checkpoint.interval较小则每次新检查点触发时都会创建新文件。compress_codec [string]压缩编解码器按格式支持情况如下txt:lzononejson:lzononecsv:lzononeorc:lzosnappylz4zlibnoneparquet:lzosnappylz4gzipbrotlizstdnoneExcel 类型不支持任何压缩格式。common options接收器写入插件常用参数如max_fail_times、max_fail_retry_time、output_data_format等请参考 Sink Common Options 了解详情。Excel 专属max_rows_in_memory / sheet_max_rows / sheet_namemax_rows_in_memory内存中可以缓存的最大数据项数仅 excel 格式。sheet_max_rows每个工作表可写入的最大行数默认1048576。sheet_name写入的工作表名默认Sheet${随机数}。CSV 专属csv_string_quote_modeCSV 的字符串引用模式ALL所有字符串字段都将被引用。MINIMAL仅当字段包含特殊字符如字段分隔符、引号字符或行分隔符中的任意字符时才加引号。NONE从不引用字段。当分隔符出现在数据中时会以转义符作为前缀若未设置转义符格式校验将抛出异常。XML 专属xml_root_tag / xml_row_tag / xml_use_attr_formatxml_root_tag指定 XML 文件中根元素的标记名默认RECORDS。xml_row_tag指定数据行的标记名称默认RECORD。xml_use_attr_format指定是否使用标记属性格式处理数据。Parquet 专属parquet_avro_write_timestamp_as_int96 / parquet_avro_write_fixed_as_int96parquet_avro_write_timestamp_as_int96支持将时间戳写入 Parquet INT96。parquet_avro_write_fixed_as_int96支持将 12 字节字段写入 Parquet INT96。encoding [string]仅当file_format_type为 json、text、csv、xml 时使用。指定要写入文件的编码该参数由Charset.forName(encoding)解析默认UTF-8。merge_update_event [boolean]仅当file_format_type为canal_json、debezium_json、maxwell_json时使用。设置为true时序列化数据中UPDATE_AFTER与UPDATE_BEFORE会合并为UPDATE事件设置为false时不合并。single_file_mode [boolean]每个并行处理只输出一个文件启用后batch_size不再生效输出文件名没有文件块后缀。create_empty_file_when_no_data [boolean]当上游没有数据同步时仍然生成相应的数据文件默认false。配置示例文本文件格式分区 自定义文件名 列裁剪CosFile { path /sink bucket cosn://seatunnel-test-1259587829 secret_id xxxxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxxxx region ap-chengdu 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 格式分区 列裁剪CosFile { path /sink bucket cosn://seatunnel-test-1259587829 secret_id xxxxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxxxx region ap-chengdu have_partition true partition_by [age] partition_dir_expression ${k0}${v0} is_partition_field_write_in_file true file_format_type parquet sink_columns [name,age] }ORC 格式最小配置CosFile { path /sink bucket cosn://seatunnel-test-1259587829 secret_id xxxxxxxxxxxxxxxxxxx secret_key xxxxxxxxxxxxxxxxxxx region ap-chengdu file_format_type orc }CDC 管道中的 Schema 演变示例LocalFile { path /tmp/cdc/${table_name} file_format_type parquet schema_evolution_enabled true have_partition true partition_by [updated_at_month] }该示例展示了文件 Sink 在 CDC 管道中开启schema_evolution_enabled的通用写法CosFile 作为同族文件 Sink 同样适用只需替换为CosFile并补齐 bucket / secret_id / secret_key / region。schema_evolution_enabledCDC 管道的 Schema 演变设置为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 变更过程中必须保持稳定。当schema_evolution_enabled false默认值时若上游 CDC Source 配置了schema-changes.enabled true且 Sink 收到AlterTableEvent作业会立即抛出如下错误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 Source 配置schema-changes.enabled false的用户不受影响。已知限制Schema 变更与 Checkpoint 不是原子操作。若作业在文件轮转与 Schema 元数据更新之间的窗口期崩溃恢复后写入的数据行可能使用变更前的 Schema。这是与其他 SeaTunnel Sink 共同存在的已知架构限制。源码层面Schema 事件的入口是 BaseFileSinkWriter 实现的SupportSchemaEvolutionSinkWriter.applySchemaChange(SchemaChangeEvent)方法它把事件直接委托给WriteStrategy处理文件关闭与重新打开逻辑CosFile 通过继承BaseFileSink免费获得该能力。源码解析Hadoop 文件系统适配CosFile 连接器本身并不直接调用腾讯云 SDK而是复用 SeaTunnel 的 Hadoop 文件系统抽象。核心链路如下CosFileSink继承BaseFileSink仅做两件事通过CosConf.buildWithReadonlyConfig(pluginConfig)构建 Hadoop 配置并返回插件名CosFile来自FileSystemType.COS。CosConf继承HadoopConf定义了 COS 的 HDFS 实现类与 schemaprivate static final String HDFS_IMPL org.apache.hadoop.fs.CosFileSystem; private static final String SCHEMA cosn;并把用户配置的secret_id、secret_key、region三个必填项映射为hadoop-cos的CosNConfigKeyscosOptions.put(CosNConfigKeys.COSN_USERINFO_SECRET_ID_KEY, readonlyConfig.get(CosFileBaseOptions.SECRET_ID)); cosOptions.put(CosNConfigKeys.COSN_USERINFO_SECRET_KEY_KEY, readonlyConfig.get(CosFileBaseOptions.SECRET_KEY)); cosOptions.put(CosNConfigKeys.COSN_REGION_KEY, readonlyConfig.get(CosFileBaseOptions.REGION)); hadoopConf.setExtraOptions(cosOptions);这解释了为什么bucket必须写成cosn://形式——Hadoop 根据cosnscheme 路由到org.apache.hadoop.fs.CosFileSystem。选项声明在 CosFileBaseOptions 中它继承FileBaseSourceOptions新增secret_id、secret_key、region、bucket四个无默认值的必填选项其余文件写入相关选项batch_size、分区、自定义文件名、压缩等都复用FileBaseSinkOptions公共定义这也是为什么 CosFile、LocalFile、OssFile 等文件 Sink 的选项表高度相似。CosFileSinkFactory通过AutoService(Factory.class)自动注册其optionRule()明确了配置校验逻辑path、bucket、secret_id、secret_key、region五项必填field_delimiter、row_delimiter、压缩参数、XML 参数、encoding等选项均按file_format_type的取值条件生效file_name_expression、filename_time_format在custom_filenametrue时条件生效partition_by等在have_partitiontrue时条件生效。配置缺项或不匹配会在作业启动阶段被拦截。源码解析精确一次提交如何实现is_enable_transaction默认true对应实现是 BaseFileSinkWriter 围绕WriteStrategy构建的“tmp 目录写入 事务提交/中止”流程写入阶段每个 Writer 构造时若从 checkpoint 状态恢复fileSinkStates非空会扫描tmp_path下该作业的前缀目录findTransactionList对不在恢复状态中的未提交事务执行writeStrategy.abortPrepare(transaction)清理再从checkpointId 1开启新事务全新作业则从事务1开始。这正是tmp_path选项的意义——结果文件先写到 tmp 路径再以 “mv” 提交到目标目录。提交阶段prepareCommit()将当前事务的数据文件FileCommitInfo上抛给聚合 Committer在 checkpoint 完成后由 Committer 执行重命名/移动到path指定的目标目录文件名前缀${transactionId}_由此而来。状态快照snapshotState(checkpointId)把当前transactionId与uuidPrefix存入FileSinkState故障恢复时依据它决定哪些事务可安全中止、哪些必须续写从而保证不丢不重。多表与 Schema 演变BaseFileSinkWriter同时实现SupportMultiTableSinkWriterCDC 多表场景共享一个 Writer 实例和SupportSchemaEvolutionSinkWriter前述applySchemaChange入口。启动前校验preCheckConfig会对高风险组合快速失败例如 binary 格式 自定义文件名 多并发但表达式不含${transactionId}/${uuid}或single_file_mode在多并发下文件名表达式不含${transactionId}等情况直接抛出IllegalArgumentException。对使用方而言实践建议是保持is_enable_transaction true的默认值并确认tmp_path默认/tmp/seatunnel在目标桶中可写若作业频繁 checkpoint可通过调大checkpoint.interval配合batch_size控制最终文件大小。变更日志该连接器的变更日志文件 connector-file-cos.md 目前为空说明当前仓库记录中暂无需要特别提示的破坏性变更记录。小结CosFile 连接器把 SeaTunnel 统一的文件写入框架格式序列化、分区目录、2PC 提交、Schema 演变与 Hadoop 的CosNFileSystem结合用户只需提供bucket、secret_id、secret_key、region、path五个必填项即可将 text/csv/parquet/orc/json/excel/xml/binary 及三种 CDC JSON 格式的数据可靠地落地到腾讯云 COS。理解tmp_pathtransactionId的提交模型后你可以从容调整batch_size、checkpoint.interval、custom_filename与分区表达式来匹配自身的存储组织与下游如 Hive的读取约定。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表