Lance Query Best Practices
Use this guide to choose a per-BE cache budget and index segment layout for your Lance queries. Start with your workload, estimate the content that needs to stay cached, and validate it under the intended concurrency. The examples are planning starting points, not benchmarked optima or BE memory limits.
Lance Catalog is experimental and supported starting from Apache Doris 4.2. First configure access using Lance Catalog and check its reader compatibility. Doris reads existing Lance indexes; use compatible Lance tooling to build or maintain them.
For scalar filters combined with vector or full-text retrieval, see Lance Hybrid Search Performance Best Practices for filter placement, aligned index construction, and profile diagnosis.
1. Choose a Tuning Path
| Your workload | Start with | What to validate |
|---|---|---|
| Vector search with low single-query latency | Retain frequently probed vector partitions; compare balanced segment counts relative to eligible BEs using the segment guide | P95 latency at comparable recall; slowest scanner; peak memory |
| High-concurrency vector search | Budget for the union of hot partitions across different query vectors; leave memory and CPU for concurrent searches instead of maximizing splits for each query | Throughput and P95 latency at target concurrency, cache reloads, per-BE memory |
| BTree range/point filters or Bitmap/LabelList filters | Estimate hot pages or bitmap entries, including lookup structures; verify predicate/index use | Page loads, selectivity, temporary row-ID sets, and output-column reads |
| Full-text or phrase search | Estimate dictionaries, document metadata, postings, and positions when enabled; include diverse terms in warmup | Common-term queries, phrase-query memory, and candidate merging |
| Vector, scalar, and full-text queries sharing BEs | Add their shared working sets and test them together | Whether one query family evicts another's entries |
| Output-column fetches or scans dominate | Inspect data-file reads and the separate disk data cache | Useful bytes versus read amplification; whether decoding or fetches dominate |
A scalar filter may also be part of a vector query's Prefilter, so a single query can need both scalar and vector index content. Choose all applicable rows rather than treating the scenarios as mutually exclusive.
2. Collect Inputs and Set a Memory Budget
Collect these inputs for each BE or for a representative BE under realistic placement:
- Workload: target QPS/concurrency, latency target, required recall,
top_k, filters, phrase queries, and output columns. - Resources: BEs eligible for the queries, CPU capacity, and memory available after other Doris workloads and process/OS needs.
- Indexes: types, versions, physical segment counts, coverage, vector dimensions/encoding, IVF partition counts, and scalar value widths or string lengths.
- Reuse: distinct partitions, pages, terms, or bitmap entries accessed across a representative query window. Use the working set across queries, not one query's returned rows.
- Query peaks: memory for active partition loads, Prefilter/result row-ID sets, candidate heaps, decoding, refinement, and second-phase fetches at the intended concurrency.
Use the index formulas to estimate payload, then check that the combined budget fits:
available_Lance_cache_budget_on_BE
= BE memory budget
- other BE memory needs
- peak concurrent query working memory
- safety headroom
index_cache_target + metadata_cache_target
<= available_Lance_cache_budget_on_BE
This is a planning check, not an enforced memory reservation. Cache capacities do not cap query memory or process RSS. Do not multiply a reusable cached working set by the number of concurrent queries; do reserve for their additional query working memory, and account for a larger union of cached content when concurrent queries access different data.
The full parameter definitions and defaults remain in BE Configuration. Set the lance_* BE parameters in each BE's be.conf and restart that BE. They are not catalog properties or SQL session variables. Index segment counts and IVF construction parameters are set through Lance index-building tooling, not these cache parameters.
3. Estimate the Index Working Set
Size the index cache for the index partitions that each BE needs to reuse. The total index on storage, the cached working set, and peak query memory are different quantities. The estimates below cover the six vector index types supported by Doris; they are planning estimates, not exact cache accounting or memory limits.
Vector Payload by Index Type
Use the following symbols. All formulas return bytes; 1 GiB = 1073741824 bytes.
| Symbol | Meaning |
|---|---|
N | Number of indexed vectors in the partitions being estimated. For an ordinary single-vector column this is the indexed row count. For multi-vector columns, count the indexed subvectors and allow separately for row mappings and refinement. |
D | Vector dimension. |
w | Bytes per stored, unquantized vector element: 2 for Float16, 4 for Float32, or 8 for Float64. For binary uint8 vectors, use the actual byte length as D and w=1. |
m, b | PQ subvector count and bits per subvector code. These are index-build settings, not query parameters. |
G | Memory occupied by the loaded HNSW graph, including all levels and auxiliary structures. |
Add all applicable components, not just the vector codes. The first table separates the per-vector payload; the following tables list shared structures and HNSW graph components. All formulas estimate loaded buffers. Rows marked Measure still need a budget; they are not zero-cost items. Count each shared allocation once.
| Index type | Vector values / codes | UInt64 row IDs | HNSW graph | Other required budget |
|---|---|---|---|---|
IVF_FLAT | N * D * w | 8 * N | None | IVF structures + common overhead below |
IVF_SQ | N * D for SQ8 | 8 * N | None | IVF structures + SQ metadata + common overhead |
IVF_PQ | N * ceil(m * b / 8) | 8 * N | None | IVF structures + PQ codebook + common overhead |
IVF_HNSW_FLAT | N * D * w | 8 * N | Add G | IVF structures + common overhead |
IVF_HNSW_SQ | N * D for SQ8 | 8 * N | Add G | IVF structures + SQ metadata + common overhead |
IVF_HNSW_PQ | N * ceil(m * b / 8) | 8 * N | Add G | IVF structures + PQ codebook + common overhead |
For PQ with b=8, the code occupies m bytes per vector; with b=4, it occupies m/2 bytes for supported, even m. For example, 1 million vectors with m=96 and b=8 have 96000000 bytes of codes and 8000000 bytes of row IDs: 104000000 bytes (about 99.2 MiB) combined. The row IDs must not be omitted when estimating strongly compressed indexes.
| Additional component | Applies to | Loaded payload / how to budget |
|---|---|---|
| IVF centroids | All six types | FP32: num_partitions * D * 4 per retained centroid set; use the actual element width for other layouts |
| IVF partition directory / sub-index metadata | All six types | Measure retained partition offsets, lengths, and reader structures; distinct from vector values and centroids |
| PQ codebook | IVF_PQ, IVF_HNSW_PQ | Standard FP32 codebook: 2^b * D * 4 per retained copy |
| SQ quantizer metadata | IVF_SQ, IVF_HNSW_SQ | Measure retained quantizer bounds/parameters; do not assume one fixed layout across versions |
| Norms, row-coverage structures, additional mappings | Where retained by the reader | Measure actual arrays/sets; multi-vector row mappings need their own budget |
| Index objects, Arrow validity buffers, spare capacity and copies | All six types | Measure allocations beyond the logical payload; count shared buffers once |
Independent index segments can retain separate centroid sets and codebooks. Do not multiply shared allocations by every partition or assume all copies are shared. The extra structures depend on the embedded Lance version.
HNSW Graph Components
The graph is additional to vector codes and row IDs. For an Arrow-backed Lance graph with UInt32 node/neighbor IDs and optional retained Float32 edge distances, let E be the stored directed neighbor entries across all loaded levels, V the node occurrences across those levels, and A the number of independently retained graph batches. A node present at several levels contributes several times to V.
| Graph component | Main payload, in bytes | Counting rule |
|---|---|---|
| Neighbor IDs | 4 * E | UInt32 per directed neighbor entry; count all levels |
| Stored edge distances | 4 * E | Float32 per neighbor entry; zero when the reader skips the distance column |
| Graph node IDs | 4 * V | UInt32 per node occurrence; these are distinct from dataset UInt64 row IDs |
| Neighbor-list offsets | 4 * (V + A) | Int32 list offsets, including one terminal offset per batch |
| Distance-list offsets | 4 * (V + A) | A second Int32 list-offset array; zero when the reader skips the distance column |
| Subtotal with distances | 8 * E + 12 * V + 8 * A | Approximately 8 * E + 12 * V for large batches; not the complete G |
| Subtotal without distances | 4 * E + 8 * V + 4 * A | Use when only node IDs and neighbor lists are retained; select one subtotal, not both |
| Level offsets, entry points and node lookup structures | Measure | Include retained metadata and lookup allocations for each loaded graph |
| Validity buffers, graph objects and spare allocation capacity | Measure | Add actual allocations beyond the buffers above, without double-counting shared buffers |
Complete graph budget G | Main subtotal + measured graph extras | Add this to the vector/row-ID payload and other index structures |
This is a layout-specific estimate, not a fixed bytes-per-node guarantee. HNSW's m controls graph connectivity and is unrelated to PQ's m. Do not substitute N * m * 4 for G: actual degree, upper levels, retained distances, and the reader version matter. Calibrate by loading representative partitions and measuring their retained allocations.
For example, suppose the cached IVF_HNSW_PQ partitions contain 1 million vectors, PQ m=96, b=8, E=32000000, V=1100000, and A=1. Assume this reader retains the distance column. These graph counts are illustrative inputs, not values derived from HNSW's m:
| Component | Example payload |
|---|---|
| PQ codes | 96000000 bytes |
| Dataset row IDs | 8000000 bytes |
| Main graph buffers | 8 * 32000000 + 12 * 1100000 + 8 = 269200008 bytes, about 256.7 MiB |
| Codes + row IDs + main graph buffers | 373200008 bytes, about 355.9 MiB |
| Alternative: main graph buffers when distances are skipped | 4 * 32000000 + 8 * 1100000 + 4 = 136800004 bytes, about 130.5 MiB |
| Alternative total: codes + row IDs + graph buffers without distances | 240800004 bytes, about 229.6 MiB |
| Still to add in either case | Graph extras, IVF centroids/directory, PQ codebook and common overhead from the tables above |
Select the example matching the reader's retained columns. Neither subtotal is a complete cache-capacity recommendation.
The upstream Lance Performance Guide explains vector payload and index-build memory estimates. The Lance vector index format describes the separation of IVF partitions, sub-indexes, and quantized storage. Their latest content may describe features newer than the Lance version embedded in Doris; use the Lance Catalog support matrix for compatibility.
Scalar and Full-Text Index Payloads
Scalar indexes can also be estimated, but the inputs differ by index family. Count the entries actually retained on one BE, not just the rows returned by a query. The following methods describe Lance storage layouts; they do not extend the Doris predicate-pushdown or index compatibility guarantees. Confirm that the embedded Lance version supports the index and that the query uses it.
BTree. A loaded leaf page contains values and UInt64 row IDs. Let R be the number of records in the cached leaf pages, H the number of those pages, w the fixed-width value size, and L the average string length in bytes in those pages:
| Leaf-page component / representation | Value buffer | Offsets | UInt64 row IDs |
|---|---|---|---|
| Int32 / Float32 | 4 * R | None | 8 * R |
| Int64 / Float64 / 64-bit timestamp | 8 * R | None | 8 * R |
| Other fixed-width values | w * R; use actual bitmap buffers for bit-packed Boolean | None | 8 * R |
| Arrow Utf8 / Binary | R * L | 4 * (R + H) for Int32 offsets | 8 * R |
| Arrow LargeUtf8 / LargeBinary | R * L | 8 * (R + H) for Int64 offsets | 8 * R |
The leaf-page subtotal is the sum of the three payload columns. Add these components for each opened segment; P below counts all leaf pages in that segment, not only cached pages:
| Additional BTree component | Main payload / how to budget |
|---|---|
| Lookup minimum values | P * w for fixed-width values; actual value/offset buffers for variable-width values |
| Lookup maximum values | P * w for fixed-width values; actual value/offset buffers for variable-width values |
| Lookup NULL counts | 4 * P for the UInt32 layout |
| Lookup page numbers | 4 * P for the UInt32 layout |
| Fixed-width lookup-buffer subtotal | P * (2 * w + 8); excludes the structures below |
| NULL-page lists and lookup objects | Measure retained lists, objects and implementation-specific copies |
| Leaf/lookup validity buffers, Arrow objects and spare capacity | Measure allocations beyond logical values/offsets/row IDs |
String formulas assume the listed Arrow layouts; dictionary or view representations need their actual buffers. Disk compression does not determine decoded page size. The default BTree page size is 4096 records, not bytes; use the index's actual page size. The lookup belongs to the index working set, not the BE metadata-cache budget merely because it describes pages.
For example, 100 million Int64 values have 100000000 * 16 = 1600000000 bytes (about 1.49 GiB) of leaf value/row-ID payload. If cached pages cover 10 million records, that payload is about 152.6 MiB. With 4096 records per page, the full index has about 24415 pages; the lookup batch's min/max/count/page-number payload is about 572 KiB before extra structures. Returning ten matching rows does not imply caching only ten records: the reader loads pages.
Bitmap, LabelList and NGram. Sum the applicable component rows below. Membership bitmaps do not automatically cost eight bytes per matching row. LabelList counts membership per distinct label; a row can appear in multiple label bitmaps. NGram counts distinct gram/row memberships; repeated occurrences in one row do not add memberships. This NGram index differs from an FTS index using an n-gram tokenizer.
| Index family | Component | Loaded payload / how to budget |
|---|---|---|
| Bitmap | Key lookup | Measure retained key values, entry references and lookup objects |
| LabelList | Label lookup / wrapper | Measure distinct label keys, entry references and wrapper objects |
| NGram | Gram lookup | Measure retained grams, entry references and lookup objects |
| All three | Membership bitmap: array containers | About 2 * c per conventional Roaring array container with c UInt16 entries in a 65536-value range |
| All three | Membership bitmap: dense containers | 8192 bytes per conventional dense container covering a 65536-value range |
| All three | Membership bitmap: run / full-fragment representations | Measure actual representation; use the representation selected for each container, not all alternatives added together |
| All three | Bitmap directories / capacity | Measure fragment/high-bit directories, container objects and spare capacity |
| Where present | NULL / null-list tracking | Measure retained NULL membership structures separately |
| Each family total | Lookup + its loaded membership representations + directories + NULL tracking | Count shared allocations once; serialized bitmap size is not a resident-memory guarantee |
For example, 256 loaded dense containers have 2 MiB of bitset payload, before keys, directories, NULL tracking and object overhead.
Other scalar and full-text indexes. Each row identifies a separate component; subtotal rows summarize preceding components and must not be added again. Count only retained content on the BE. Measure means the reader's actual allocations are required, not that the component can be omitted.
| Index family | Component | Loaded payload / how to budget |
|---|---|---|
| FTS / Inverted | Term dictionary | Measure loaded terms, dictionary structure and lookup allocations |
| FTS / Inverted | Document metadata | Measure retained document mappings, lengths and other reader metadata |
| FTS / Inverted | Plain posting row IDs | 8 * T for T UInt64 term/document entries |
| FTS / Inverted | Plain posting frequencies | 4 * T for Float32 frequencies; IDs + frequencies total 12 * T |
| FTS / Inverted | Positions, when enabled | 4 * O for O stored Int32 occurrences |
| FTS / Inverted | Position offsets | Actual retained offset-array bytes; separate from 4 * O |
| FTS / Inverted | Compressed posting layout | Use actual retained compressed blocks instead of the corresponding plain-layout formulas |
| FTS / Inverted | Scoring / block metadata and allocation overhead | Measure retained block descriptors, scoring data, objects and spare capacity |
| ZoneMap | Minimum values | Z * w logical bytes for Z loaded zones and fixed-width values of width w |
| ZoneMap | Maximum values | Z * w logical bytes; min + max total 2 * Z * w |
| ZoneMap | Counts and zone boundaries | Measure retained counts and boundaries; count zones per Fragment, including partial zones |
| ZoneMap | NULL tracking and boxed-value/object overhead | Measure optional NULL row sets and actual scalar allocations; these can dominate logical min/max bytes |
| BloomFilter | Filter bit arrays | Sum actual allocated filter bytes; Z * F for Z equal-sized filters of F bytes each |
| BloomFilter | Zone descriptors / filter objects | Measure retained descriptors, filter parameters and object overhead |
| BloomFilter | NULL tracking | Measure optional NULL row sets |
| RTree | Bounding-box coordinates | 32 * Q for four Float64 coordinates per 2D entry; count Q entries across all cached tree levels |
| RTree | Row / child-page IDs | 8 * Q for UInt64 IDs; coordinates + IDs total 40 * Q |
| RTree | Tree metadata and page offsets | Measure retained tree descriptors and page-offset buffers |
| RTree | NULL tracking, Arrow/objects and spare capacity | Measure actual additional allocations |
| JSON-path wrappers | Underlying index | Use the complete component budget of the selected scalar index family |
| JSON-path wrappers | Path and wrapper structures | Measure retained paths and wrapper allocations in addition to the underlying index |
Bloom filters require the builder's actual block-rounded allocations; row count and false-positive probability alone do not describe every layout. For unfamiliar or experimental layouts, measure loaded allocations instead of applying another family's formula. A small number of matching rows can still require substantial dictionary, tree, or document metadata.
Keep query working memory separate from all cache tables: temporary intersections/unions, result row-ID sets, rechecks of approximate filters, original geometry reads for exact RTree checks, and output-column reads need their own budget.
See the upstream BTree format, Bitmap format, and FTS format for the underlying structures. The formulas above estimate loaded payload, not the compressed file size or total query RSS.
From Index Size to a Per-BE Working Set
For segment s, let N_s be its indexed vector count, P_s its IVF partition count, and H_s the number of distinct partitions to retain on a particular BE. With reasonably balanced partitions:
cached_vectors_on_BE ≈ sum over segments (N_s * H_s / P_s)
index_cache_target ≈ cached vector/row-ID payload
+ cached HNSW graphs
+ centroids, codebooks and other index allocations
+ measured headroom
Use actual partition row counts when partitions are skewed. H_s describes the working set across queries, not just the nprobes of one query. Diverse queries can eventually touch every partition; repeated or concurrent queries can share the same cached partition. Count distinct (dataset, index segment, partition) entries and include all tables, indexes, and versions accessed by that BE. Cache residency on one BE does not warm another BE, and scheduling does not guarantee that every index segment stays on one BE. Do not divide the total index size by the number of BEs unless the observed placement and reuse justify it.
Query-time memory is additional: in-flight partition loads, objects still held by queries after cache eviction, Prefilter state, distance tables, candidate heaps, decoded batches, refinement, and second-phase fetches. top_k, ef, refine_factor, scan concurrency, and nprobes affect this memory or the accessed working set; they do not change the stored bytes per vector. A configured cache capacity is not a hard cap on Lance or BE RSS. Increasing the cache also does not remove distance computation or per-query filtering work.
Configuration Example and Tuning
Assume an FP32, 768-dimensional IVF_PQ index with 100 million vectors, m=96, b=8, and 4096 balanced partitions in one index segment. A representative workload on one BE repeatedly accesses 1024 distinct partitions:
Full code/row-ID payload = 100000000 * (96 + 8) = 10400000000 bytes ≈ 9.69 GiB
Cached payload = 10400000000 * 1024 / 4096 = 2600000000 bytes ≈ 2.42 GiB
IVF centroids = 4096 * 768 * 4 = 12582912 bytes = 12 MiB
PQ codebook = 256 * 768 * 4 = 786432 bytes = 0.75 MiB
If this is the only substantial index working set on the BE and its memory budget permits, 4 GiB is a reasonable initial index-cache allocation for this example, leaving room above the calculated payload for additional index allocations. It is an initial value to validate, not a universal overhead ratio. If the workload instead needs the entire index resident, 4 GiB is insufficient, and even the default 10 GiB is close to the 9.69 GiB payload before other allocations.
For a BE where 4 GiB of index cache and 1 GiB of metadata cache fit alongside peak query memory, other Doris workloads, and process/OS headroom, set:
lance_index_cache_size_bytes = 4294967296
lance_metadata_cache_size_bytes = 1073741824
Restart that BE for these settings to take effect. Apply the budget separately to each BE; these are BE configuration parameters, not catalog properties or SQL session variables.
- Warm up with representative queries and concurrency, including different query vectors, filters, and tables. Repeating one vector alone understates a diverse workload's working set.
- Compare
doris_be_lance_session_index_cache_usage_byteswithdoris_be_lance_session_index_cache_capacity_bytes, the interval increases ofhits_total/misses_total, and the per-queryLanceIndexPartitionCacheMissLoads. Persistent misses with usage near capacity suggest cache churn; misses can also come from new partitions, index versions, or different BEs. Increase capacity only when reuse and available memory justify it. - Start metadata-cache sizing from the default 1 GiB when the BE memory budget permits. Adjust using its
usage_bytesand hit/miss metrics. File/Fragment counts, schema/page metadata, deletion information, and row-ID mappings determine demand; it is not a fixed percentage of vector-index size. The BE metadata cache is separate from the FE table-access cache. - Check BE process memory and query peaks under the intended concurrency before increasing either cache. A high cache hit rate with high latency can indicate scoring, Prefilter, decoding, or fetch costs; additional cache is not necessarily useful.
- Size
lance_data_cache_disk_capacity_bytesindependently for repeated data-file reads such as refinement and output-column fetches. Increasing it does not enlarge the index cache or compensate for index-cache misses: index files bypass this disk data cache. Keeplance_data_cache_read_block_size_bytesat its 1 MiB default initially, then change it only after measuring useful reads versus read amplification.
4. Budget for Shared Caches
On each BE, Lance uses a shared session and one index-cache capacity. Vector partitions, BTree pages, Bitmap/LabelList entries, and full-text index content can compete for this space, including indexes from different tables and catalogs accessed on the same BE. Cache keys separate their identities; they do not reserve capacity. The Lance BE cache settings provide no per-table or per-index-type quota or pinning control.
Only loaded content consumes cache space; merely creating several indexes does not load them all. Once queries use them, loading one index can evict another index's entries. For example, a broad BTree range scan or diverse FTS queries can replace vector partitions that were warm. A later vector query then reloads them. Changes of index version can also introduce new entries while older entries remain resident. Eviction does not immediately free an object still referenced by a running query.
Plan the mixed working set, without counting shared allocations twice:
index_cache_target_on_BE ≈ vector working set
+ BTree working set
+ Bitmap / LabelList / NGram working sets
+ FTS and other index working sets
+ measured headroom
For example, suppose measurements and payload estimates give 2.5 GiB for vectors, 0.25 GiB for BTree, 0.5 GiB for Bitmap, and 0.75 GiB for FTS, including their index structures. Their combined working set is 4 GiB. A 5 GiB initial index cache leaves 1 GiB for growth and estimation error; it does not guarantee every workload will fit. If the BE also has room for 1 GiB of metadata cache and peak query memory, other Doris workloads, and process/OS headroom:
lance_index_cache_size_bytes = 5368709120
lance_metadata_cache_size_bytes = 1073741824
The index and metadata caches have separate capacities: metadata entries do not directly evict index entries, and increasing one does not enlarge the other. Both still consume the same BE process memory. The disk data cache is a third budget and does not cache index files. Separate catalogs on the same BE do not provide index-cache isolation; workloads requiring isolation need separate BE resources and routing supported by the deployment.
To diagnose contention, first warm one query family and record latency and cache-metric deltas. Run a second family, then repeat the first. Renewed misses or reloads with cache usage near capacity suggest eviction pressure, provided index versions and BE placement stayed comparable. Session hit/miss metrics aggregate workloads, so combine them with query profiles. Prewarming every index can itself evict the working set you wanted to retain. Prefer representative mixed-workload warmup, and increase capacity only when reuse and available memory justify it.
5. Choose Index Segment Granularity
An Index Segment, a data Fragment, an IVF partition, and a BTree leaf page are different units. In the current Doris vector-search path, each selected physical Index Segment becomes one scan Split; uncovered Fragments add Flat Search splits. FTS also uses physical index segments. One vector segment's IVF partitions are not separate Doris splits distributed across BEs. Lance can still execute work internally in parallel. For ordinary scalar scans, segment-based grouping depends on the selected index and plan; do not assume every scalar index follows the vector-search split policy.
Let B be the number of BEs actually eligible for this query and S the number of selected vector segments. If S is smaller than B, there are too few indexed scan splits to use all those BEs for that query. Increasing S creates scheduling opportunities, but it does not guarantee that all splits run simultaneously. Scanner limits, Lance's internal concurrency, CPU, I/O, cache residency, and concurrent queries also matter. Increasing Doris pipeline instances alone does not subdivide one index segment.
For a large dataset, compare B, 2 * B, and 4 * B reasonably balanced vector segments as an initial experiment, not an official optimum or a requirement for small datasets. For an illustrative 100-million-vector dataset and eight eligible BEs:
| Segments | Average vectors per segment | What to evaluate |
|---|---|---|
| 8 | 12.5 million | Enough splits to cover the BEs, with fewer per-segment operations |
| 16 | 6.25 million | More scheduling flexibility and less sensitivity to one large segment |
| 32 | 3.125 million | More tasks, but potentially greater search, loading, and merge overhead |
Balance estimated search work and loaded index bytes, not just Fragment counts. Equal rows are a useful starting point for vectors with the same dimension and encoding; HNSW graph size, IVF skew, scalar selectivity, and document length can change the balance. There is no universal rows-per-segment or file-size threshold. High-throughput workloads may favor fewer tasks per query to preserve capacity for concurrent queries; single-query latency may benefit from more splits.
More segments also mean more per-segment setup, metadata/quantizer objects, and local candidate generation. Each vector/FTS split retains up to top_k + offset candidates; Doris does not divide this bound by S. Thus the scan-side candidate bound is S * (top_k + offset), subject to available matches. With top_k=1000, no offset, and 32 splits, that is up to 32000 candidates before Doris reductions. Local TopN can reduce network traffic, so this is not a network-row guarantee. Candidate processing becomes especially relevant for large top_k.
Tune segment count together with each segment's IVF num_partitions and query nprobes. For balanced partitions, a rough probed-vector count is sum(N_s * nprobes_s / P_s), capped at N_s for each segment. This is not an HNSW comparison-count formula. Rebuilding segments or changing partition counts can change recall even at the same nprobes; compare latency at comparable recall. HNSW ef and refinement settings also affect the comparison.
Build and maintain segments with Lance tooling, then commit them as the intended logical index. Several separately named full-table indexes are not a substitute for one index with multiple physical segments. Prefer balanced groups of whole Fragments, avoid accumulating many tiny segments from small appends, and use the embedded Lance version's supported maintenance APIs. Verify the resulting physical segment count and uncovered-fragment count in Doris EXPLAIN/Profile; do not infer it from the number of rows in SHOW INDEX. Merging every segment into one can reduce Doris scan parallelism. See Lance distributed indexing for segment construction and commit concepts.
For BTree point lookups or selective Bitmap queries, splitting into many segments can add lookups and opens while little work is saved per segment. FTS adds per-segment dictionary/posting searches and candidate merging. Evaluate these workloads separately instead of applying a vector segment count to all index families. Splitting also does not remove unnecessary per-row Prefilter construction; verify the reader's filtering behavior independently.
6. Validate the Mixed Workload
- Confirm index selection, physical splits, and index coverage. Include fallback scans and second-phase fetches in the latency analysis.
- Estimate each index's retained payload and add the shared working sets. Reserve query memory separately, including parallel loads, row-ID sets, candidate heaps, and decoding.
- Compare cold-cache and warmed mixed workloads with realistic concurrency. Track P95 latency, throughput, recall where applicable, per-BE peak memory, cache usage and hit/miss deltas, index-load I/O, and scan-time imbalance.
- Change one dimension at a time: cache capacity, segment count, or search parameters. If segment rebuilding changes the model, recheck recall. Stop increasing splits or cache when throughput, latency, or memory ceases to improve.
- Repeat after substantial appends, index rebuilds, BE scaling, or changes in query mix. A configuration calibrated for one vector or table is not a stable budget for all workloads.
7. Diagnose Before Increasing Resources
| Observation | Check first | Action to evaluate |
|---|---|---|
| Repeated index misses with cache usage near capacity | Index/version changes, BE placement, and the combined working set | Increase lance_index_cache_size_bytes if memory permits, reduce unnecessary warmup, or isolate workloads on separate BE resources |
| Vector queries slow down after scalar or FTS queries | Repeat the cross-workload eviction check in shared-cache budgeting | Budget for both query families together; separate catalogs on the same BE do not isolate the cache |
| Few BEs scan a large indexed dataset | Physical segment Split count, eligible BEs, and scanner scheduling | Compare more balanced segments; increasing pipeline instances alone cannot split a segment |
| More segments increase CPU, memory, or merge time | Per-split candidate count, top_k + offset, and query concurrency | Reduce segment count or avoid requesting more candidates than the application needs; recheck recall when search settings change |
| Cache hits are high but latency remains high | Scoring, Prefilter work, decoding, and second-phase fetch time | Optimize the dominant stage; a larger index cache does not remove computation or fetches |
| Metadata reloads persist | Metadata-cache usage and hit/miss deltas; Fragment/file counts and version churn | Adjust lance_metadata_cache_size_bytes independently if reuse and memory permit |
| Process memory is high despite cache usage below capacity | Concurrent query allocations, active loads, retained objects, and other BE workloads | Lower concurrency or cache budgets based on measured peaks; cache capacity is not a process-memory cap |
| Data-file I/O stays high | Disk-cache counters, reuse, read block size, and query output columns | Adjust the data-cache budget or read granularity; index files do not use this disk data cache |
Use the profile counters and BE metrics in Checking Cache Effectiveness. An aggregate hit rate alone does not identify which index was evicted or which stage dominates latency. Change one setting at a time and keep a record of the workload, index version, BE placement, latency, throughput, recall, and peak memory for each comparison.