You need to enable JavaScript to run this app.
文档中心
文档控制台
注册
E-MapReduce

E-MapReduce

复制全文
下载 pdf
最佳实践
使用 Spark 为 Lance 数据集加列与改列最佳实践
复制全文
下载 pdf
使用 Spark 为 Lance 数据集加列与改列最佳实践

在大数据场景中,数据表往往需要随业务迭代频繁新增或修改列,而传统表需要重写整张表,数据量越大成本越高。Lance 的列式 + 分片(Fragment)存储则可以解决此痛点——新增或修改某列时,只需写入单列数据文件,原有数据文件完全不动,大幅降低变更成本。
本文聚焦增加列(ADD COLUMNS FROM)和修改列(UPDATE COLUMNS FROM)两个高频场景,提供可复用的 SQL 范式及落地最佳实践,适用于需要频繁迭代列的机器学习特征回填(Feature Backfill)场景。

背景信息

为什么 Lance 能做到低成本

Lance 是面向 AI/机器学习场景的现代列式数据格式,与主流表格式一样,Lance 支持传统的 Schema 演进(增加、删除、修改列),且其中大部分操作都不需要重写数据文件,因此已非常高效。
在此之上,Lance 进一步支持数据演进,不仅能改 Schema,还能在不重写原有数据文件的前提下,为存量行回填新列的数据。这一能力的底层原因在于 Lance 的列式 + 分片(Fragment)存储模型:

  • 列文件相互独立。 新增/修改一列时,只需为该列写一份新的列数据文件,原有列的数据文件原封不动。
  • 用行地址而非主键定位。 Lance 通过 _rowaddr(行地址)和 _fragid(分片 ID)精确定位目标数据集中每一行的物理位置,从而把新列的数据"对齐"回填到正确的行上。
  • 低写放大。 相比 MERGE INTOUPDATE SET 这类会重写整行甚至整批受影响的操作,Lance 列级更新只动一列,写入量显著降低。

关于 lance-spark connector

Lance-spark connector 是 Lance 官方提供的 Spark SQL DML 插件,作用是让 Spark 能直接读写 Lance 格式的数据集,它提供的核心 DML 语句是:

  • ALTER TABLE ADD COLUMNS FROM :为存量数据回填新列
  • ALTER TABLE UPDATE COLUMNS FROM :对已有列进行批量更新

配合 lance-spark connector,可以直接用 Spark SQL 完成对 Lance 数据集的大规模加列与改列。

准备工作

必须先启用 Lance Spark SQL Extension 才能使用lance-spark-connector。请在提交 Spark 作业(提交方式包括spark-submit / spark-sql / SparkSession等)时配置对应的 SQL 扩展,详细配置操作请参见Spark SQL Extensions

实践操作

一句话说明核心:​通过业务主键把"目标表"与"数据源表"关联,构造一个携带 _rowaddr_fragid 的临时视图,再用 ALTER TABLE ... ADD/UPDATE COLUMNS ... FROM 落库。Lance 只写入新列文件,不重写原数据、保留行地址、写放大极低

三个必须理解的标识

在开始实践操作前,理解下面三个标识是正确使用的前提。

标识

含义与作用

_rowaddr

行地址。Lance 用它来定位目标数据集中某一行的物理位置。临时视图中必须包含该列,新列数据才能被精确写回对应行。

_fragid

分片 ID(Fragment ID)。标识行所属的数据分片,与 _rowaddr 共同完成"行 → 物理位置"的寻址。临时视图中同样必须包含。

biz_key

业务主键。它是目标表 a数据源表 b之间的关联键,用于把外部计算/外部表中的列值,对齐到目标表的每一行。它本身不直接参与寻址,而是通过 JOINtag_0 (要添加/修改的列) 关联到带有 _rowaddr/_fragid 的目标行上。

注意

易错点: 寻址靠的是 _rowaddr + _fragid,而不是 biz_keybiz_key 只是用来在 SQL 里把数据源 b 的列匹配到目标表 a 的行;最终写回时,Lance 实际需要的是行地址。因此临时视图里既要 JOIN 出 tag_0,又要带上 a_rowaddr/_fragid

增加列:ADD COLUMNS FROM

适用场景

  • 机器学习特征工程中的特征回填:模型迭代时频繁新增特征列,且不希望重写历史数据。
  • 为已有大表新增一个特征列/标签列,数据来自另一张表或一段计算逻辑。
  • 任何"加一列、即时可查、不动原数据"的需求。

基本语法

ALTER TABLE 
<table> ADD COLUMNS <col1>, <col2>, ... FROM <source_view>;  

使用示例

例如,要为表 a 添加一个新列 tag_0,该列数据保存在表 b,且表 a 与表 b 通过业务主键 biz_key 关联:

-- 第一步:构造临时视图(临时视图) 
-- 从目标表 a 取行地址(_rowaddr、_fragid),从数据源 b 取新列值 
-- 注意:_rowaddr 和 _fragid 必须来自目标表 a,不能用 b 的 
CREATE TEMPORARY VIEW tmp_view 
AS 
SELECT 
 a._rowaddr, -- 行地址,用于精确定位目标表中的物理行 
 a._fragid, -- 分片 ID(Fragment ID),配合 _rowaddr 寻址 
 b.tag_0 AS tag_0 -- 新列值,来自数据源 b 
FROM a 
INNER JOIN b ON a.biz_key = b.biz_key; -- 用业务主键(biz_key)关联,只有匹配行才会被加列 
 
-- 第二步:执行加列操作 
-- 将 tmp_view 中的 tag_0 写入表 a,无需重写原有数据文件 
ALTER TABLE a ADD COLUMNS tag_0 FROM tmp_view;

执行后,tag_0 作为一个新列被即时写入,没有表重写或数据搬迁,并立刻可被查询。

注意事项

  • 临时视图必须包含 _rowaddr_fragid。 因为这两列被用来寻址目标数据集中要写入新列数据的行。
  • 新列名不能与已有列冲突。 ADD COLUMNS 是新增列,目标列在表中应不存在。
  • JOIN 的覆盖度决定回填范围。 使用 INNER JOIN 时,只有在 b 中能匹配上 biz_key 的行才会拿到 tag_0 值;未匹配的行其新列值需结合实现语义确认(通常为 NULL 或不回填)。如需保证所有行都有值,应确认数据源 ba 的主键覆盖是否完整。

修改列:UPDATE COLUMNS FROM

适用场景

  • 目标表与数据源表都很大,需要对某一列做大规模批量更新
  • 重算了某个特征/标签,需要用新值覆盖旧值,但不想触发整行重写。
  • 需要保留 _rowaddr、追求低写放大的列级批量更新。

说明

Lance 通过行地址 _rowaddr 来匹配要更新的行,只有匹配上的行才会被更新,数据源中未匹配的行会被忽略。它底层调用的是 Lance 的 fragment.updateColumns() API,保留 _rowaddr_fragid(不重写整行)、写放大更低,这正是它区别于 MERGE INTO 的核心优势。

基本语法

ALTER TABLE <table> UPDATE COLUMNS <col1>, <col2>, ... FROM <source_view>;  

使用示例

例如,假设要修改表 a 的列 tag_0,新值保存在表 b,两表通过业务主键 biz_key 关联:

-- 第一步:构造临时视图(临时视图) 
-- 从目标表 a 取行地址,从数据源 b 取更新后的列值 
-- 注意:_rowaddr 和 _fragid 必须来自目标表 a 
CREATE TEMPORARY VIEW tmp_view 
AS 
SELECT 
 a._rowaddr, -- 行地址,精确定位目标表 a 中要更新的物理行 
 a._fragid, -- 分片 ID(Fragment ID),配合 _rowaddr 寻址 
 b.tag_0 AS tag_0 -- 新的列值,来自数据源 b 
FROM a 
INNER JOIN b ON a.biz_key = b.biz_key; -- 只有 biz_key 匹配的行才会被更新 
 
-- 第二步:执行改列操作 
-- 用 tmp_view 中的 tag_0 覆盖表 a 中对应行的 tag_0 列,原有数据文件不重写 
ALTER TABLE a UPDATE COLUMNS tag_0 FROM tmp_view;

UPDATE COLUMNS 也支持在一条语句里同时更新多列,只要临时视图中包含对应的列即可,示例:

-- 第一步:构造临时视图 
-- 直接从目标表 users 取行地址,同时用常量/表达式计算新列值 
-- 只处理 id = 2 的行(WHERE 过滤) 
CREATE TEMPORARY VIEW update_source AS 
SELECT 
 _rowaddr, -- 行地址,来自目标表 users 本身 
 _fragid, -- 分片 ID(Fragment ID) 
 200 AS value, -- 新的 value 值(常量覆盖) 
 'updated' AS name -- 新的 name 值(常量覆盖) 
FROM users 
WHERE id = 2; -- 只更新 id=2 这一行,其他行不受影响 
 
-- 第二步:一次性更新两列 
-- value 和 name 同时被覆盖写入,无需重写整张表 
ALTER TABLE users UPDATE COLUMNS value, name FROM update_source;

注意事项

  • 数据源必须包含 _rowaddr_fragid。 用于定位目标数据集中要更新的行。
  • UPDATE COLUMNS 中指定的列必须已存在于目标表。 这是它与 ADD COLUMNS 的根本区别:一个改已有列,一个加新列。
  • 未匹配的行被忽略。 仅更新能通过 _rowaddr 匹配上的行;数据源中无法匹配到目标行的记录不会产生影响。

操作要点

视图构建规范

  • 临时视图必须带上 _rowaddr_fragid。 这是这两类操作能成功落库的硬性前提,遗漏会导致无法寻址。
  • 从目标表 a 取行地址。 JOIN 时,_rowaddr_fragid 必须来自目标表 a(被加/改列的那张表),而不是数据源 b

数据质量要求

  • 明确 JOIN 语义与主键覆盖度。 用 INNER JOIN 时只有匹配行被处理;先确认 ba 主键的覆盖是否符合预期,避免出现非预期的 NULL 或漏更新。
  • 保证 biz_key 在数据源侧唯一。 若 bbiz_key 存在重复,JOIN 会产生行膨胀,导致一个目标行被映射到多条值,结果不可控。回填前建议对 b 的主键做去重或唯一性校验。

命令选择优化

  • 区分加列与改列。 列不存在用 ADD COLUMNS(新增);列已存在用 UPDATE COLUMNS(覆盖)。
  • 大表更新优先于 MERGE INTO。 当目标表和数据源都很大、且只动个别列时,UPDATE COLUMNS FROM 的低写放大优势最明显。
  • 善用同表表达式回填。 新列数据不一定来自外部表,也可以是同表上的计算(如 hash(name)),同样遵循"带 _rowaddr/_fragid"的规则。例如,以如下方式创建临时视图,计算字段 name 的 hash 值,并添加到源表。
    -- 第一步:构造临时视图 
    -- 直接从目标表 users 取行地址,同时用同表表达式计算新列值 
    -- 无需关联外部表,_rowaddr 和 _fragid 本身就来自目标表 users 
    CREATE TEMPORARY VIEW tmp_view AS 
    SELECT 
     _rowaddr, -- 行地址,精确定位 users 表中每一行的物理位置 
     _fragid, -- 分片 ID(Fragment ID),配合 _rowaddr 完成寻址 
     hash(name) AS hash_name -- 新列值:对同表 name 列计算哈希(Hash)值 
    FROM users; 
     
    -- 第二步:执行加列操作 
    -- 将 tmp_view 中的 hash_name 写入 users 表,原有数据文件不重写 
    ALTER TABLE users ADD COLUMNS hash_name FROM tmp_view;
    

常见问题

现象

可能原因与排查方向

ALTER TABLE ... COLUMNS ... FROM 语法报错/无法解析

未启用 Lance Spark SQL 扩展。
请检查 Spark 作业是否正确配置了 Lance SQL Extension。

新列/更新值大量为 NULL 或部分行数据未生效

INNER JOIN 下数据源 b 未覆盖到这些 biz_key,或视图遗漏了 _rowaddr/_fragid 导致寻址异常。
请核对主键覆盖度与视图列。

结果行数异常/值出现重复错乱

数据源 bbiz_key 不唯一,JOIN 产生行膨胀。
请对 b 做去重后再回填。

UPDATE COLUMNS 报列不存在

目标列在表 a 中尚不存在,而改列要求列已存在。
若是新列,请改用 ADD COLUMNS

ADD COLUMNS 报列冲突

目标列已存在,而新增列要求列不存在。
若是覆盖已有列,请改用 UPDATE COLUMNS

最近更新时间:2026.07.16 14:49:22
这个页面对您有帮助吗?
有用
有用
无用
无用