尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel OssJindoFile 连接器实战指南:通过 Jindo SDK 对接阿里云 OSS 的 Source/Sink 配置、原理与示例

SeaTunnel OssJindoFile 连接器实战指南:通过 Jindo SDK 对接阿里云 OSS 的 Source/Sink 配置、原理与示例 SeaTunnel OssJindoFile 连接器实战指南通过 Jindo SDK 对接阿里云 OSS 的 Source/Sink 配置、原理与示例【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 OssJindoFile 连接器connector-file-jindo-oss基于阿里云 EMR 的 Jindo SDK通过 HDFS 协议访问阿里云 OSS 对象存储同时提供 Source读取与 Sink写出两种能力支持 text、csv、parquet、orc、json、excel、xml、binary 等多种文件格式并内置 exactly-once 语义。本文以 changelog 文档 与 Source 官方文档、Sink 官方文档 为主线结合 connector-file-jindo-oss 模块源码 讲解其原理与完整实战配置读完即可在 Spark、Flink 或 SeaTunnel Zeta 引擎上完成 OSS 数据入湖、文件迁移与二进制文件同步。连接器概览与演进历史OssJindoFile 连接器最早以 Add oss jindo source sink connector (#3456)插件历史变更与本仓库的 connector-file-jindo-oss.md 变更日志 中几个关键里程碑如下版本关键变更2.3.0新增 OssJindo Source 与 Sink 连接器2.3.2修复 file-oss 配置检查 bug修正 file-oss-jindo 的 factoryIdentifier2.3.3优化 jindo oss 连接器新增file_filter_pattern文件过滤参数2.3.4统一文件 Source/Sink 选项并更新文档支持自定义行分隔符写入文本文件2.3.5为 SFTP、FTP、LocalFile、HdfsFile 等文件连接器统一支持 XML 文件类型2.3.6支持任意文件的传输binary 模式parquet 支持 timestamp/fixed 写为 int962.3.9支持null_format文本空值配置Read/WriteStrategy 由setSeaTunnelRowTypeInfo改为setCatalogTable2.3.10支持filename_extension读写参数、文件 Sink 单文件模式、无数据时创建空文件2.3.11文本文件 Sink 支持row_delimiter选项2.3.12文本文件处理支持自定义行分隔符maxcompute sink writer 支持 timestamp 字段类型从 plugin-mapping.properties 可以看到seatunnel.source.OssJindoFile与seatunnel.sink.OssJindoFile两个插件标识统一映射到connector-file-jindo-oss这一个 Maven 模块即同一模块同时提供读取与写出能力。支持引擎与核心特性根据官方文档该连接器同时支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎。Source 侧特性batch 批处理支持批式读取multimodal 多模态使用 binary 文件格式可读取任意格式文件视频、图片、压缩包等exactly-once在单次 pollNext 调用中读取一个 split 的全部数据split 读取情况保存于 snapshot 中保证精确一次parallelism 并行度支持并行读取支持的文件格式textcsvparquetorcjsonexcelxmlbinarymarkdownpdfSink 侧特性multimodal以 binary 格式写出任意格式文件exactly-once默认使用 2PC两阶段提交保证精确一次多表写入支持 multiple table write支持的文件格式textcsvparquetorcjsonexcelxmlbinarycanal_jsondebezium_jsonmaxwell_json注意 Sink 侧额外支持canal_json、debezium_json、maxwell_json三种 CDC 事件格式Source 侧则额外支持markdown与pdf两种文档解析格式Markdown 只支持读取、不支持写出。环境依赖与 Jindo SDK 安装该连接器通过 Jindo阿里云 EMRSDK 走 HDFS 协议访问 OSS因此依赖 Jindo 与 Hadoop 相关 jar 包官方文档给出了明确的安装前提下载jindosdk-4.6.1.tar.gz解压后将其lib目录下的jindo-sdk-4.6.1.jar与jindo-core-4.6.1.jar复制到${SEATUNNEL_HOME}/lib下且每个运行作业的节点都需要放置。如果使用 Spark/Flink需要确保集群已集成 Hadoop官方测试的 Hadoop 版本为 2.x。如果使用 SeaTunnel Engine安装 SeaTunnel Engine 时已自动集成 Hadoop jar可在${SEATUNNEL_HOME}/lib下确认。由于为支持更多文件类型而内部走 HDFS 协议访问 OSS连接器只支持 Hadoop 版本2.9.X。从 pom.xml 可以看出该模块声明了hadoop-common 2.9.2依赖且作用域为provided运行时由运行环境提供并依赖同级的connector-file-base基础模块。核心原理OssConf 如何把 Jindo 桥接到 Hadoop 抽象从源码层面看整个连接器只做了两件事把 OSS 配置翻译成 Hadoop 文件系统配置以及复用connector-file-base中通用的文件读写实现。OssConf.java 继承自HadoopConf核心逻辑如下HDFS_IMPL com.aliyun.emr.fs.oss.JindoOssFileSystem即 Jindo 提供的 OSS Hadoop 文件系统实现类SCHEMA oss即路径前缀oss://buildWithReadonlyConfig()方法把配置项翻译成 Hadoop 配置项SeaTunnel 配置项Hadoop 配置项说明bucket作为hdfsNameKey传入构造器OSS bucket 地址如oss://tyrantlucifer-image-bedaccess_keyfs.oss.accessKeyIdOSS 访问密钥 IDaccess_secretfs.oss.accessKeySecretOSS 访问密钥 Secretendpointfs.oss.endpointOSS 服务端点固定值fs.AbstractFileSystem.oss.implcom.aliyun.emr.fs.oss.OSS固定值fs.oss.implcom.aliyun.emr.fs.oss.JindoOssFileSystem固定值fs.oss.upload.thread.concurrency上传线程并发数固定 20固定值fs.oss.upload.queue.size上传队列大小固定 100而 OssFileBaseOptions.java 中声明了四个必填的 OSS 专属选项均无默认值ACCESS_KEY // access_key: OSS bucket access key ACCESS_SECRET // access_secret: OSS bucket access secret ENDPOINT // endpoint: OSS endpoint BUCKET // bucket: OSS bucketOssFileSource.java 与 OssFileSink.java 分别继承BaseFileSource与BaseFileSink仅实现initHadoopConf()与getPluginName()两个方法插件名来自FileSystemType.OSS_JINDO值为OssJindoFile见 FileSystemType.java其余文件格式解析、分区处理、事务提交等全部复用 connector-file-base 的通用实现。另外OssJindoFactoryTest.java 验证了OssFileSourceFactory与OssFileSinkFactory的optionRule()均能正常构造可作为连接器注册与选项规则的自检入口。Source 配置详解从 OSS 读取文件必选参数参数类型必填默认值说明pathstring是-源文件路径file_format_typestring是-文件类型bucketstring是-OSS bucket 地址如oss://tyrantlucifer-image-bedaccess_keystring是-OSS 访问密钥access_secretstring是-OSS 访问密钥 Secretendpointstring是-OSS 服务端点如oss-cn-beijing.aliyuncs.com常用可选参数参数类型默认值说明read_columnslist-读取列清单实现字段投影delimiter / field_delimiterstringtext 为\001csv 为,字段分隔符delimiter将在 2.3.5 之后废弃请改用field_delimiterrow_delimiterstring\n行分隔符parse_partition_from_pathbooleantrue是否从文件路径解析分区键与分区值date_formatstringyyyy-MM-dd字符串转 date 的格式datetime_formatstringyyyy-MM-dd HH:mm:ss字符串转 datetime 的格式time_formatstringHH:mm:ss字符串转 time 的格式skip_header_row_numberlong0跳过前 N 行仅 text 与 csvschemaconfig-上游数据结构定义file_filter_patternstring-基于正则的文件过滤模式filename_extensionstring-按扩展名过滤文件compress_codecstringnone压缩编解码器lzo/noneorc/parquet 自动识别archive_compress_codecstringnone归档压缩格式ZIP/TAR/TAR_GZ/GZ/NONEencodingstringUTF-8文件编码null_formatstring-表示 null 的字符串如\Nfile_filter_modified_start / file_filter_modified_endstring-按修改时间过滤yyyy-MM-dd HH:mm:ssquote_charstringCSV 字段包围字符escape_charstring-CSV 转义字符recursive_file_scanbooleantrue是否递归扫描子目录sort_files_by_modification_timebooleanfalse是否按修改时间倒序排序文件其中关键选项的用法细节如下schema与field_delimiter配合当file_format_type为 text 且上游是tyrantlucifer#26#male这类定界文本时不配 schema 会整行作为content字段配置field_delimiter #与 schema 后即可拆成name/age/gender三列。json 类型必须配 schema 才能解析parquet/orc 则能自动识别元数据中的 schema。parse_partition_from_path读取形如oss://hadoop-cluster/tmp/seatunnel/parquet/nametyrantlucifer/age26的路径时会自动为每条数据附加nametyrantlucifer、age26两个分区字段注意不要在 schema 中重复定义分区字段。该选项的默认值为true见 FileBaseSourceOptions.java。file_filter_pattern标准正则表达式仅按文件名过滤时直接写文件名正则如abc.*如需同时匹配目录则表达式需以path开头如/data/seatunnel/20241007/abc[h,g].*。XML 安全限制出于 XXE 加固考虑file_format_type xml的文件若包含!DOCTYPE ...声明即使是仅定义内部实体的良性声明会以FILE_READ_FAILED错误拒绝读取且没有配置项可以恢复旧行为。若 XML 文件由导出工具生成了 DOCTYPE 头需先移除或预处理再接入。markdown/pdf 文档解析Source 侧可将 markdown 与 pdf 解析为结构化文档元素行element_id、element_type、heading_level、text等字段并通过markdown_rag_metadata_enabled/pdf_rag_metadata_enabled追加source_uri、document_id、chunk_id、chunk_index、content_hash等 RAG 元数据列注意仅支持单栏自上而下排版的 PDF。Source 配置示例读取 ORC 文件OssJindoFile { path /seatunnel/orc bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type orc }读取 JSON 文件并指定 schemaOssJindoFile { path /seatunnel/json bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type json schema { fields { id int name string } } }二进制文件同步multimodal源与 Sink 同时使用binary格式即可把图片、压缩包等任意格式文件从 OSS 迁移到 S3、HDFS 等其他存储env { parallelism 1 job.mode BATCH } source { OssJindoFile { bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com path /seatunnel/read/binary/ file_format_type binary } } sink { OssJindoFile { bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com path /seatunnel/read/binary2/ file_format_type binary } }按文件名过滤读取env { parallelism 1 job.mode BATCH } source { OssJindoFile { bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com path /seatunnel/read/binary/ file_format_type binary // 文件示例 abcD2024.csv file_filter_pattern abc[DX]*.* } } sink { Console { } }Sink 配置详解向 OSS 写出文件必选参数参数类型必填说明pathstring是Sink 写入的目标目录不存在会自动创建bucketstring是OSS bucket 地址access_keystring是OSS 访问密钥access_secretstring是OSS 访问密钥 Secretendpointstring是OSS 服务端点常用可选参数参数类型默认值说明tmp_pathstring/tmp/seatunnel结果文件先写临时路径再用mv提交到目标目录需为 OSS 目录custom_filenamebooleanfalse是否自定义文件名file_name_expressionstring${transactionId}仅custom_filenametrue时生效支持${now}、${uuid}变量filename_time_formatstringyyyy.MM.dd仅custom_filenametrue时生效${now}的时间格式file_format_typestringcsv支持 text/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_jsonfilename_extensionstring-覆盖默认文件扩展名如.xml、.json、datfield_delimiterstringtext 为\001csv 为,仅 text/csv 格式生效row_delimiterstring\n仅 text/csv/json 格式生效have_partitionbooleanfalse是否处理分区partition_byarray-仅have_partitiontrue时生效按选定字段分区partition_dir_expressionstring${k0}${v0}/${k1}${v1}/.../${kn}${vn}/分区目录表达式is_partition_field_write_in_filebooleanfalse分区字段是否同时写入数据文件写 Hive 数据文件时应为 falsesink_columnsarray空写入文件的列字段顺序决定实际写入顺序为空则写全部列is_enable_transactionbooleantrue开启后保证数据不丢失不重复文件名自动加${transactionId}_前缀batch_sizeint1000000单个文件的最大行数与checkpoint.interval共同决定文件切分compress_codecstringnone压缩编解码器max_rows_in_memoryint-仅 excel 格式内存中缓存的最大数据条数sheet_max_rowsint1048576仅 excel 格式每 sheet 最大行数sheet_namestringSheet随机数仅 excel 格式csv_string_quote_modeenumMINIMAL仅 csv 格式ALL/MINIMAL/NONExml_root_tag / xml_row_tag / xml_use_attr_format-RECORDS / RECORD / -仅 xml 格式single_file_modebooleanfalse每个并行度只输出一个文件开启后batch_size失效输出文件名无文件块后缀create_empty_file_when_no_databooleanfalse上游无数据时仍生成对应数据文件parquet_avro_write_timestamp_as_int96 / parquet_avro_write_fixed_as_int96-false / -仅 parquet 格式encodingstringUTF-8仅 json/text/csv/xml 格式merge_update_eventbooleanfalse仅 canal_json/debezium_json/maxwell_json将 UPDATE_AFTER 与 UPDATE_BEFORE 合并为 UPDATE 事件schema_evolution_enabledbooleanfalse为 CDC 管道启用 schema 演进binary 格式不支持关键选项深入说明file_name_expression与事务前缀${now}表示当前时间格式由filename_time_format控制${uuid}表示随机 UUID当is_enable_transactiontrue时会自动在文件名头部追加${transactionId}_。filename_time_format支持的常用符号包括y年、M月、d日、H时、m分、s秒。is_enable_transaction为true时通过 2PC 保证写入目标目录的数据不丢失、不重复。官方文档明确该参数当前仅支持true。注意文件最终扩展名取决于file_format_type其中 text 文件的后缀为txt。batch_size与 checkpoint 的关系在 SeaTunnel Engine 中文件行数由batch_size与checkpoint.interval共同决定——若 checkpoint 间隔足够大sink writer 会持续写入直到文件行数超过batch_size若 checkpoint 间隔很小则每次 checkpoint 触发时都会新建文件。compress_codec支持范围txt / json / csvlzo、noneorclzo、snappy、lz4、zlib、noneparquetlzo、snappy、lz4、gzip、brotli、zstd、noneexcel不支持任何压缩格式schema_evolution_enabledCDC 场景置为true后文件 Sink 可在运行期处理 CDC schema 变更事件ADD/DROP/RENAME/MODIFY COLUMN每次 schema 变更会关闭当前输出文件并按新 schema 打开新文件无需重启作业。支持除binary外的所有格式binary 会在作业启动时报配置校验错误当have_partitiontrue时不允许删除partition_by中的分区列若保持默认false而上游 CDC 源开启了schema-changes.enabledtrue则AlterTableEvent到达 Sink 时会立即抛出可操作的错误提示Received AlterTableEvent but schema_evolution_enabledfalse at this sink. ...默认的 CDC 源配置schema-changes.enabled false不受影响。已知限制schema 变更与 checkpoint 不是原子的作业恰好在文件轮转与 schema 元数据更新之间的窄窗口崩溃时恢复后可能按变更前 schema 写行。CDC 管道示例OssJindoFile { path /tmp/cdc/${table_name} bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type parquet schema_evolution_enabled true have_partition true partition_by [updated_at_month] }Sink 配置示例text 格式 分区 自定义文件名env { parallelism 1 job.mode BATCH } sink { OssJindoFile { path/seatunnel/sink bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxx access_secret xxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com 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 格式 指定写出列env { parallelism 1 job.mode BATCH } sink { OssJindoFile { path /seatunnel/sink bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type parquet sink_columns [name,age] } }orc 格式最小配置env { parallelism 1 job.mode BATCH } sink { OssJindoFile { path/seatunnel/sink bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxx access_secret xxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type orc } }canal_json 格式 CDC 事件合并env { parallelism 1 job.mode BATCH } sink { OssJindoFile { path /seatunnel/sink bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxx access_secret xxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type canal_json merge_update_event true } }常见问题与使用建议Hadoop 版本限制连接器内部走 HDFS 协议访问 OSS仅支持 Hadoop 2.9.X使用 Spark/Flink 时需自行保证集群已集成对应版本的 Hadoop使用 SeaTunnel Zeta 时lib目录已内置 Hadoop jar。Jindo jar 缺失作业启动报类加载失败时优先检查每个节点${SEATUNNEL_HOME}/lib下是否存在jindo-sdk-4.6.1.jar与jindo-core-4.6.1.jar。事务与文件前缀开启is_enable_transaction后文件名会带${transactionId}_前缀若对下游文件名有严格要求可结合file_name_expression设计命名规则。分区字段重复定义Source 侧开启parse_partition_from_path后不要在 schema 中重复定义从路径解析出的分区字段。XML 文件被拒若 XML 含!DOCTYPE ...声明会读取失败属于 XXE 加固的安全设计需在接入前去除 DOCTYPE。如需完整的参数列表与更多示例可直接查阅仓库中的 OssJindoFile Source 文档、OssJindoFile Sink 文档、Sink Common Options 与 Source Common Options并在 connector-file-jindo-oss 模块中查看具体实现。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表