yy782 commented on issue #66492:
URL: https://github.com/apache/doris/issues/66492#issuecomment-6011275985
# Design Proposal — Lance file TVF path expansion and multi-dataset support
---
## 0. Goals and Phase-1 boundaries
### 0.1 Goals
- Support glob wildcard expansion to **multiple** Lance datasets.
- Read **multiple** datasets in a single TVF call (explicit multi-path /
glob / mixed).
- Support this on the `s3()` entry point, and **automatically cover** OSS /
COS / TOS / GCS / Azure (whose URIs are normalized to `s3://`).
### 0.2 Explicitly out of scope (boundary shrink)
- **Multi-dataset `local()` is out of scope for this phase; single-dataset
keeps the existing capability from baseline #65730 and stays untouched this
phase.** `local()` pins the BE that holds the file via `backend_id`: the FE
never touches the local disk, it only produces one whole-dataset split
(`version=0` → latest) routed back to that BE, which opens it with local FS
permissions. The FE therefore cannot enumerate multiple local datasets —
multi-dataset would require an architecture change (BE-side enumeration or
explicit multi-path, with more complex `backend_id` constraints), so it is out
of scope for this phase (see §7 Q1).
- **`hdfs()` blocked upstream**: lance-c does not support the hdfs backend
yet; revisit once upstream is ready.
---
## 1. TVF surface and glob syntax
### 1.1 Supported entry points
| TVF / backend | Status in this phase |
| --- | --- |
| `s3()` | Main landing point; multi-dataset Lance reading is supported here
first |
| OSS / COS / TOS / GCS / Azure | URI is normalized to `s3://`, goes through
the same function, **covered automatically** |
| `file()` | Pure delegation to the TVFs above, **covered automatically** |
| `local()` | **excluded (not in Phase-1)**, see §0.2 / §7 Q1 |
| `hdfs()` | Blocked upstream (lance-c does not support it yet) |
| `http()` / `http_stream()` / `cdc_stream()` / `group_commit()` | Do not
support the Lance format, not involved |
### 1.2 glob syntax scope
| Syntax | Verdict | Basis |
| --- | --- | --- |
| `*` | **Supported** | Both DuckDB and ClickHouse support it |
| `?` | **Supported** | Present in DuckDB's official glob table; also in
ClickHouse's official list |
| `[]` (character class) | **Supported** | DuckDB's official glob table has
`[abc]`, `[a-z]` |
| `{}` (enumeration) | **Supported** | ClickHouse documents `{abc,def}` for
url parameters |
| `**` (recursive across levels) | **Not supported** | With level-by-level
expansion (§1.3), every level costs one list call and the request count grows
with the total number of directories traversed; dataset directories also have
no uniform naming suffix, so they can only be determined one by one (§1.4) and
cannot be pruned up front |
> The syntax comparison above all comes from the **plain-file (Parquet /
CSV) reading** scenario of the two systems: neither ClickHouse nor DuckDB
supports glob expansion for Lance.
### 1.3 Directory expansion
Expand level by level with `fileSystem.listDirectories`, then filter once
and keep the paths that qualify.
### 1.4 Determining whether an expanded result is a dataset
The candidate directories produced by expansion are not all Lance datasets
(the same level may be mixed with ordinary directories such as `_tmp/`,
`xxx_bak/`), so each one has to be determined individually.
My current approach: **just try `Dataset.open()` — if it opens, it is a
valid dataset; if it fails, skip it.**
```java
Dataset.open()
```
This matches DuckDB's thinking: after expanding a glob, DuckDB also performs
no format pre-check and simply opens the file — "whether it can be opened" is
itself the determination.
There is one difference from DuckDB, however, and it affects the failure
semantics in §3: DuckDB **fails the whole query** when an open fails, whereas
here a candidate directory that cannot be opened is **skipped**. The reason is
that expanded directories may legitimately be non-dataset neighbours, so
"cannot open" is an expected, normal case and should not take down the entire
query.
### 1.5 Expansion boundaries and failure handling
- **Zero matches**: fail directly, but distinguish the two cases — "the
directory listing produced nothing" and "things were listed but none is a
dataset" — so the user can diagnose it themselves.
- **Count cap**: cap the number of results from one expansion at **1000**.
The concern is a pattern like `s3://bucket/*` expanding into tens of thousands
of datasets and dragging down the FE; exceeding the cap fails directly and asks
the user to write a narrower glob.
---
## 2. Dataset boundaries and the multi-dataset model
### 2.1 Ways to specify multiple datasets
| Means | Example | Scope of this phase |
| --- | --- | --- |
| glob wildcard | `s3://bucket/datasets/*` | **Supported** |
| explicit multiple paths | `"uri" = "s3://bucket/a, s3://bucket/b"` |
**Supported** (comma separation has an ambiguity risk, see §2.1.1) |
| mixed form | `"uri" = "s3://bucket/a, s3://bucket/20*"` | **Supported**
(see §2.1.1) |
### 2.1.1 Multi-path separation: ambiguity caused by S3 keys that may
contain commas
S3 keys may contain commas, so `"uri" = "s3://bucket/a, s3://bucket/b"` is
inherently ambiguous — it may be two datasets, or it may be a single key that
contains a comma.
Both other systems deliberately avoid this form: DuckDB uses an **array
literal** (`read_parquet(['a.parquet','b.parquet'])`, internally a
`LIST(VARCHAR)`), ClickHouse uses `{}` enumeration (§1.2). Neither is ambiguous.
But Doris TVF parameters are `Map<String, String>` string pairs with **no
array literal available**, so my plan is:
- **Keep comma separation** for now, but state explicitly in the
documentation that **paths containing commas are not supported**, consistent
with the existing glob syntax.
### 2.2 version pinning
A fixed version is required by snapshot semantics: all splits of one query
must read the same version. The problem is that **the latest versions of
different datasets are simply not the same number** — there is no way to derive
a single common version.
So I lean towards keeping the `(uri, version, fragments)` triple
**separately per dataset**: each dataset reads along its own version axis
rather than forcing a merge.
### 2.3 Schema compatibility: merge rules
After expanding N datasets, one question cannot be avoided: **how do N
schemas become 1**.
The basic convention I have settled on is **union of columns, missing
columns filled with NULL** — this is exactly the semantics of DuckDB's
`union_by_name=true`, quoted from the official docs:
> The `union_by_name` option can be used to unify the schema of files that
have different or missing columns. **For files that do not have certain
columns, NULL values are filled in**
There is a prerequisite here: for a missing column to be filled with NULL,
the merged column must be **nullable**. If some dataset declares that column
NOT NULL, the merge has to relax it to nullable — otherwise "this dataset does
not have the column at all" simply cannot be represented.
### 2.4 Conflicts that need a ruling
> This section (including §2.4.1–§2.4.3) refers to **DuckDB's and
ClickHouse's multi-file merge for "plain files"** — the rules they apply when
reading plain files such as Parquet / JSON / CSV and combining several schemas
into one.
| Conflict | Concrete case | Handling |
| --- | --- | --- |
| Column name case | `AGE` on one side, `age` on the other | **Merge into
one column** (Doris column names are case-insensitive, both refer to the same
column); the types of the two are then decided by the whitelist rules in
§2.4.1, and anything outside the whitelist errors out |
| Different column types | `int32` vs `int64`, `int64` vs `uint64`, `int` vs
`string` | **Unify via named whitelist rules, error out outside the whitelist**
(ClickHouse style, see §2.4.1). For example `int32` vs `int64` can be unified;
`int64` vs `uint64` and `int` vs `string` are not supported as column type
conversions under Lance semantics, so they **error out** |
| Nullability | nullable on one side, non-nullable on the other | **Merge to
nullable** (ClickHouse style, see §2.4.2). Nested types are pushed down
recursively, but the `map` key is not relaxed |
| Different nested fields | `struct<a:int>` vs `struct<a:int,b:string>` |
**Union by field name**: merge into `struct<a:int,b:string>`, with the missing
sub-field filled with NULL; same-named sub-fields recursively look for a common
type and only error out if none is found (see §2.4.3) |
#### 2.4.1 Type conflict resolution: I prefer the ClickHouse style over the
DuckDB style
DuckDB's `union_by_name` is "pick the higher-scored / wider type, never
throw" — `int` and `string` are silently merged into `string`, with integers
converted to strings and **no hint whatsoever**. That suits us poorly: the
merged type is eventually pushed down to the BE to materialize Arrow data (§4),
so "picking a wider type" does not really solve anything; it only moves the
failure from FE to BE and makes it harder to diagnose.
ClickHouse is **named and controlled**: it only unifies under a few
whitelisted rules (promotion within the numeric family, integer↔float, and so
on), and for nested types it **recursively** looks for a common type and gives
up when none is found — it is not a general implicit cast. When unification is
impossible it **errors out directly**.
So my plan for landing this:
- Make the merge rules a **whitelist of named rules**, not a general "take
the wider one": only allow a few cases such as promotion within the numeric
family (`int8/16/32 → int64`, `float → double`); cross-family cases (`int` vs
`string`, signed vs unsigned) **always error out**.
- Include **both sides of the conflict** in the error message: which
dataset, and the column name and type on each side, so the user can locate the
offending dataset.
- With this, the BE **does not need a general type adaptation layer** —
incompatibilities are rejected at the FE, and the BE only needs to support the
narrowing/widening materialization corresponding to the few whitelisted
conversions (see §4).
One more note on column name case: ClickHouse treats `age` / `AGE` as two
independent columns because its identifiers are case-sensitive. That route
relies on the precondition of "case-sensitive identifiers", which **Doris does
not have** (Doris column names are case-insensitive); copying it would produce
ambiguous references, so I do not adopt it and instead **merge into one
column** per Doris semantics. But the two cases must be distinguished: if
**within one single dataset** there really are both `age` and `AGE` columns,
then the data genuinely has two fields and merging would lose data — that
should **error out explicitly** with an explanation.
#### 2.4.2 Nullability merge: mixed means nullable (ClickHouse style)
When one side is nullable and the other is not, the merged result is relaxed
to nullable.
The basis is ClickHouse's merge strategy: as soon as any `Nullable` appears
among the types, all other types that can be wrapped in `Nullable` are relaxed
as well. This is consistent with the basic convention in §2.3 ("union of
columns, missing filled with NULL") — if dataset A declares a column NOT NULL
while dataset B does not have that column at all, the merged column must be
able to represent "the value does not exist".
Nested types are relaxed by the same **recursive push-down** rather than
only at the top level; but one constraint has to be honoured along with it — a
`Map` key can never be nullable (only the value is relaxed).
#### 2.4.3 Nested type merge: union by field name
**Example A: struct missing a field → union, missing sub-field filled with
NULL**
Suppose `s3://bucket/db*` expands into two datasets, db1 and db2.
db1's schema:
```
id: int64
addr: struct<city: string, zip: string>
```
db2's schema:
```
id: int64
addr: struct<city: string, country: string>
```
In `addr`, db1 has `zip` and db2 has `country`. After merging (union):
```
id: int64
addr: struct<city: string, zip: string, country: string>
```
The rows actually read:
| Source | id | addr.city | addr.zip | addr.country |
| --- | --- | --- | --- | --- |
| row from db1 | 1 | Beijing | 100000 | **NULL** |
| row from db2 | 2 | Shanghai | **NULL** | CN |
That is "missing sub-field filled with NULL" — the same thing as filling
NULL for a missing top-level column.
**Example B: two levels of nesting, recursing level by level**
db1: `addr: struct<geo: struct<lat: float>>`
db2: `addr: struct<geo: struct<lat: double, lng: double>>`
- Level 1 `addr`: both sides have `geo` → keep recursing.
- Level 2 `geo`: `lat` exists on both sides → `float` vs `double` → take
`double`; `lng` exists only in db2 → keep it as is.
- Merged result: `addr: struct<geo: struct<lat: double, lng: double>>`
When reading db1, `geo.lng` is filled with NULL.
Two kinds of "inconsistency" are easy to confuse, and they are handled
differently:
| Case | Example | Result |
| --- | --- | --- |
| Field **name** differs (the focus of this section) | `addr<zip>` vs
`addr<country>` | **Union**, no error |
| Same field name but **incompatible type** | `addr<zip:int32>` vs
`addr<zip:string>` | **Error** (recursive common-type search fails, decided by
the whitelist in §2.4.1) |
#### 2.4.4 Multi-dataset constraints
Cross-provider / cross-credential multi-dataset queries are not supported in
this phase.
---
## 3. Execution model (FE planning cost / failure semantics / memory)
### 3.1 FE planning cost
Every dataset needs one `Dataset.open()` for metadata loading (determination
plus fetching version / fragments / schema, see §1.4). The planning cost grows
linearly with N datasets, so allocator / memory usage has to be evaluated on
the FE side and tied to the cap of 1000 in §1.5, so that planning does not drag
the FE down.
> I previously considered using `Dataset.listManifestLocations` to only list
manifests and save a full open, but the `lance-core:12.0.0` pinned by Doris
does not have that method yet (it arrives in a newer version). So for this
phase we stick to `Dataset.open()` and rely on the count cap as the safety net.
If the lance dependency is upgraded later, this optimization can be revisited.
### 3.2 Shape of the scan tasks dispatched by the FE
Add a multi-dataset branch to the existing `TVFScanNode#getLanceSplits`:
https://github.com/apache/doris/blob/78aa9c996accd114a6e0f5943d7882ee1a34011a/fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java#L195
```java
fe/fe-core/src/main/java/org/apache/doris/datasource/tvf/source/TVFScanNode.java
→ TVFScanNode#getLanceSplits() (multi-dataset branch)
```
Roughly:
```java
private List<Split> getLanceSplits() throws UserException {
if (!sessionVariable.enableFileScannerV2) {
throw new UserException("Lance TVF requires
enable_file_scanner_v2=true");
}
if (tableValuedFunction.getTFileType() == TFileType.FILE_LOCAL) {
// A local dataset is visible to its selected BE, not to FE. Keep
exactly one
// whole-dataset split; BE resolves version zero to latest when it
opens the dataset.
return Collections.singletonList(
LanceSplit.scanLatestDataset(tableValuedFunction.getFilePath()));
}
List<LanceTableMetadata> datasets =
tableValuedFunction.getLanceDatasets();// newly added API, returns multiple
datasets
if (datasets.isEmpty()) {
long version = tableValuedFunction.getLanceDatasetVersion();
if (version <= 0) {
throw new UserException(
"S3 Lance TVF metadata was not initialized with a
fixed dataset version");
}
return
LanceSplitBuilder.buildFragmentSplits(tableValuedFunction.getFilePath(),
version,
tableValuedFunction.getLanceFragments(), 1);
}
List<Split> splits = Lists.newArrayList();
for (LanceTableMetadata dataset : datasets) {
splits.addAll(LanceSplitBuilder.buildFragmentSplits(
dataset.getDatasetUri(), dataset.getVersion(),
dataset.getFragments(), 1));
}
return splits;
}
```
---
## 4. BE-side relaxation (execution side of schema compatibility)
https://github.com/apache/doris/blob/78aa9c996accd114a6e0f5943d7882ee1a34011a/be/src/format_v2/lance/lance_record_batch_converter.cpp#L124
```
be/src/format_v2/lance/lance_record_batch_converter.cpp
→ LanceRecordBatchConverter::_bind_schema(...)
```
It currently requires **strict two-way agreement between the columns
returned by lance-c and the FE projected columns**: one column short → `Lance
did not return requested column`; one column extra → `Lance returned unknown
column`.
With a single dataset the two are naturally identical and this never causes
trouble, but after schema merging the FE pushes down a **union**, and when
reading any one dataset it necessarily lacks part of that union. So it has to
be relaxed to: columns absent from the projection are **filled with NULL**
rather than raising an error.
Nested fields need attention here: nested `struct` also takes a union of
fields (§2.4.3), so "fill missing with NULL" must **recurse into nested
fields** and not only handle top-level columns — if a dataset's `struct` is
missing a sub-field, that sub-field must be filled with NULL too.
Correspondingly, the type merge side must also **recursively** find a common
type for same-named sub-fields (both DuckDB and ClickHouse merge nested types
recursively, see §2.4.3). This means the BE-side NULL filling and type
unification both have to support nesting, which is more work than "top-level
columns only" and should be evaluated together.
---
## 5. Security and path handling
- **Prefix constraint**: glob expansion results must still fall inside the
user's authorized path prefix, preventing `../` or an out-of-range prefix from
escaping to unauthorized datasets.
- **local constraint**: `local()` multi-dataset involves `backend_id`
pinning and local permissions; it is not extended in this phase and needs to be
documented as unsupported for `local()` (see §0.2).
---
## 6. Testing
- **Scope 1 (glob)**: expansion of `*`/`?`/`[]`/`{}`; `**` not supported;
non-dataset directories skipped; both error wordings for "zero matches" and
"all failed"; error when the cap of 1000 is exceeded.
- **Scope 2 (multi-dataset)**: explicit multi-path / glob / mixed; schema
union with missing columns filled with NULL; case merging; unification within
the type whitelist and errors outside it (the error carries dataset + column
information); nullability relaxation; nested recursive union / NULL filling;
error when same-named sub-field types are incompatible
- **Scope 3 (entry points)**: `s3()` as the main target;
OSS/COS/TOS/GCS/Azure normalized to `s3://` and covered automatically; `file()`
delegation; `local()` single dataset keeps the existing behavior; `hdfs()`
unsupported.
---
## 7. Open Questions for Reviewers
- **Q1**: Should `local()` be included in this phase, or strictly excluded
from Phase-1? (see §0.2)
- **Q2**: Are the two error wordings for "zero matches" diagnostic enough?
(see §1.5)
---
--
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]