Repository navigation
Replies: 1 comment
|
Thanks for starting this discussion! I think there are 4 prerequisites for supporting liquid clustering at format level, which as of today we only have implemented 1.5:
Outside that, I don't think we want to add the actual clustering implementation in lance, because this is clearly going beyond single-node execution. Clustering would require a large scale distributed sorting according to the hilbert curve, so I think the best way to do that is to add the 4 format features above in Lance, and we choose an engine, spark or ray (for anyone interested feel free to pick more 😆) and get the clustering implemented via the engine. I would say let's do it in spark, since people in lance-spark have wanted storage partitioned join for a long time, and Spark is likely the best one to get that supported along the way and it has enough hooks to plugin related information. @yanghua what do you think about this? |
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Related issues: #1434 (space-filling curve
cluster_bywrite param), #1045 (EPIC: statistics and data skipping), #952 (implicit partitioning, closed not-planned), #6803 (recluster task for row-id healing).1. Background and motivation
Lance today has the read-side and write-side building blocks for data skipping, but not the layer
that would tie them into a clustering feature:
rows_per_zonerows and prunes zones at scan time (ZoneMapIndexinrust/lance-index/src/scalar/zonemap.rs). On append to a current-format dataset, an existing zonemap index with seeding enabled can already be seeded inline throughIndexSeedWriter.WriteParamshas no notion of clustering or sort order. Any clustering must be arranged by the caller before data reaches the writer. In the Spark connector,PARTITIONED BYapproximates single-key clustering: Spark shuffles/sorts by the key, andLanceDataWriterrolls a fresh fragment whenever the key changes.compact_filesandplan_compactiondeliberately do not reorder data — they use the ordinary planner and try to preserve insertion order. There is no reorder-aware compaction entry point today.HilbertSorterinrust/lance-index/src/scalar/rtree/sort/hilbert_sort.rs). It is not a general multi-column encoder.The core of liquid clustering — re-clustering that re-sorts data by a multi-column key — cannot be implemented in a connector (Spark/LanceDB) alone, because compaction lives in the Lance core. We therefore propose to implement it primarily in the Rust core, with thin binding/connector surfaces on top.
Why not "just Z-ORDER + partitioning"
Legacy
ZORDER BY+ Hive-style partitioning requires the user to pick partition columns up front, suffers small-file / skew problems, and needs full rewrites to change layout. Liquid clustering would replace both: clustering keys are declared once and can change without an immediate rewrite. Normal append and overwrite writes would inherit the active declaration, while an explicit clustering compaction would bring older fragments, and fragments changed by other write paths, into the currentlayout.
2. Goals and non-goals
Goals
Dataset::set_clustering), persisted in table metadata. SQLCLUSTER BY (a, b, ...)connector support is future work.OPTIMIZE-driven task that picks up under-clustered data and merges it into the clustered layout. Optional source budgets should be able to bound the amount of data touched per run.Non-goals (for the first iteration)
3. Proposed design overview
We propose three cooperating pieces, all in the Rust core:
namespace is atomic: no key in that namespace means the declaration is absent, while any key in it requires all five
known keys to be present and valid in the resulting manifest. Partial or malformed declarations are rejected. A
clustering feature flag gates readers and writers that predate this contract.
SpaceFillingEncoderwould turn N key columns into a single 1-D ordering value using a space-filling curve. Sorting by this value gives multi-dimensional locality.4. Clustering key encoder
We propose a new clustering module in the
lance-indexcrate that does not reuse the geoHilbertSorter(which is 2D and bbox-specific). ASpaceFillingEncoderwould map N key columns into a single fixed-size binary ordering value; aClusteringCurveselector would choose Z-order or Hilbert encoding.The initial encoder would accept only top-level columns with these Arrow types:
Boolean,Int8/Int16/Int32/Int64,UInt8/UInt16/UInt32/UInt64,Float32, andFloat64.Strings, binary, dictionaries, temporal values, decimals, nested values, and other types would be rejected.
Proposed encoding behavior:
The production write and compaction path obtains both passes from a memory-first replay spill, so the empirical model is task-wide without collecting the full input in memory. This avoids the severe precision collapse that fixed-domain high-bit truncation causes for common narrow
Int64and floating-point ranges. The standaloneSpaceFillingEncoderretains fixed-domain encoding when no fitted model is supplied. Null maps to the maximum coordinate in its own dimension; a multi-dimensional curve does not promise row-levelNULLS LAST.Precision/whitening (per-column bit width, handling skew) is an open question — see §10.
5. Format / metadata changes (additive, no stable-format break)
The complete desired layout would be stored in
manifest.config. It should deliberately not reuse the existinglance-schema:unenforced-clustering-key:positionmarker: that stable marker asserts an already-achieved physical ordering for query-engine optimizations, while liquid clustering is a policy to which existing fragments may only converge incrementally.lance.clustering.columnslance.clustering.curvehilbert|zorder)lance.clustering.versionlance.clustering.bits_per_dimRationale:
configis already an additiveHashMap<String,String>mutated through the existingUpdateConfigoperation (rust/lance/src/dataset/metadata.rs), so changing the desired layout would be a cheap metadata-only commit. The complete declaration would be updated atomically. Any change to columns, curve, or bit width would have to increase the layout version so existing fragments become eligible for incremental reclustering. Increasing only the version would also be allowed and wouldforce all fragments with an older stamp to become eligible again. Clearing the declaration would retain existing fragment stamps; re-enabling clustering would therefore need a version greater than the maximum retained stamp. Every manifest build would validate the complete declaration against the resulting schema, so a schema change could not leave an active clustering column missing, nested, or of an unsupported type. Manifests carrying the declaration would set both a clustering reader and writer feature flag, so older implementations do not open the dataset while ignoring the new invariants.
Per-fragment clustering state. To make re-clustering incremental, the implementation records both the active declared clustering
versionand aclustering_group_idshared by every fragment produced by one sorted write or compaction task. A fragment with an absent, older, or unexpected future version is under-clustered. A current-version group whose aggregate live-row count is below the target remains eligible to merge with an adjacent partial group. A lone partial group is left alone until another candidate arrives, avoiding a rewrite loop. A one-shot sort on an undeclared dataset remains unstamped and ungrouped.The scalar
DataFragment.clustering_versionfield (0 = unset) and stringDataFragment.clustering_group_idfield (empty = unset) persist this state. In Rust, decoded versions and group IDs are kept in privateManifestsidecars aligned positionally withManifest::fragments; they are not exposed as fields on the publicFragmenttype. A reader-and-writer feature flag fences older readers and writers whenever an active clustering declaration or any fragment clustering state is present.Because the proposed field would be a stable scalar encoding whose zero value means "unset", active clustering versions must be positive, and that zero sentinel could not later be reinterpreted as a valid version. A future format needing explicit presence or version zero would require a new field or encoding.
A future quality heuristic could additionally infer overlap from zonemap ranges, but it would not replace the version stamp as the authoritative current-layout predicate.
6. Write path
We propose to keep clustering out of
WriteParams. Normal dataset writes would accept explicit column names through a newInsertBuilder::with_cluster_by_columns; multi-fragment writes could useFragmentCreateBuilder::with_cluster_by_columns. The builders would resolve those columns to a privateClusteringSpecpassed separately to the internal writer.Proposed write behavior:
InsertBuilderwould resolve explicit columns, or inherit the dataset's declared spec when columns are omitted, for both append and overwrite. Overwrite would inherit the declaration because overwrite commits preserve dataset config. On a dataset with an active declaration, explicit columns must exactly match it. On a new or undeclared dataset, explicit columns would request a one-shot sort and would not declare or stamp a persistent layout. Low-level multi-fragment writes would have to pass explicit columns throughFragmentCreateBuilder; omitting them would not inherit the declaration. Update and merge-insert would intentionally omit the clustering spec: they would preserve stamps only for fragments whose layout metadata is unchanged, while changed or new fragments would be left unstamped for a later clustering compaction. A note on transport: a public fragment-write API that returns onlyVec<Fragment>cannot carry the private manifest-sidecar version to a later manual commit. We would need a hidden transport (e.g. awrite_fragments_with_clustering_versionvariant) for distributed callers, and Python could instead write through an uncommitted-stream path, extract the reserved version from the resulting transaction, and carry it alongside exported fragment metadata. Any caller that separates fragment writing from commit must preserve the transaction's reserved clustering-version marker; otherwise the sorted fragments would be committed without a stamp.cluster_sort_streamhelper: aSortExecover a hiddenFixedSizeBinaryordering column, spilling to disk for large inputs. Oversized input batches would be split and deep-copied before they reach the sort, because DataFusion cannot spill a single batch that already exceeds its memory reservation; a single row larger than the cap would be rejected. For the streaming writer this is a sort within the write unit; global ordering across concurrent writers would not be guaranteed (that is what re-clustering converges).IndexSeedWriterpath can seed an already-declared zonemap index when that index has seeding enabled. Clustering maintenance creates and trains a missing seed-enabled zonemap for each clustering column and fills missing fragment coverage. This deferred creation keeps the declaration update metadata-only while ensuring that the first clustering optimize closes the data-skipping loop.Connectors would keep their current role: Spark's
RequiresDistributionAndOrderingcan still pre-sort at the engine for scale; the core sort would be the correctness backstop when the engine does not.7. Compaction / re-clustering
This is the part that must live in core because ordinary
compact_filesdeliberately does not reorder. We propose to keep the ordinarycompact_filesandplan_compactionentry points order-preserving and to leaveCompactionModeunchanged (the nonbreaking three-variant enumReencode,TryBinaryCopy,ForceBinaryCopy). Reordering would be exposed through separate, newly introduced APIs instead of being folded into the existing ones.Dedicated clustering operation. We propose a
ClusteringCompactionPlannerand aplan_clustering_compactionentry point that produce aClusteringCompactionPlan. Its tasks would be executed asClusteringCompactionTaskvalues, which resolve the declaration at the task's immutable read version and re-sort rows through a private clustering rewrite strategy before writing fragments viacluster_sort_stream. Each task would yield aClusteringRewriteResult; acommit_clustering_compactionstep would validate and commit those results, then record the task's clustering version after final fragment IDs are assigned. Clustering would require the existingCompactionMode::Reencodebehavior and would not use binary copy. Output would still respect thewriter semantics of
target_rows_per_fragment/max_bytes_per_file: the row target would be passed asWriteParams::max_rows_per_file, while the byte limit would remain the writer's existing soft limit.A single-process convenience entry point,
compact_files_with_clustering, wraps that same flow: plan withplan_clustering_compaction, execute the dedicated tasks, and commit withcommit_clustering_compaction. Distributed callers use those three stages directly. Both the no-op path and clustering commit path create missing clustering-column zonemaps and refresh their fragment coverage. If post-commit zonemap maintenance fails, the data rewrite remains valid and uncovered fragments continue through the scanner's normal fallback.Incremental planner. The dedicated
ClusteringCompactionPlannerselects under-clustered fragments — those whose per-fragment version does not equal the dataset's currentClusteringSpec::version— plus adjacent current-version groups whose aggregate live-row count is still below the target. It groups adjacent whole source fragments until the accumulated live rows reach or exceedtarget_rows_per_fragment, and honors optionalmax_source_fragments,max_source_rows, andmax_source_bytesbudgets. Because source fragments would not be split during planning, a task could exceedtarget_rows_per_fragment; the option would be separately forwarded to the output writer asmax_rows_per_fileduring rewriting. All three source budgets would default toNone, so clustering compaction would be unbounded by source volume unless a caller or dataset config sets at least one. When several are set, all would be hard upper bounds: planning stops before the next eligible fragment would exceed any one. Rows mean live rows; bytes include source data and overlay files but exclude separately stored Blob payloads. If the first eligible fragment exceeds a configured budget, the plan is empty and the budget must be raised. A stable current-version group breaks adjacency. A single partial group is also a no-op until another adjacent candidate arrives, so repeated optimize runs do not rewrite the same group forever.Distributed operation separation. We propose that
ClusteringCompactionTaskandClusteringRewriteResultuse dedicated, version-tagged serialized envelopes rather than the ordinaryCompactionTaskandRewriteResultpayloads. The tags would prevent accidental cross-routing between order-preserving and row-reordering work. Task execution would check out the declared read version and validate the clustering declaration and clustering-specific restrictions at that snapshot. The initial design does not require comparing caller-providedTaskDatafragment descriptors with the manifest's fragment descriptors at that version;commit_clustering_compactionwould require one read version and clustering version across all results and check that version against the declaration before reserving fragment IDs, but would not verify complete original-fragment identity or that output fragments preserve the input live-row count. The serialized task/result boundary would therefore be a trusted-worker boundary, not an attestation of the referenced files or rewrite contents. Distributed callers would be required to provide valid, non-empty fragment descriptors corresponding to the stated snapshot and preserve the tagged execution results unchanged. Hardening this boundary (descriptor validation against the snapshot, rejecting empty tasks/results, row-count preservation checks) is called out in §10 and would be part of the delivery plan rather than deferred indefinitely.Not yet implemented: using persisted key or curve ranges to pull in non-adjacent or otherwise stable current-version
groups whose value ranges overlap new data. The current group-size policy handles adjacent partial writes but does
not yet measure range overlap.
Changing keys without full rewrite. Bumping the clustering version while changing the columns or other layout parameters would mark all existing fragments as under-clustered lazily; they would be re-clustered by later
OPTIMIZEruns. With a source budget configured, a large table could converge over several bounded runs. With the default unbounded source budgets, one run could select all eligible fragments. Old data would stay readable throughout.Interaction with existing compaction machinery. Ordinary compaction preserves row order, so its positional old→new row mapping (for index remap and stable-row-id rechunk) holds. Reclustering reorders rows, which breaks that positional assumption. Rather than silently corrupt row ids or a secondary index, the dedicated clustering path currently rejects datasets that use stable row ids or carry a remappable secondary index other than zonemap (drop the index, recluster, rebuild), as well as
defer_index_remap=true. Carrying row identity through the sort and rebuilding the mapping from the sorted order is the follow-up that would lift this restriction.8. Read path
No new pruning mechanism is introduced. Clustering maintenance ensures every clustering column has a zonemap, creating missing declarations with seeding enabled and refreshing coverage after rewritten fragments are committed. Existing scan-time pruning then exploits the resulting value locality. Before the first clustering optimize, a newly declared clustering spec may still have no zonemap; this keeps declaration itself a cheap metadata-only operation. Query-planning exposure for cost estimates and connector pushdown remains future work.
9. Proposed API surface (thin bindings)
We would centralize resolution and validation in Rust while keeping the binding-level names consistent (
cluster_by,clustering). All names below are proposals subject to review.InsertBuilder::with_cluster_by_columnsandFragmentCreateBuilder::with_cluster_by_columnsfor explicit write columns;Dataset::set_clustering/Dataset::clustering_spec/Dataset::clear_clusteringto declare/read/drop the config-backed spec. Reclustering would useClusteringCompactionPlanner,plan_clustering_compaction,ClusteringCompactionTask::execute,commit_clustering_compaction, or the single-processcompact_files_with_clusteringhelper.CompactionModewould gain no clustering variant; a separateCompactionRequest::Clusteringdiscriminator would carry options normalized toCompactionMode::Reencode.WriteParamswould contain no clustering field. Since a publicFragmentCreateBuilder::write_fragmentsreturningVec<Fragment>cannot expose the declaration version needed by a separate manual commit, a hiddenwrite_fragments_with_clustering_versiontransport would exist for binding transports.write_dataset(..., cluster_by=["a", "b"]);dataset.set_clustering(columns, *, curve, version, bits_per_dim)/dataset.clustering_spec() (returns a dict or None) / dataset.clear_clustering();dataset.optimize.compact_files(compaction_mode="cluster"). The string"cluster"is a binding-level request discriminator, not a RustCompactionMode: Python normalizes the core mode toReencodeand dispatches execute/plan/commit to the dedicated clustering APIs. When a clustering declaration is active, omitting compaction_mode also selects clustering maintenance; an explicit non-clustering mode overrides that default.WriteParams.Builder.withClusterBy(List<String>)as the public Java binding surface, with JNI forwarding its column list to the Rust write builder;Dataset.setClustering(ClusteringSpec)/Dataset.getClusteringSpec()/Dataset.clearClustering(); and aCompactionMode.CLUSTERvalue. Java'sCLUSTERwould likewise be a binding-level request discriminator: the distributedCompactionAPI would pass it to JNI, which maps it toCompactionRequest::Clustering, normalizes the core mode toReencode, and uses the dedicated clustering plan, tagged task/result, and commit path. It would not be a RustCompactionModevariant; the JavaDataset.compactconvenience method dispatches the same binding-level value tocompact_files_with_clustering. As in Python, an active clustering declaration selects clustering maintenance when no mode is specified, while an explicit non-clustering mode overrides it.ClusteringSpec/ClusteringCurveare thin value types mirroring the Rust/Python shape.CLUSTER BY (a, b)to the persisted spec;OPTIMIZEtriggers the incremental recluster; keepLanceScanBuilderpruning as-is. (Connector lives outside this repo.)10. Open questions
rows_per_zonealignment. Clustering quality is only useful if zone granularity is fine enough; should the default zone size be coupled to the clustering config?11. Proposed phased delivery
ClusteringSpectype, the complete five-key declaration inlance.clustering.*config, andDataset::set_clustering/clustering_spec/clear_clustering.SpaceFillingEncoderinlance-index::clustering, with task-wide quantile-rank normalization, bounded deterministic sampling, spill-backed replay, and encoder tests.WriteParams, resolve the private spec from the dataset on append and overwrite, and sort withcluster_sort_stream. Commits stamp fragments governed by the active declaration with version and group state. An explicit one-shot sort on a new or undeclared dataset is not stamped. Update and merge-insert do not inherit the sort; changed/new fragments from those paths are left under-clustered. Low-level Rust fragment writes must preserve the separately returned clustering-version transport when committing if their outputs are to be stamped.ClusteringCompactionPlanner/plan_clustering_compactionpath selects under-clustered fragments and emitsClusteringCompactionTaskvalues whose execution reorders rows by the clustering key. It also selects adjacent current-version partial groups and honors configured source budgets, which remain unbounded by default.ClusteringRewriteResultvalues are committed throughcommit_clustering_compaction; version and group state are persisted throughDataFragment.clustering_versionandDataFragment.clustering_group_id. Missing clustering-column zonemaps are created and zonemap coverage is refreshed after either a rewrite or a no-op clustering optimize. Deferred work includes stable-row-id support, remappable non-zonemap indices, overlap-aware group selection, and stronger distributed task/result validation.cluster_byon the write path,set_clustering/clustering_spec/clear_clusteringfor declaration, and the"cluster"/CLUSTERbinding-level request discriminators. An active declaration also selects clustering maintenance when optimize mode is omitted. Deferred: SparkCLUSTER BY+OPTIMIZEintegration, which lives outside this repo.13. Alignment with upstream
cluster_bywrite work tracked by Implement space-filling curvecluster_bywrite param #1434 and advance the [EPIC] Statistics and data skipping #1045 data-skipping EPIC (zonemap consumption).All reactions