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]

Reply via email to