SeaTunnel Lance Sink 连接器完全指南:配置、数据类型映射与写入模式实战
2026/9/18 12:38:24 网站建设 项目流程

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 语义)
  • ❌ 定时刷新

说明:连接器同时实现了SupportMultiTableSinkSupportMultiTableSinkWriter(见 LanceSink.java),因此支持在一个作业中把多张上游表分别写入各自的 Lance 数据集。

依赖声明

如果需要在 Maven 工程中直接依赖 Lance 相关组件,可声明如下依赖(与当前仓库 connector-lance/pom.xml 所用版本一致):

<dependency> <groupId>com.lancedb</groupId> <artifactId>lance-core</artifactId> <version>0.33.0</version> </dependency> <dependency> <groupId>com.lancedb</groupId> <artifactId>lance-namespace-core</artifactId> <version>0.0.14</version> </dependency>

其中lance-core提供DatasetWriteParamsTransaction等核心 API,lance-namespace-core提供目录 namespace(DirectoryNamespace)等实现。实际部署时无需手工管理这些依赖,连接器模块已被 SeaTunnel 发行包内置。

Sink 配置项总览

下表汇总了 Lance Sink 的全部配置项(源码定义见 LanceCommonOptions.java 与 LanceSinkOptions.java):

名称类型是否必填默认值说明
dataset_pathstring/test.lanceLance 数据集路径。目录 namespace 下通常是本地数据路径。
namespace_typestringdirLance namespace 类型。当前仅支持dir
namespace_idstring""Lance namespace ID。
namespace_idslist[]解析目标表 namespace 时使用的 namespace 路径片段。
root_namespace_pathstring/tmpLance namespace 的根路径。
tablestringtest目标 Lance 表名。设置后会覆盖上游表名。
lance.write.max-rows-per-fileint10单个 Lance 文件最多写入的行数。
lance.write.max-rows-per-groupint20单个 Lance row group 最多写入的行数。
lance.write.max-bytes-per-filelong20480单个 Lance 文件最多写入的字节数。
lance.write.modestringCREATELance 写入模式,会传给 LanceWriteParams.WriteMode
lance.write.enable.stable.row.idsbooleantrue写入 Lance 时是否启用稳定 row ID。
lance.write.storage.optionsmap{}传给 Lance 的额外存储参数。
multi_table_sink_replicaint1多表写入时的 sink 并行副本数。

注意:从 LanceSinkFactory.optionRule() 可以看到,dataset_pathnamespace_type虽然带有默认值,但在工厂的 OptionRule 中被声明为必填项(required),并且附带notBlank非空校验条件——即显式配置时不能为空字符串或仅含空白字符。省略时仍会使用默认值,但按工厂规则推荐显式声明。

dataset_path:数据集落盘位置

Lance 数据的目录或数据集路径。使用本地目录模式时,请确保 SeaTunnel 运行环境有权限创建并写入该路径。

  • 默认值为/test.lance
  • 显式配置时该值不能为空字符串或仅包含空白字符(对应工厂中的notBlank校验);
  • 实际写入时,LanceSinkWriter会以该路径直接调用Dataset.create(...)Dataset.open(...)(见 LanceSinkWriter.java#L87-L121)。

namespace_type:namespace 类型

Lance namespace 类型,当前连接器仅支持dir(目录模式)。默认值即dir,显式配置时同样不能为空字符串。

源码层面的支撑:连接器内置了完整的 namespace 类型枚举 LanceNamespaceType.java,包含restdirhive2hive3glue五种候选类型及其对应实现类,但当前连接器仅对接了dir(DirectoryNamespace)。加载过程见 LanceCatalogLoader.loadNamespace():它把root_namespace_path作为root属性传入LanceNamespaces.connect(...),从而构造出指向本地根目录的 namespace。

namespace_id 与 namespace_ids

  • namespace_id:目录 namespace 实现使用的 namespace 名称,本地目录模式下可以填写类似root的简单名称。它会在 LanceSinkFactory.renameCatalogTable() 中作为 Catalog 名称兜底(当上游表没有 catalog 名时)。
  • namespace_ids:解析目标表 namespace 时使用的额外路径片段。如果直接写入根 namespace,可以保持为空。当显式配置时,其第一个元素会被用作目标表的 namespace(覆盖上游表的 schema 名)。

root_namespace_path:namespace 根目录

Lance namespace 的根目录,默认/tmp。SeaTunnel 运行用户需要有权限在该目录下创建和写入文件。从 LanceCatalogLoader 源码可见,该值会作为root属性传给DirectoryNamespace,最终决定数据集所在的根位置。

table:目标表名

目标 Lance 表名,默认值test。不设置时,如果上游存在表名,连接器会使用上游表名(见renameCatalogTableStringUtils.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(...)传入 Lance(LanceSinkWriter.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 文件最多写入的行数,默认 10;
  • lance.write.max-rows-per-group:单个 Lance row group 最多写入的行数,默认 20;
  • lance.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(),其中只设置了maxBytesPerFilemaxRowsPerFilemodestorageOptions不包含 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 / Arrow

Lance 使用 Apache Arrow 类型系统,sink 会根据上游 SeaTunnel 表结构创建 Lance schema。当前映射把所有整数类型(TINYINTSMALLINTINTBIGINT一律收窄为 Arrowint32,因此超出有符号 32 位范围的BIGINT值会被截断。

SeaTunnel 数据类型Lance / Arrow 数据类型
BOOLEANbool
TINYINTint32
SMALLINTint32
INTint32
BIGINTint32(超出有符号 32 位范围的值会被截断)
FLOATfloat32
DOUBLEfloat64
DECIMALdecimal128
NULLnull
BYTESbinary
DATEdate32
TIMEtime32(毫秒精度)
TIMESTAMPtimestamp(微秒精度,Asia/Shanghai 时区)
STRINGutf8
ARRAYlist
MAPmap

上述规则在源码中有两处直接印证:

  1. schema 创建SchemaUtils.convertSchema()中所有整数类型(TINYINT/SMALLINT/INT/BIGINT)统一构造为new ArrowType.Int(32, true)TIMESTAMP构造为new ArrowType.Timestamp(TimeUnit.MICROSECOND, "Asia/Shanghai")(SchemaUtils.java#L64-L143);
  2. Catalog schema 转换LanceTypeMapper.convertJsonArrowType()中对 TINYINT/SMALLINT/BIGINT/INT 统一设置type = "int32",ARRAY 映射为带element字段的list,MAP 映射为entries结构体形式的map(LanceTypeMapper.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,再通过 LanceTransactionAppend操作提交)。在流式模式下,Writer 会在每个 checkpoint(prepareCommit())把内存中的行缓冲写入 Lance。

:::

写入流程与底层实现原理

理解底层实现有助于合理配置参数与排查问题。LanceSinkWriter(源码)的完整写入链路如下:

  1. 惰性初始化数据集:收到第一条记录时,initializeDataset()先尝试Dataset.open()打开已有数据集并复用其 schema;若打开失败(数据集不存在),则用首条记录推导 Arrow schema,通过Dataset.create()WriteParams(含 mode、行数/字节数阈值、存储参数)创建数据集后再重新打开。
  2. 内存行缓冲:Writer 内部维护默认容量 1000 行的batchBufferDEFAULT_BATCH_SIZE = 1000),write()逐行入队,达到阈值即触发flushBatch()
  3. 事务式批量追加flushBatch()将缓冲区内每行通过FragmentConverter.reconvert()转换为 LanceFragmentMetadata,然后用dataset.newTransactionBuilder().operation(Append.builder().fragments(...))提交一个 Lance 事务,提交后重新打开数据集以获得最新版本。
  4. Checkpoint 与关闭时刷新prepareCommit()close()都会强制flushBatch()abortPrepare()则清空缓冲。配合引擎的 checkpoint 机制,流式作业实现"按 checkpoint 落盘"的效果。
  5. 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 数据到 Lance(CREATE 建表写入)

批处理模式下,从 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-filelance.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 = 134217728(128 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.4+,Flink 引擎下无法使用该连接器。
  • 写入权限dataset_pathroot_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询