拓冰建站拓冰建站
首页 / 资讯中心 / 正文

SeaTunnel Lance Sink 连接器完全指南:配置、数据类型映射与写入模式实战

SeaTunnel Lance Sink 连接器完全指南配置、数据类型映射与写入模式实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel导读本指南基于 Apache SeaTunnel 开源仓库中的 Lance Sink 连接器文档系统讲解如何将 SeaTunnel 数据写入 Lance 数据集。你将掌握 Lance Sink 的全部配置项及其含义、SeaTunnel 与 Lance/Arrow 的数据类型映射规则、三种典型写入场景批处理建表写入、APPEND 追加大数据量、流式按 Checkpoint 刷新的完整配置并深入理解连接器底层的目录 namespace 机制、批量事务写入与 schema 恢复等实现原理可直接用于构建面向向量数据库生态的落库作业。概述Lance Sink 能做什么Lance 是一个基于列式存储与 Apache Arrow 内存格式的现代数据格式由 lancedb 项目发展而来常被用作大规模向量检索与 AI 应用的数据底座。SeaTunnel 的 Lance Sink 连接器负责把上游数据写入 Lance 数据集它可以根据上游 SeaTunnel 表结构自动创建 Lance 表schema 由 SeaTunnel 类型转换得到可以按配置的Lance 写入模式CREATE / APPEND / OVERWRITE创建新数据集或向已有数据集追加数据当前支持基于目录dir的 Lance namespace即本地文件系统路径上的数据集。从实现上看连接器对应源码位于 connector-lance 模块插件标识为Lance核心写入逻辑在 LanceSinkWriter类型转换在 LanceTypeMapper。支持的引擎与特性引擎支持情况SeaTunnel Zeta✅ 支持Spark✅ 3.4 及以上版本Flink❌ 暂不支持连接器特性支持矩阵连接器 v2 特性说明✅批处理Batch✅流处理Streaming✅多表写入Multi-Table Sink❌ 精确一次Exactly-Once❌ CDC仅支持按行追加无 CDC 语义❌ 定时刷新说明连接器同时实现了SupportMultiTableSink与SupportMultiTableSinkWriter见 LanceSink.java因此支持在一个作业中把多张上游表分别写入各自的 Lance 数据集。依赖声明如果需要在 Maven 工程中直接依赖 Lance 相关组件可声明如下依赖与当前仓库 connector-lance/pom.xml 所用版本一致dependency groupIdcom.lancedb/groupId artifactIdlance-core/artifactId version0.33.0/version /dependency dependency groupIdcom.lancedb/groupId artifactIdlance-namespace-core/artifactId version0.0.14/version /dependency其中lance-core提供Dataset、WriteParams、Transaction等核心 APIlance-namespace-core提供目录 namespaceDirectoryNamespace等实现。实际部署时无需手工管理这些依赖连接器模块已被 SeaTunnel 发行包内置。Sink 配置项总览下表汇总了 Lance Sink 的全部配置项源码定义见 LanceCommonOptions.java 与 LanceSinkOptions.java名称类型是否必填默认值说明dataset_pathstring否/test.lanceLance 数据集路径。目录 namespace 下通常是本地数据路径。namespace_typestring否dirLance namespace 类型。当前仅支持dir。namespace_idstring否Lance namespace ID。namespace_idslist否[]解析目标表 namespace 时使用的 namespace 路径片段。root_namespace_pathstring否/tmpLance namespace 的根路径。tablestring否test目标 Lance 表名。设置后会覆盖上游表名。lance.write.max-rows-per-fileint否10单个 Lance 文件最多写入的行数。lance.write.max-rows-per-groupint否20单个 Lance row group 最多写入的行数。lance.write.max-bytes-per-filelong否20480单个 Lance 文件最多写入的字节数。lance.write.modestring否CREATELance 写入模式会传给 LanceWriteParams.WriteMode。lance.write.enable.stable.row.idsboolean否true写入 Lance 时是否启用稳定 row ID。lance.write.storage.optionsmap否{}传给 Lance 的额外存储参数。multi_table_sink_replicaint否1多表写入时的 sink 并行副本数。注意从 LanceSinkFactory.optionRule() 可以看到dataset_path与namespace_type虽然带有默认值但在工厂的 OptionRule 中被声明为必填项required并且附带notBlank非空校验条件——即显式配置时不能为空字符串或仅含空白字符。省略时仍会使用默认值但按工厂规则推荐显式声明。dataset_path数据集落盘位置Lance 数据的目录或数据集路径。使用本地目录模式时请确保 SeaTunnel 运行环境有权限创建并写入该路径。默认值为/test.lance显式配置时该值不能为空字符串或仅包含空白字符对应工厂中的notBlank校验实际写入时LanceSinkWriter会以该路径直接调用Dataset.create(...)或Dataset.open(...)见 LanceSinkWriter.java#L87-L121。namespace_typenamespace 类型Lance namespace 类型当前连接器仅支持dir目录模式。默认值即dir显式配置时同样不能为空字符串。源码层面的支撑连接器内置了完整的 namespace 类型枚举 LanceNamespaceType.java包含rest、dir、hive2、hive3、glue五种候选类型及其对应实现类但当前连接器仅对接了dirDirectoryNamespace。加载过程见 LanceCatalogLoader.loadNamespace()它把root_namespace_path作为root属性传入LanceNamespaces.connect(...)从而构造出指向本地根目录的 namespace。namespace_id 与 namespace_idsnamespace_id目录 namespace 实现使用的 namespace 名称本地目录模式下可以填写类似root的简单名称。它会在 LanceSinkFactory.renameCatalogTable() 中作为 Catalog 名称兜底当上游表没有 catalog 名时。namespace_ids解析目标表 namespace 时使用的额外路径片段。如果直接写入根 namespace可以保持为空。当显式配置时其第一个元素会被用作目标表的 namespace覆盖上游表的 schema 名。root_namespace_pathnamespace 根目录Lance namespace 的根目录默认/tmp。SeaTunnel 运行用户需要有权限在该目录下创建和写入文件。从 LanceCatalogLoader 源码可见该值会作为root属性传给DirectoryNamespace最终决定数据集所在的根位置。table目标表名目标 Lance 表名默认值test。不设置时如果上游存在表名连接器会使用上游表名见renameCatalogTable中StringUtils.isNotEmpty(sinkConfig.getTable())的分支逻辑配置了就用配置值否则取上游tableId.getTableName()。设置后覆盖上游表名。lance.write.mode写入模式控制 Lance 的写入方式默认值CREATE。该值需要是 LanceWriteParams.WriteMode支持的值CREATE创建新数据集若数据集不存在则创建已存在时行为取决于 Lance 底层实现APPEND保留已有数据集并追加写入新行OVERWRITE覆盖已有数据集内容。连接器在 LanceSinkConfig 构造时通过WriteParams.WriteMode.valueOf(...)将其解析为枚举随后在LanceSinkWriter.initializeDataset()中通过WriteParams.Builder().withMode(...)传入 LanceLanceSinkWriter.java#L100-L110。lance.write.max-rows-per-file / max-rows-per-group / max-bytes-per-file文件分片控制这三个参数直接映射到 LanceWriteParams控制单个 Lance 文件fragment与 row group 的规模lance.write.max-rows-per-file单个 Lance 文件最多写入的行数默认 10lance.write.max-rows-per-group单个 Lance row group 最多写入的行数默认 20lance.write.max-bytes-per-file单个 Lance 文件最多写入的字节数默认 20480即 20 KB源码中定义为2048 * 10L。在追加大量数据时调大这些阈值可以减少产生的 Lance fragment 数量降低后续扫描与压缩开销。lance.write.enable.stable.row.ids稳定 row ID已知缺口写入 Lance 时是否启用稳定的 row ID默认true。连接器会把该选项读入LanceSinkConfig.enableStableRowIds字段并通过getEnableStableRowIds()暴露但当前实现中该值仅被解析还未真正传入底层的 LanceWriteParams查看 LanceSinkWriter.initializeDataset() 构造的WriteParams.Builder()其中只设置了maxBytesPerFile、maxRowsPerFile、mode、storageOptions不包含 stable row IDs 开关。因此当前切换该配置对写入路径没有可见效果这是一项已记录的缺口需要后续连接器提交来补齐。在官方修复前建议保持默认值不要依赖该配置改变行为。lance.write.storage.options额外存储参数以键值对形式传递额外的 Lance 存储参数默认{}。这些参数会原样放入WriteParams传给 Lance。示例lance.write.storage.options { key1 value1 key2 value2 }multi_table_sink_replica多表并行副本数多表写入时的 sink 并行副本数默认 1。当一个多表作业写入大量 Lance 表、单个副本成为瓶颈时调大该值。这是 SeaTunnel 的 Sink 通用选项详见 Sink 通用选项。数据类型映射SeaTunnel → Lance / ArrowLance 使用 Apache Arrow 类型系统sink 会根据上游 SeaTunnel 表结构创建 Lance schema。当前映射把所有整数类型TINYINT、SMALLINT、INT、BIGINT一律收窄为 Arrowint32因此超出有符号 32 位范围的BIGINT值会被截断。SeaTunnel 数据类型Lance / Arrow 数据类型BOOLEANboolTINYINTint32SMALLINTint32INTint32BIGINTint32超出有符号 32 位范围的值会被截断FLOATfloat32DOUBLEfloat64DECIMALdecimal128NULLnullBYTESbinaryDATEdate32TIMEtime32毫秒精度TIMESTAMPtimestamp微秒精度Asia/Shanghai 时区STRINGutf8ARRAYlistMAPmap上述规则在源码中有两处直接印证schema 创建SchemaUtils.convertSchema()中所有整数类型TINYINT/SMALLINT/INT/BIGINT统一构造为new ArrowType.Int(32, true)TIMESTAMP构造为new ArrowType.Timestamp(TimeUnit.MICROSECOND, Asia/Shanghai)SchemaUtils.java#L64-L143Catalog schema 转换LanceTypeMapper.convertJsonArrowType()中对 TINYINT/SMALLINT/BIGINT/INT 统一设置type int32ARRAY 映射为带element字段的listMAP 映射为entries结构体形式的mapLanceTypeMapper.java#L89-L183。反向转换Lance → SeaTunnel用于 Catalog 读取也已在LanceTypeMapper.convertDataType()中实现支持 bool、各类整数、字符串、decimal、浮点、日期时间、binary 等类型但struct/list/map的完整反向支持仍有 TODO 待完善见源码第 83 行注释。:::tip 写入语义提醒Sink不会按UPDATE/DELETE行类型执行 CDC 语义——每条上游记录都会按lance.write.mode追加到 Lance 数据集中代码层面LanceSinkWriter.flushBatch() 对缓冲区内每一行执行FragmentConverter.reconvert(...)生成 fragment再通过 LanceTransaction以Append操作提交。在流式模式下Writer 会在每个 checkpointprepareCommit()把内存中的行缓冲写入 Lance。:::写入流程与底层实现原理理解底层实现有助于合理配置参数与排查问题。LanceSinkWriter源码的完整写入链路如下惰性初始化数据集收到第一条记录时initializeDataset()先尝试Dataset.open()打开已有数据集并复用其 schema若打开失败数据集不存在则用首条记录推导 Arrow schema通过Dataset.create()以WriteParams含 mode、行数/字节数阈值、存储参数创建数据集后再重新打开。内存行缓冲Writer 内部维护默认容量 1000 行的batchBufferDEFAULT_BATCH_SIZE 1000write()逐行入队达到阈值即触发flushBatch()。事务式批量追加flushBatch()将缓冲区内每行通过FragmentConverter.reconvert()转换为 LanceFragmentMetadata然后用dataset.newTransactionBuilder().operation(Append.builder().fragments(...))提交一个 Lance 事务提交后重新打开数据集以获得最新版本。Checkpoint 与关闭时刷新prepareCommit()与close()都会强制flushBatch()abortPrepare()则清空缓冲。配合引擎的 checkpoint 机制流式作业实现按 checkpoint 落盘的效果。schema 恢复支持Writer 实现了SupportSchemaEvolutionSinkWriter在处理RestoreTableSchemaEvent时先用旧 schema 刷新剩余行再切换到 checkpoint 状态中的运行时 schema该路径在单元测试LanceSinkTest.restoreRuntimeSchemaBeforeWritingRowsWithCheckpointLayout()中有直接验证见 LanceSinkTest.java。此外连接器的 Catalog 实现 LanceCatalog.java 负责 namespace 级别的表管理createTable/dropTable/listTables/tableExists/getTable并在createTable时通过 Arrow IPC 流写入 schema 元数据主键、表注释、选项、列注释等以seatunnel.*前缀保存为 schema 元数据同时按root_namespace_path dataset_path table .lance的规则计算数据集落盘路径见getDatasetPath()。任务示例示例一写入 FakeSource 数据到 LanceCREATE 建表写入批处理模式下从 FakeSource 读取多类型数据并写入 Lance 数据集。运行后会按上游 schema 自动创建数据集env { parallelism 1 job.mode BATCH } source { FakeSource { row.num 100 schema { fields { 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 } } plugin_output fake } } sink { Lance { dataset_path /tmp/seatunnel_mnt/lanceTest/lance_sink_table namespace_type dir namespace_id root table lance_sink_table } }要点说明dataset_path指向最终 Lance 数据集的落盘目录运行前需确保该路径所在目录对 SeaTunnel 进程可写namespace_id root表示使用根 namespace此时root_namespace_path默认/tmp可作为相对基准但dataset_path使用绝对路径时以绝对路径为准table显式指定表名lance_sink_table否则会沿用上游表名示例中的c_bigint等整数字段在 Lance 侧会收窄为int32请避免写入超出 32 位有符号范围的值。示例二使用 APPEND 模式并调大文件分片APPEND模式会保留已有数据集并写入新行。把lance.write.max-rows-per-file和lance.write.max-bytes-per-file调大可以减少追加大批量数据时产生的 Lance fragment 数量env { parallelism 1 job.mode BATCH } source { FakeSource { row.num 1000000 schema { fields { c_string string c_int int } } plugin_output fake } } sink { Lance { dataset_path /tmp/seatunnel_mnt/lanceTest/lance_sink_table namespace_type dir namespace_id root table lance_sink_table lance.write.mode APPEND lance.write.max-rows-per-file 100000 lance.write.max-rows-per-group 5000 lance.write.max-bytes-per-file 134217728 } }调参建议max-bytes-per-file 134217728128 MB与默认的 20 KB 相比提升了 3 个数量级适合百万行量级的批量追加max-rows-per-group也相应放大到 5000让每个 row group 装下更多行减少元数据开销。需要注意这些阈值最终会原样传给 LanceWriteParams具体效果取决于 Lance 版本对文件切分的实际执行。示例三流式追加并按 Checkpoint 刷新流式模式下Writer 在每个 checkpoint 将内存中的行缓冲写入 Lance。下面的配置把checkpoint.interval设为 30 秒意味着最多 30 秒的数据会暂存在内存中随后以事务方式一次性追加env { parallelism 2 job.mode STREAMING checkpoint.interval 30000 } source { FakeSource { row.num 1000 schema { fields { c_string string c_int int } } plugin_output fake_stream } } sink { Lance { plugin_input fake_stream dataset_path /tmp/seatunnel_mnt/lanceTest/lance_sink_table namespace_type dir namespace_id root table lance_sink_table lance.write.mode APPEND } }要点说明plugin_input fake_stream显式声明数据来源与上游plugin_output对应lance.write.mode APPEND保证每次 checkpoint 刷新都是增量追加不会覆盖历史数据除了 checkpoint 周期Writer 内部的 1000 行缓冲阈值DEFAULT_BATCH_SIZE也会触发中途刷新两者共同决定实际落盘粒度。常见问题与使用建议BIGINT 截断由于整数统一映射为int32超出 ±2,147,483,647 范围的BIGINT值会静默截断。需要完整 64 位精度时当前版本需在上游预处理如拆分为字符串/多列或等待连接器后续版本修正映射。Flink 不可用Lance Sink 仅支持 SeaTunnel Zeta 与 Spark 3.4Flink 引擎下无法使用该连接器。写入权限dataset_path与root_namespace_path对应的目录必须对运行 SeaTunnel 的用户可写否则Dataset.create()会抛出TABLE_DATASET_PATH_OPEN_EXCEPTION对应错误码见 LanceConnectorErrorCode.java。CDC 语义缺失上游UPDATE/DELETE行不会被特殊处理所有记录一律按追加写入如有去重/更新需求需在上游或 transform 阶段先行处理。稳定 row ID 配置暂未生效lance.write.enable.stable.row.ids目前只被解析不参与实际写入参数构造改动它不会改变行为升级时请留意更新日志。更新日志连接器的变更记录见 connector-lance 更新日志升级 SeaTunnel 版本后建议核对该页确认写入模式、类型映射等行为是否有调整。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门