You need to enable JavaScript to run this app.
文档中心
E-MapReduce

E-MapReduce

复制全文
下载 pdf
数据治理
Lance Compact
复制全文
下载 pdf
Lance Compact

本文以 Lance 的存储模型和 Compact 生命周期为主线,说明为什么需要 Compact、一次 Compact 实际做了什么、资源压力从哪里产生,以及如何规划参数和验收结果。最后以 Daft 为例说明分布式执行参数。除 Daft 示例章节外,其余原则同样适用于 Lance 原生 Python SDK、Ray、Spark 等执行方式。

为什么需要 Compact

Lance数据集频繁小批量追加会不断增加 Fragment。大量小 Fragment 会提高打开表、读取 Manifest、规划扫描和对象存储请求的固定开销。删除操作则可能留下 Deletion File,读取时仍需扫描物理行并过滤被删除记录。

Compact 的主要收益

  • 把相邻小 Fragment 合并为更合适的 Fragment,减少 Fragment 和小 Data File 数量。
  • 把已标记删除的行真正从新 Data File 中去掉,后续读取无需再过滤这些行。
  • 如果表已删除某些列,Compact 在 Rewrite 时不会再把这些旧列写入新 Data File。
  • 降低 Manifest 级操作和对象存储请求的成本。

Compact 不解决什么

  • Compact 不是历史 Version 清理;它创建新 Version,但不会立即删除旧 Version 仍引用的 Data File。
  • Compact 不是“Version 越多就必须执行”。Version 数和 Fragment 布局是两类问题。
  • Compact 不能突破对象存储、网络或内存上限;不合理的目标大小和并发反而会放大写入成本。

先理解 Lance 表的存储结构

Lance Dataset 不是“一个大 Data File”,而是由版本化 Manifest 管理的 Fragment、Data File、Deletion File 和 Index 集合。

对象

作用

与 Compact 的关系

Manifest / Version

描述某个不可变快照所引用的 Fragment、Data File、Index 和事务信息

Compact 最终通过一次 Commit 发布新 Manifest,也就是创建一个新 Version

Fragment

表的行分片,承载一段连续行;一个 Fragment 可以包含覆盖不同列集合的多个 Data File,但这些 Data File 共同描述同一批行

Compact 的基本规划对象。通常把若干相邻小 Fragment Rewrite 为较少的大 Fragment

Data File

保存 Fragment 中一个或多个列的数据

Rewrite 阶段读取旧 Data File,并写出新的 Data File

Deletion File

记录某个 Fragment 中哪些行已删除;此时旧 Data File 尚未 Rewrite

Compact 时可以只把未删除的行写入新 Fragment,之后读取就不再需要额外过滤

Index

向量、标量、全文等二级 Index

如果 Index 保存 Row Address,Compact 改变 Row Address 后需要 Remap

因此,Fragment、Data File 和 Version 是三个不同层次:Fragment 是行布局单元,Data File 是物理存储对象,Version 是某一时刻对 Fragment、Data File、Index 等对象的引用集合。

说明

全文术语约定:​表快照 - Version;物理数据对象 - Data File;数据搬运阶段 - Rewrite;任务计划 - Compaction Plan;执行单元 - Compaction Task;物理行定位 - Row Address;索引地址更新 - Remap;Daft 调度分片 - Daft Partition;一次提交范围 - Commit Batch。软件版本仍称“软件版本”,不与 Lance Version 混用。

Image

Compact 的三阶段原理

Compaction Plan:决定哪些 Fragment 要 Rewrite

Lance 根据当前 Version 的 Fragment 顺序、行数、删除比例和 CompactionOptions 生成 Compaction Plan。Compaction Plan 中包含若干 Compaction Task;每个 Compaction Task 描述一组源 Fragment 以及目标 Rewrite 方式。

Rewrite:真正搬运数据

执行 Compaction Task 时,Lance 分批读取并解码源 Fragment。如果开启 materialize_deletions,它会跳过已标记删除的行,只把其余行编码并写入新的 Data File,最后返回 RewriteResult。此时旧 Data File 仍然存在,旧 Version 仍可正常读取。

Commit:原子发布新 Version

Driver 会汇总每个 Compaction Task 返回的 RewriteResult,并构造事务。Compact 把行写入新的 Fragment 后,Row Address 会变化;如果 Index 保存旧 Row Address,Commit 需要执行 Remap,使 Index 能继续定位移动后的行。最后,Lance 提交新的 Manifest;只有 Commit 成功,新 Version 才对后续读者可见。
Image

参数分层:从 Compaction Plan 到进程 I/O

理解 Compact 的三个阶段后,可以把参数按四层区分:Lance 如何生成 Compaction Plan、单个 Compaction Task 如何 Rewrite、Daft 如何调度与 Commit、每个 Lance 进程如何执行 I/O。先定位瓶颈属于哪一层,再调整该层参数;作用范围不同的参数不能互相替代。

说明

执行引擎说明:​本文涉及分布式调度与分批 Commit 时,以 Daft 为例。表中的 num_partitionsconcurrencymicro_commit_batch_size 是 Daft daft.io.lance.compact_files() 接口的执行或提交参数,不属于 Lance CompactionOptions,也不是 Lance 原生 Compact API 的通用参数。使用 Lance-Ray、Spark 或其他执行引擎时,应改用对应引擎的调度与提交参数。

参数

控制层

准确含义

主要风险

target_rows_per_fragment

Lance Compaction Plan / Rewrite

输出 Fragment 的目标行数,同时影响 Lance 如何组合源 Fragment

过大时单个 Compaction Task 更重;过小时输出 Fragment 和元数据更多

max_bytes_per_file

Lance Rewrite

单个 Data File 的字节软上限

版本不支持时无效;只看行数、不看字节数仍可能产生大文件

materialize_deletions

Lance Rewrite

Rewrite 时只把未删除的行写入新的 Data File

本次 Rewrite 数据量增加,但后续读取无需再过滤这些删除记录

materialize_deletions_threshold

Lance Compaction Plan

删除比例达到阈值时,将 Fragment 纳入 Rewrite 候选

阈值越低,因删除而进入 Compaction Plan 的 Fragment 越多

batch_size

Lance Rewrite

单个 Compaction Task 每次扫描多少行组成一个 Arrow RecordBatch

调大可能提高吞吐,也会提高单个 Compaction Task 的内存峰值

num_threads

Lance Rewrite

单个 Compaction Task 内部的本地计算线程数

调大可能增加 CPU 争抢和同时存在的中间 Buffer

num_partitions

Daft 调度

当前 Commit Batch 被划分成多少个 Daft Partition;不是 Compaction Task 数

过小限制并行,过大增加调度开销

concurrency

Daft 调度

最多同时执行多少个 Daft 工作单元

过大时节点内存、网络和对象存储请求同时增加

micro_commit_batch_size

Daft 提交

一个 Commit Batch 最多包含多少个 Compaction Task

过小增加 Commit 和 Version 数;过大增加 Driver 结果汇总与 Commit 峰值

max_source_fragments

Lance Compaction Plan

限制一次 Compaction Plan 最多纳入多少个源 Fragment;默认不限制

过小会使本次 Compact 只处理很少的 Fragment,需要多次调用;取值应至少能容纳一个完整 Compaction Task

LANCE_DEFAULT_IO_BUFFER_SIZE

Lance 进程 I/O

限制扫描可以提前读取的数据量

过大抬高 Worker 内存;过小降低 I/O 吞吐

LANCE_IO_THREADS

Lance 进程 I/O

每个 Lance 进程的 I/O 并行度

会与 Daft 的 concurrency叠加;过大可能造成限流、排队和更多 Buffer

LANCE_DEFAULT_FRAGMENT_READAHEAD

Lance 进程 I/O

控制跨 Fragment 的预读深度

过大抬高预读内存;过小降低流水线吞吐

Lance 如何生成 Compaction Task

Compaction Task 数不能用“总行数 ÷ target_rows_per_fragment”精确计算。这个除法只能估算最终输出 Fragment 的数量级;Compaction Task 数由 Lance 根据相邻 Fragment 的组合、删除比例和 CompactionOptions 生成。

哪些因素会改变 Compaction Task 数

  • 现有 Fragment 的数量和大小:​相邻小 Fragment 越多,通常越容易被组合成 Compaction Task,因此 Compaction Task 数可能增加。
  • target_rows_per_fragment:目标值越大,Lance 通常会让一个 Compaction Task 吸收更多源 Fragment,Compaction Task 数可能减少,但单个 Compaction Task 更重;目标值越小,输出 Fragment 和 Compaction Task 数可能增加。
  • materialize_deletions_threshold:阈值越低,越多含删除记录的 Fragment 会进入 Compaction Plan,因此 Compaction Task 数和 Rewrite 数据量都可能增加。
  • max_bytes_per_file:Data File 达到字节软上限时会提前切分。它主要改变输出 Data File / Fragment 数,不一定按相同比例改变 Compaction Task 数。
  • max_source_fragments:Lance 原生 Compaction Plan 参数。它限制本次 Compaction Plan 纳入的源 Fragment 总数;Lance 按顺序加入完整的 Compaction Task,如果加入下一个 Compaction Task 会超过上限,就不再加入。默认值为None,即不限制。

一个 Compaction Task 可以包含多个 Fragment

可以,而且这是常见情况。Lance 会把若干相邻小 Fragment 组合到同一个 Compaction Task;执行后也可能产生一个或多个输出 Fragment。为保持行的相对顺序,Lance 通常不会跨过未参与 Compact 的 Fragment,把其两侧的 Fragment 随意组合。

一个 Compaction Task 也可能只包含一个 Fragment

可能。例如某个 Fragment 因删除比例超过 materialize_deletions_threshold而需要物化删除,但它附近没有适合一起合并的 Fragment,这个 Compaction Task 就可能只有一个源 Fragment。判断 Compaction Task 大小不能只看“包含几个 Fragment”,还要看这些 Fragment 的总行数、总字节数和列宽。

Lance 参数如何影响 Compaction Task

Lance 参数

直接控制什么

对 Compaction Task 的影响

target_rows_per_fragment

输出 Fragment 希望达到的行数

调大时通常会合并更多源 Fragment,Compaction Task 数可能减少,但单个 Compaction Task 的行数和字节数可能增加

max_bytes_per_file

单个 Data File 的字节软上限

主要限制输出 Data File 大小;可能让一个 Compaction Task 产生多个输出 Fragment

materialize_deletions

Rewrite 时是否丢弃已删除行

开启后需要读取源数据并写出未删除行,会增加本次 Rewrite 工作量

materialize_deletions_threshold

Fragment 因删除比例进入 Compaction Plan 的门槛

阈值越低,进入 Compaction Plan 的 Fragment 越多,Compaction Task 数或单个 Compaction Task 覆盖范围可能增加

batch_size

单个 Compaction Task 一次扫描的行数

不改变 Compaction Plan 中的 Compaction Task 数;调大通常提高吞吐,也提高单个 Compaction Task 的内存峰值

num_threads

单个 Compaction Task 内部的本地计算线程数

不改变 Compaction Plan 中的 Compaction Task 数;调大可能提高单个 Compaction Task 的计算并行,也可能增加 CPU 争抢和中间 Buffer

设置target_rows_per_fragment要同时看行数和字节数

目标行数 ≈ 目标 Fragment 字节数 ÷ 采样行宽

对 String、Binary、List、Struct 等变长列,应查看平均值和高分位行宽。Blob 列可把大对象与普通行数据分离,Compact 时不必反复复制大块内联 Payload;普通 Binary 列则可能在整行 Rewrite 时产生明显写放大。

batch_sizenum_threads控制的是单个 Compaction Task

batch_size控制一个 Compaction Task 每次从源 Fragment 读入多少行组成 Arrow RecordBatch。它越大,一次参与解码、转换和编码的数据越多;吞吐可能提高,但内存峰值也会提高。
num_threads控制同一个 Compaction Task 内部可以使用多少本地计算线程。它不是 Daft 的集群并发;Daft 的concurrency决定同时运行多少个工作单元,而num_threads决定每个工作单元内部的计算并行。两者同时调大,会放大同机 CPU 争抢和内存压力。

“流式执行”为什么仍然可能 OOM

本文所说的“流式执行”,是指执行器把输入拆成有限大小的Batch,按“读取 → 解码 → 处理 → 编码 → 写出”的流水线逐批推进,而不是一次把整张表加载进内存。这里描述的是数据在算子中的处理方式,不是在说某个框架一定属于流式计算引擎:Daft、Ray 上的任务或 Spark Batch 都可以采用分批迭代和 Spill。流式执行只能把内存从“整表大小”降到“若干在途 RecordBatch 与 Buffer”,不等于常量内存或零 Buffer。

  1. 正在读取的压缩数据与对象存储 I/O Buffer。
  2. 已解码的 Arrow RecordBatch,包括变长 String 的 Offsets 和 Values Buffer。
  3. 已预读但尚未消费的 RecordBatch 或 Fragment。
  4. 删除过滤、列转换和编码过程中的中间 Buffer。
  5. 正在写出的压缩页、Data File 元数据和 Multipart Upload Buffer。

因此,单个 Compaction Task 的峰值由扫描 RecordBatch、预读深度、行宽、编码器和写出 Buffer 共同决定;节点峰值还要考虑同机并发的 Compaction Task 数。

节点峰值 ≈ 同机并发 Compaction Task 数 × 单个 Compaction Task 峰值 + 进程运行时、缓存和安全余量

即使数据列主要是 URL 等短字符串,一个 Compaction Task 处理大量行时,字符串内容、偏移数组、空值标记,以及解码、编码阶段的中间 Buffer 仍会同时占用内存。因此,流式执行只能避免加载整张表,不能排除单个 Compaction Task OOM。

用三个环境变量限制单进程预读与 I/O 并行

单个 Lance 进程的 I/O 预读和并行度过高,也会抬高 Worker 瞬时内存。遇到 Worker OOM 时,可以用下面三个进程级环境变量收紧单进程 I/O 压力。

环境变量

Lance 7.0.0 附近版本的默认行为

OOM 时如何调整

LANCE_DEFAULT_IO_BUFFER_SIZE

默认 I/O Buffer 为 2 GiB

调小。它限制单个扫描可提前读取的数据量;调小可降低内存峰值,但过小会降低 I/O 吞吐。

LANCE_IO_THREADS

云对象存储默认 64;本地存储默认 8

调小。它限制每个 Lance 进程同时进行的 I/O 工作;调小可减少并行请求和相关 Buffer,但过小会使读取变慢。

LANCE_DEFAULT_FRAGMENT_READAHEAD

Legacy 扫描路径默认 4;部分新扫描路径会结合 I/O 并行度计算

调小。它限制同时预读的 Fragment 数;这是收紧跨 Fragment 预取内存最直接的参数,但也可能降低流水线并行。

这三个值不是越小越好,而是用吞吐换取更稳定的内存边界。总 I/O 压力仍然取决于两层并发的乘积:

全集群潜在 I/O 并行 ≈ 并发 × 每个 Lance 进程的 I/O 并行

例如concurrency=80LANCE_IO_THREADS=16时,理论上仍可能形成很高的集群 I/O 并行度;把LANCE_IO_THREADS再提高到 80,通常不会让吞吐线性增长,反而可能因对象存储限流、排队和更多 Buffer 而变慢。
Image

Index、Row ID 与 Commit

未启用 Stable Row ID 时,一行数据通常由“Fragment ID + Fragment 内行号”组成的 Row Address 定位。Compact 会把行写入新的 Fragment,Row Address 因而发生变化;如果二级 Index 保存旧 Row Address,Commit 就需要执行 Remap,把 Index 中的旧 Row Address 更新为新 Row Address。Stable Row ID 是不随 Compact 搬迁而变化的逻辑行标识;它属于表格式能力,不是给旧 Dataset 临时增加一个 Compact 参数。

策略

Compact 时做什么

主要权衡

立即 Remap

Commit 时生成“旧 Row Address → 新 Row Address”的映射,并更新受影响 Index 中保存的 Row Address

Compact 完成后查询可直接使用更新后的 Index;代价是 Commit 更慢、读写更多 Index 数据,并占用更多 Driver 内存

延迟 Remap

defer_index_remap=True时先不更新业务 Index,而是保存 Fragment Reuse Index,记录旧 Row Address 如何映射到新 Row Address

本次 Compact 更省资源;查询旧 Index 时要额外做 Row Address 转换。连续延迟会累积映射层,后续仍需择机完成 Remap

Stable Row ID

Index 使用不随 Compact 搬迁而变化的逻辑 Row ID,不再依赖会变化的 Row Address

这是表格式能力,不是旧 Dataset 临时增加一个 Compact 参数即可完成的迁移;是否可用于生产取决于部署的 Lance 软件版本

为什么立即 Remap 容易造成 Driver OOM

在未启用 Stable Row ID 且未延迟 Remap 时,Lance 会在 Commit 中汇总当前 Commit Batch 的 RewriteResult,为受影响行建立“旧 Row Address → 新 Row Address”的内存映射,再用这份映射更新相关 Index。映射条目数与当前 Commit Batch 涉及的行数相关;Commit Batch 越大、删除行越多、Index 越多,Driver 同时持有的 Row Address 映射和 Index 数据就越多。因此可能出现 Worker 已全部完成,Driver 却在 Commit 阶段耗尽普通内存。

有 Index 时:立即 Remap 还是延迟 Remap

  • 优先保证查询性能:​选择立即 Remap。Compact 当次耗时和内存更高,但 Commit 后 Index 可直接使用新 Row Address。
  • 优先降低维护窗口资源峰值:​可以选择延迟 Remap。Compact 当次不更新 Index,但查询旧 Index 时要先做 Row Address 转换;连续延迟会让映射层增加。
  • 延迟 Remap 不是永久免除 Remap。应监控查询延迟和映射层数,并在资源充足的维护窗口完成 Remap 或重建 Index。

没有 Index 时

没有二级 Index 时,没有 Index Data File 需要更新,defer_index_remap本身也没有查询收益。不过,部分 Lance 版本在 Commit 中仍会根据 RewriteResult 构造或汇总 Row Address 映射;分布式执行器还可能一次收集大量 RewriteResult。因此Indices=[]可以排除 Index Data File 更新,却不能自动排除 Driver 内存问题。

版本兼容必须实测

CompactionOptions随 Lance 软件版本变化。只有当已安装的 Lance Python API 接受defer_index_remap,并且 Daft 能把它传给 Lance 时,延迟 Remap 才能使用。若 Lance 7.0.0 返回Invalid compaction option: defer_index_remap,当前镜像就不支持通过这条调用链启用它,需要升级兼容的 Lance/Daft;同理,io_buffer_size不是 Compact 参数时也不能写入compaction_options

说明

不要让 Compact 与 Index 构建同时修改同一个 Dataset,除非已经确认当前软件版本的冲突处理与 Fragment Reuse 能力。并发写入冲突可能让已完成的工作重新执行。

Version、历史 Data File 与 Cleanup

Compact 成功后,Lance 发布引用新 Fragment 和新 Data File 的 Version。为了支持 MVCC、Time Travel 和正在读取旧快照的客户端,旧 Version 及其引用的旧 Data File 不会立即删除。

Compact 与 Cleanup 的边界

  • compact_files():改善当前 Fragment 布局,并创建新 Version。
  • cleanup_old_versions():删除超过保留期限的旧 Version,并回收不再被任何保留 Version 引用的 Data File。

只有当旧 Data File 不再被当前 Version、保留的历史 Version 或受保护标签引用时,Cleanup 才能安全回收它。Cleanup 会缩短 Time Travel 窗口,并可能使回滚不可恢复。

推荐顺序

  1. Compact 并确认新 Version 的 Commit 成功。
  2. 打开最新 Version,校验 Schema、行数、业务抽样、Fragment 和 Index 状态。
  3. 保留足够的回滚窗口,等待长时间读者结束。
  4. 再独立执行 cleanup_old_versions()

说明

不要为了立即释放空间而把保留期设为 0,并同时使用激进的未验证文件删除选项;并发写入或尚未提交完成的操作可能仍需要这些文件。

选择执行方式

执行方式

适用场景

调度与 Commit 特点

单机 Lance

数据量和维护窗口允许在单进程内完成

接口直接,但 Worker 的 CPU、内存和 I/O 都受单机资源限制

Lance-Ray

希望使用 Lance 官方 Ray 分布式工作流

通过 Ray Worker 分发 Compaction Task,主要暴露 Worker 数和 Ray 资源参数

Lance Spark Connector

已有 Spark SQL 或批处理平台,希望通过 OPTIMIZE 维护 Lance 表

通过 Lance Spark SQL Extension 暴露 OPTIMIZE;Lance 负责 Compact 语义,作业资源与执行行为受 Spark Driver、Executor 和集群配置影响

Daft 集成

已有 Daft/Ray 环境,希望用 Daft 调度 Lance Compaction Task

Lance 负责 Compaction Plan、Rewrite 与 Commit;Daft 负责 Daft Partition、调度和 RewriteResult 收集

无论使用哪种执行器,数据格式和 Compact 语义仍由 Lance 决定。执行器主要改变 Compaction Task 在哪里运行、同时运行多少个,以及 RewriteResult 如何送回 Driver。

使用 Daft 分布式执行:参数与代码示例

先看执行代码

from daft.io import lance

lance.compact_files(
    DATASET_URI,
    io_config=io_config,
    compaction_options={
        "target_rows_per_fragment": 2_000_000,
        "materialize_deletions": True,
        "max_source_fragments": 100,
    },
    num_partitions=200,
    concurrency=80,
    micro_commit_batch_size=20,
)

这些数值只是示例,不是通用推荐。是否能安全运行,取决于每个 Worker 可容纳的单个 Compaction Task 峰值、同机工作单元数量,以及每个 Lance 进程内部的 I/O 与 CPU 并行。

代码内部如何执行

  1. Driver 打开 Lance Dataset,并调用Compaction.plan()生成完整 Compaction Plan;Compaction Plan 中包含若干 Compaction Task
  2. Daft 按micro_commit_batch_size从完整 Compaction Plan 中取出当前 Commit Batch
  3. Daft 把当前 Commit Batch 中的 Compaction Task 分配到num_partitionsDaft Partition
  4. Daft Actor 调用CompactionTask.execute(dataset)执行 Rewrite;concurrency限制同时运行的工作单元数量。
  5. df.to_pandas()把本批次的分布式 RewriteResult 收集成 Driver/Head 进程中的单机 pandas DataFrame。它方便后续 Commit,但会让本批次结果在 Driver/Head 侧物化。
  6. Driver 调用Compaction.commit()提交当前 Commit Batch,生成新的 Version,然后处理下一批。

Image

三个 Daft 参数分别控制什么

参数

准确语义

过小 / 过大的影响

num_partitions

当前 Commit Batch 中的 Compaction Task 被划分成多少个 Daft Partition;不是 Compaction Task 数

过小会限制可并行调度的工作单元并放大长尾;过大增加调度和 Actor 管理开销

concurrency

Daft 最多同时执行多少个工作单元

过小可能用不满资源;过大使同机内存、网络和对象存储请求同时增加

micro_commit_batch_size

一个 Commit Batch 最多包含多少个 Compaction Task

过小增加 Commit 和 Version 数;过大增加 Driver/Head 的 RewriteResult 汇总与 Commit 峰值

在当前 Daft 实现中,num_partitions未设置时有效值为 1,并且不会超过 Compaction Task 总数;micro_commit_batch_size未设置时等于 Compaction Task 总数,即整个 Compaction Plan 只形成一个 Commit Batch。日志中“默认显示 80”通常意味着该 Compaction Plan 恰好包含 80 个 Compaction Task,而不是固定默认值为 80。

Lance 进程级并行

Daft 控制跨进程、跨节点的调度;Lance 仍在每个 Actor 进程内完成扫描、解码、编码和对象存储 I/O。因此总压力可能呈乘法关系:

全集群潜在 I/O 并行 ≈ Daft 并发 Actor × 每个 Lance 进程的 I/O 并行

Daft 控制跨节点的 Actor 并发,Lance 控制每个 Actor 内部的 I/O 并行。两层并发会相乘。

环境变量必须进入 Ray Worker

这些变量只影响真正执行 Lance 读取的进程。在 Driver 中执行 os.environ,并不能修改已经启动的 Ray Worker。应通过作业 Runtime Env、容器环境或集群级环境变量传递,并确保 Worker 在启动时已经读取到这些值。如果在 Python 中设置,必须在初始化 Ray 和创建 Worker 之前完成。

说明

这些变量用于控制 Worker Rewrite 阶段,不属于 CompactionOptions。因此不要把 io_buffer_size写入 compaction_options;Lance 7.0.0 会将其判定为无效 Compact 参数。

Worker OOM、Head OOM 与 Ray Object Store

  • Worker OOM:​通常发生在 Rewrite 阶段,来源是单个 Compaction Task 的 Buffer、同机多个工作单元、I/O 预读和进程内并行。
  • Head OOM:​通常发生在 Worker 完成后,来源是 RewriteResult 汇总、事务构造、Row Address 映射、Remap,或一个 Commit Batch 包含过多 Compaction Task。
  • Ray Object Store OOM:​只涉及 Ray 对象存储空间;Arrow/Rust/Python 进程 RSS 超限属于普通内存 OOM。两者不能混为一谈。

Image

使用 Lance 原生参数做增量 Compact

对于超大 Dataset,可以用 Lance 原生max_source_fragments限制一次调用生成的 Compaction Plan 大小。它先限制本次涉及的源 Fragment;Daft 随后执行这些 Compaction Task,并可继续用micro_commit_batch_size把结果拆成多个 Commit Batch。两者作用阶段不同,不冲突。

在 Daft 中传入 Lance 原生参数

max_source_fragments应放在compaction_options中。它由 Lance 在生成 Compaction Plan 时读取,不需要自定义选择 Compaction Task 的 Python 控制器。

lance.compact_files(
    DATASET_URI,
    io_config=io_config,
    compaction_options={
        "target_rows_per_fragment": 2_000_000,
        "materialize_deletions": True,
        "max_source_fragments": 100,
    },
    num_partitions=200,
    concurrency=80,
    micro_commit_batch_size=20,
)

说明

上面的 100 和 20 只是为了说明两个参数的关系,不是通用推荐值。一次调用只处理当前 Compaction Plan 中的 Fragment;若max_source_fragments截断了 Compaction Plan,需要再次调用 Compact 才会继续处理剩余 Fragment。

max_source_fragmentsmicro_commit_batch_size如何同时生效

两者不矛盾,并且可以同时设置。max_source_fragments先按源 Fragment 数限制 Lance 本次生成的 Compaction Plan;micro_commit_batch_size再按Compaction Task 数把这个 Compaction Plan 拆成多个 Commit Batch。一个 Compaction Task 可能包含一个或多个源 Fragment,因此两个数不能直接换算。
例如 Lance 原本可生成三个 Compaction Task,分别包含 3、4、5 个源 Fragment。设置max_source_fragments=10后,本次 Compaction Plan 只纳入前两个 Compaction Task,共 7 个源 Fragment;第三个加入后会达到 12,因此留给下一次 Compact。若同时设置micro_commit_batch_size=1,Daft 会把这两个 Compaction Task 分成两个 Commit Batch,分别 Commit。
Image

参数

生效阶段

计数单位

直接控制什么

max_source_fragments

Lance 生成 Compaction Plan 时

当前 Compaction Plan 涉及的源 Fragment 数

本次 Compact 最多规划多大的 Fragment 范围

micro_commit_batch_size

Compaction Plan 生成之后,由 Daft 分批 Commit 时

一个 Commit Batch 包含的 Compaction Task 数

Driver 每次汇总多少个 RewriteResult 并执行一次 Commit

max_source_fragmentsmicro_commit_batch_size都可以单独设置,也可以同时设置。前者不限制时,Lance 会规划全部符合条件的 Fragment;后者不设置时,Daft 默认将当前 Compaction Plan 中的 Compaction Task 作为一个 Commit Batch。若当前 Compaction Task 数小于等于micro_commit_batch_size,本次仍然只会 Commit 一次。

性能慢与 OOM 的定位方法

现象

更可能的阶段

如何理解和检查

进度未完成,某个 Worker 被杀并反复重试

Rewrite

先看该 Worker 的进程内存、同机 Actor 数、batch、预读、I/O Buffer,以及单个 Task 涉及的 Fragment 和数据量。

Worker 很快达到 N/N,随后 Head 长时间无进展或 OOM

结果汇总 / Commit

检查 micro_commit_batch_sizeto_pandas()产生的结果、RewriteResult 体积、Row Address 映射和索引处理。

Ray Dashboard 中节点或进程的 Memory 很高,但 Object Store Memory 很低

Worker 或 Driver 的进程内存

说明内存主要被 Lance、Arrow 或 Python 进程占用,而不是被 Ray Object Store 占用。先找出哪个进程的 RSS 持续增长,再结合日志判断发生在 Rewrite 还是 Commit。

CPU 较低、网络较高,增加并发后吞吐不再提升

对象存储或网络瓶颈

检查请求延迟、对象存储限流、重试和集群总 I/O 并行。

提高 LANCE_IO_THREADS 后反而更慢

I/O 过度并行

检查 Actor 数与每个进程 I/O 线程数的乘积,以及对象存储限流、排队和额外内存。

为什么进度会卡在固定百分比

如果剩余 Compaction Task 总在少数节点 OOM,Ray 可能不断重试这些 Compaction Task,进度会长期停在同一百分比。另一种情况是所有 Rewrite 已完成,进度条不再变化,但 Head 正在汇总 RewriteResult 或执行 Commit。需要用日志时间线区分“Compaction Task 重试循环”和“Commit 阶段”,不能只凭进度数字判断。

为什么提高并发反而变慢

当网络、对象存储请求率或单节点 CPU 已饱和,继续增加 Actor 或LANCE_IO_THREADS只会增加排队、限流和上下文切换。提高target_rows_per_fragment还可能改变 Compaction Plan,扩大单个 Compaction Task 或总 Rewrite 字节数。因此变慢时,应分别比较 Compaction Task 数、Rewrite 字节数、网络吞吐和 Commit 耗时,而不是只看参数是否变大。

说明

一次调优应保留 Compaction Task 数、源/目标 Fragment 数、Rewrite 字节数、峰值内存、网络吞吐和 Commit 耗时。没有这些基线,只比较总耗时很容易得出错误结论。

最近更新时间:2026.08.06 14:47:29
这个页面对您有帮助吗?
有用
有用
无用
无用