yy782 commented on issue #66492:
URL: https://github.com/apache/doris/issues/66492#issuecomment-6063914307
# Revised Design — Lance File TVF Path Expansion and Multi-Dataset Support
(Phase 1)
> Responding to the review comments from Gabriel39 in #66492.
Thanks for the thorough review. The comments go well beyond what I had
considered — failure semantics, planning budgets, cancellation and cleanup,
cache identities, and the multi-path representation in particular — and they
materially tightened this design. Several items below exist only because you
raised them. The revision follows your comments point by point; where I have
taken a different route, the reasoning is stated inline so you can push back on
it.
---
## 1. Phase 1 Scope
In scope:
- Read **multiple** Lance datasets in a single TVF invocation: glob
expansion, an explicit path list, or a mix of both. Limited to datasets whose
schemas are **exactly identical**.
- Require **strict schema equality** validation across all datasets.
- One self-consistent metadata snapshot per dataset; the version is pinned
per dataset.
- Preserve the existing single-dataset behavior.
- Phase 1 lands only on `s3://`, and claims support only for combinations
that have been verified end to end. At present that is MinIO, as demonstrated
by `test_lance_s3_tvf.groovy`. OSS / COS / TOS / GCS / Azure, and any other
S3-compatible endpoint that has not been verified, are **not** claimed as
supported.
- One provider per invocation. Reject multi-provider / cross-credential
combinations: multi-dataset mode accepts a single provider and a single
credential set. Preserve the provider information needed to build storage
options through both expansion and open.
---
## 2. Path Input: `uri` and `uris`
- **`uri` keeps its current single-path semantics**, so the behavior of
existing single-dataset input is unchanged.
- Multiple paths are expressed through **`uris`**, not through comma
separation.
- `uri` and `uris` are **mutually exclusive**. Supplying both is a parse
error, reported with an explicit message naming both keys.
- Each element of `uris` is parsed **independently** under the glob grammar.
Literal elements and pattern elements may coexist in the same array.
- Multi-path and glob take effect **only** on `uris`; `uri` retains
single-dataset semantics.
- The **element count** of `uris` is itself bounded. An array of ten
thousand elements must be rejected at parse time, not discovered later during
expansion.
```sql
-- Explicit list
SELECT * FROM s3(
"uris" = '["s3://bucket/a", "s3://bucket/b"]',
"format" = "lance");
-- Literal and glob mixed
SELECT * FROM s3(
"uris" = '["s3://bucket/b", "s3://bucket/20*"]',
"format" = "lance");
```
Report an explicit parse error for each of the following: malformed JSON, a
non-string element, an empty element, an empty array, an invalid element
scheme, an element ending in `/`, and an invalid pattern.
---
## 3. Glob Grammar and Expansion
**Supported**: `*`, `?`, `[]`, `{}`.
**Not supported**: `**` (cross-level recursion). `**` is rejected with an
explicit error rather than being silently treated as `*`.
Per-element semantics:
- A pattern is split into a root (`scheme://authority`) and a sequence of
path segments.
- A **literal segment** performs no storage access; it is appended to every
prefix already reached.
- A **wildcard segment** produces one bounded directory listing (§4) per
parent prefix already reached.
- Expansion proceeds breadth first, level by level.
**Escaping and brace scope**:
- A backslash escapes the next character, so `\*` and `\?` are literals.
This is the spelling to use when a path genuinely contains a wildcard character.
- `{}` applies **only within a single path segment**, enumerating
alternatives for that segment: in `s3://bucket/{a,b}/ds`, the second segment is
expanded. A `{...}` that spans `/` is an invalid pattern, not a cross-segment
union.
- Nesting inside `{}` and `[]` is not supported and is rejected as an
invalid pattern in Phase 1.
- An unclosed `[` or `{`, and an empty `[]`, are invalid patterns, rejected
at parse time before any listing is issued.
- `[]` supports character sets, ranges, negated ranges, and negated sets. It
does not support sequences, steps, or nesting. `{}` supports only enumerating
alternatives; no nesting, sequences, or steps.
**Trailing slash**: elements of `uris` may not end with `/`. This is
rejected at parse time with a message asking for the trailing slash to be
removed, so that a given directory has exactly one spelling. The existing
behavior of `uri` is unaffected.
---
## 4. Bounded Enumeration
**(a) A bounded, pageable directory-listing primitive.** With
`delimiter="/"`, `maxKeys`, and a continuation token, it returns one page per
call so the caller can stop midway. Budgets below are checked as each page
arrives, so enumeration can be aborted in the middle.
**(b) Incremental budgets, counted across the entire TVF invocation.**
| Budget | What it counts | Default |
| --- | --- | --- |
| Directory candidates | Total directory entries listed, counted **before**
deciding whether each is a dataset | 1 000 |
| Listing requests | Individual listing calls issued, summed across all
storage backends, cumulative for the invocation | 2 000 |
| Listing time | Wall-clock time spent listing, cumulative for the
invocation | Must be below the deadline; value not fixed yet |
| Datasets | Datasets actually read, after deduplication | 500 |
| Fragments / splits | Sum of fragments / splits over all datasets | 10 000 |
| Fields per schema | Fields in one dataset's schema | 1 000 |
| Total planned fields | Sum of fields over all datasets in the plan | No
separate limit; equals the per-schema cap × the dataset cap |
| Concurrent opens | Dataset handles open at the same time | 1 |
> These defaults are reference values, not tuned values. They need to be
confirmed in the PR.
**Memory bound.** The metadata extracted from each dataset (`version` /
`schema` / `fragments`) must be retained in the plan. This part grows linearly
with the dataset count and is the real planning-time memory bound. It is
therefore constrained jointly by the dataset cap and the per-schema field cap.
---
## 5. Dataset Identity, Deduplication, and Determinism
**Dataset-set semantics.** An explicit URI and a glob may identify the same
dataset, and a dataset must be read only once; deduplicate with a set.
Deduplication happens **after** expansion and normalization and **before**
opening and version pinning, so no dataset is opened twice and two different
pinned versions of the same dataset cannot both enter the plan.
**Identity is the only deduplication key.** Deduplicate by normalized
dataset identity only, independent of data content: two datasets with identical
contents but different URIs remain two datasets, each read once, with the rows
doubled. Results are not deduplicated by row.
Normalization used for identity, in this order:
1. The bucket is preserved verbatim; no case folding.
2. The path is split on `/` and compared verbatim; no folding of any kind. A
repeated slash is a valid character in an object key — `a//b` and `a/b` are two
distinct keys — so `//` is **rejected explicitly** and never folded into a
single `/`. Folding would change the object being addressed.
3. `..` and `.` in path segments are treated as literals and are not folded:
`a/../b` and `b` are two distinct keys in object storage.
**Deterministic order.** Candidates returned by listing are sorted before
use. The sort result determines both the traversal order and the selection of
the reference schema (§8).
---
## 6. Version Pinning per Dataset
One open takes `version`, `schema`, and `fragments` from the same handle. No
common version is synthesized; each dataset keeps its own `(uri, version,
fragments)`.
> The dataset set is resolved during TVF construction and is held by an
object that is not rebuilt within a plan.
**Therefore**, within one analyzed plan — including any execution retry —
the dataset set, versions, and fragments all come from that single resolved
result; `latest` is not resolved a second time.
**Therefore**, if a pinned version is deleted or becomes unreadable, that is
an open failure: the query fails (§7) with an explicit error, and there is no
fallback to `latest`.
> The resulting consistency boundary needs to be stated in the user
documentation: per-dataset snapshot consistency does not provide a consistent
snapshot across datasets.
---
## 7. Any Open Failure Fails the Query
Regardless of explicit multi-path, mixed input, or glob matching: if **any**
candidate — including a path finally determined not to be a dataset — cannot be
opened, the whole query fails.
Each failure reports:
- the failing **path**;
- the failing **operation** (listing, dataset open, …);
- **no sensitive information**: credentials, signed-URL secrets, and
storage-option values must never appear in the message.
Two zero-result cases are reported separately:
- **No path matched**: enumeration produced zero candidates. The message
states that the pattern matched nothing and reports the pattern itself together
with the number of listing requests issued.
- **Matched but could not be opened**: candidates were produced and all
failed. The message states how many candidates were found and reports the first
failure in full.
---
## 8. Strict Schema Equality
Validation is performed entirely in FE, before any scan task is dispatched.
The comparison is **recursive** and requires exact equality of:
- field **names**;
- types and **type parameters**: decimal precision and scale, timestamp unit
and timezone, fixed-size list dimension, and comparable attributes;
- **nullability**;
- **nested structure**, compared recursively;
- field **order**, which is part of the equality definition. Any pushdown
later placed on the multi-dataset path must be designed against this premise
and must not assume that fields of different datasets are interchangeable.
- Lance internal **field IDs** need not be equal.
**Reference schema**: the first dataset under the deterministic order (§5);
all others are compared against it. The TVF output columns continue to be
derived from a single schema, exactly as today.
---
## 9. Dataset State and Cache Audit
This section is an audit. It does not assert that the existing caches are
defective.
**BE side — changes required in Phase 1**
> Today `LanceTableReader` keeps only one open dataset and errors when the
dataset key changes, while a single scanner processes several scan ranges and
reuses the same reader.
- **Chosen approach**: when the dataset key (`uri` + `version` + storage
options) changes because of a new split, close and reset the dataset-scoped
state instead of erroring. This covers the opened dataset handle and the
recorded dataset key, the Arrow schema bound in the record-batch converter, the
cached dataset schema, and the full-text search query context. Every split
carries the identity of its own dataset, so which dataset to open is explicitly
determined and requires no implicit context.
- **Profile accounting**: the two data-cache byte counters take the absolute
snapshot of the current dataset handle and **replace** rather than accumulate.
Accumulate into a local running total before switching datasets and write the
accumulated value back; otherwise only the numbers of the last closed dataset
survive.
- **Pushed expression binding**: that cache implicitly assumes the schema is
invariant within one scan local state and that the query conditions are unique.
No change is needed in this phase; ambiguity arises only once multiple schemas
are merged.
**Shared session cache.** The BE holds one process-global Lance session
owning the metadata and index caches. Both tiers are a "single global pool plus
key prefix" structure: the metadata cache is prefixed by dataset URI and
carries the version in the entry key, and the index cache adds the index
identity on top of it. Cross-dataset and cross-version false hits are therefore
already ruled out, and Phase 1 does not modify this cache.
---
## 10. Cancellation and Resource Cleanup
> A planning-time timeout already exists, but its checkpoints do not cover
the scan-node finalize segment, which is exactly where listing and open happen.
The planning path therefore has to compute its own deadline and place its
own checkpoints.
- **Deadline.** `min(query_timeout, 60s)`, consistent with the existing
Lance metadata read. It is not disabled together with the planning-timeout
switch.
- **Per-operation timeout.** A deadline only helps between checkpoints; a
real stall happens inside one listing or one open, where no checkpoint ever
runs. In addition to the overall deadline, each listing and each open therefore
needs a timeout of its own.
- **Checkpoints.** Before each listing, between pages of one listing, before
each open, and inside the submitted tasks — not only at the points where
results are awaited.
- **Concurrency.** One shared bounded thread pool: fixed thread count,
bounded queue, and rejection with an error rather than an unbounded backlog or
blocking the submitting thread. No thread pool per dataset and no unbounded
task queue.
- **Failure handling.** On the first failure: stop issuing new listing
requests and new opens, let in-flight work finish but start nothing new, then
close every handle and allocator already created. Tasks already inside a native
call are not interrupted; wait for them to return.
- **Paths to cover**: success, listing failure, open failure, timeout.
---
## 11. Path Parsing and Identity Consistency
- Listing, deduplication, and the Lance open all operate on the **same
normalized URI string**.
- Do not fold `..` and `.` for object keys: doing so changes the object that
is addressed (§5).
---
## 12. Preserving Existing Single-Dataset Behavior
- **`uri` parsing and routing are unchanged**, so `uri` semantics stay
exactly as they are today.
- **`local()` single dataset is unchanged**, and multi-dataset queries are
explicitly rejected with a clear error.
- **`file()` delegation is unchanged**: it forwards to the same TVF, so it
gains new behavior only where the underlying TVF supports it.
- The following must keep working and are covered by the test plan (§13):
empty datasets and datasets with zero fragments, `COUNT(*)`, `LIMIT`,
cancellation, `file()` delegation, and the existing single-dataset `s3()` and
`local()` paths.
---
## 13. Test Plan
| Acceptance item | Cases covered |
| --- | --- |
| Multi-dataset results are equivalent to `UNION ALL` | N datasets read in
one call, compared row by row against the `UNION ALL` of N single-dataset scans
|
| Overlapping datasets are read once | An explicit path and a glob
identifying the same dataset |
| Query behavior when a dataset fails to open | Permission error, timeout,
corrupt manifest, non-dataset directory; covered separately for explicit paths
and glob matches |
| Schema mismatch is rejected | Swapped field order, different name case,
different decimal precision and scale, different timestamp unit / timezone, and
comparable differences |
| Path parsing and glob grammar | Keys containing commas, `{a,b}` spanning
segments, `\{` and `\*`, unclosed `[` and `{`, empty `[]`, `**`, `//`, trailing
`/`, malformed JSON, `uri` and `uris` used together |
| Enumeration and planning budgets | A large listing with few matches; few
datasets with very many fragments |
| Cleanup on cancellation and partial failure | Other datasets still queued
when one fails; timeout; cancellation in the middle of listing |
| No regression in existing single-dataset behavior | Single-dataset `s3()`
and `local()`, `file()` delegation, empty dataset, zero fragments, `COUNT(*)`,
`LIMIT`, cancellation, multi-dataset `local()` |
| Traversal order and the two zero-result cases | Order consistency across
runs; the "no path matched" and "matched but could not be opened" errors |
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]