github-actions[bot] commented on code in PR #68453:
URL: https://github.com/apache/doris/pull/68453#discussion_r4085941205
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,68 +250,382 @@ public LanceTableMetadata loadBasicTableMetadata(String
dbName, String tableName
}
public Schema loadTableSchema(String dbName, String tableName) {
- return readTableSnapshot(dbName, tableName, Optional.empty(),
+ return readTableSnapshot(dbName, tableName, LanceRefSelector.latest(),
(dataset, access, metrics) -> metrics.measure(Stage.SCHEMA,
dataset::getSchema));
}
public LanceTableMetadata loadTableMetadata(String dbName, String
tableName,
Optional<TableSnapshot> tableSnapshot) {
- return loadQueryMetadata(dbName, tableName, tableSnapshot,
LanceMetadataLoader.MetadataScope.WITH_INDEXES);
+ return loadTableMetadata(dbName, tableName,
LanceRefSelector.snapshot(tableSnapshot));
+ }
+
+ public LanceTableMetadata loadTableMetadata(String dbName, String
tableName, LanceRefSelector selector) {
+ return loadQueryMetadata(dbName, tableName, selector,
LanceMetadataLoader.MetadataScope.WITH_INDEXES);
}
private LanceTableMetadata loadQueryMetadata(String dbName, String
tableName,
Optional<TableSnapshot> tableSnapshot,
LanceMetadataLoader.MetadataScope mode) {
- return readTableSnapshot(dbName, tableName, tableSnapshot,
+ return loadQueryMetadata(dbName, tableName,
LanceRefSelector.snapshot(tableSnapshot), mode);
+ }
+
+ private LanceTableMetadata loadQueryMetadata(String dbName, String
tableName,
+ LanceRefSelector selector, LanceMetadataLoader.MetadataScope mode)
{
+ return readTableSnapshot(dbName, tableName, selector,
(dataset, access, metrics) ->
LanceMetadataLoader.read(dataset, access, mode, metrics));
}
- /** Pins one resource generation, resolved table access, and the Dataset
version for the whole read. */
- private <T> T readTableSnapshot(String dbName, String tableName,
Optional<TableSnapshot> tableSnapshot,
+ /**
+ * Pins one resource generation, resolved table access, and the Dataset
version for the whole read.
+ *
+ * <p>The latest version of the main chain is opened once and every other
selector is a
+ * checkout from that handle, so the SDK resolves the ref with the same
commit handler
+ * (the namespace's, for a managed table). A tag is resolved first to the
chain and version it
+ * points at, so a tag created on a branch selects that branch. The two
shortcuts that skip the
+ * latest open are an explicit version on the main chain, and {@code FOR
TIME AS OF} on a
+ * managed table whose namespace reports commit times.
+ */
+ private <T> T readTableSnapshot(String dbName, String tableName,
LanceRefSelector selector,
SnapshotReader<T> reader) {
- LanceTableAccess tableAccess = null;
+ ReadState state = new ReadState(selector, dbName + "." + tableName);
LanceMetadataMetrics metrics =
LanceMetadataMetrics.startMetadataRead();
try {
T result;
try (BufferAllocator allocator =
namespaceAllocator.newChildAllocator(
"lance-metadata-read", 0, namespaceAllocator.getLimit())) {
- tableAccess = metrics.measure(Stage.TABLE_ACCESS,
+ state.access = metrics.measure(Stage.TABLE_ACCESS,
() -> namespaceClient.resolveTableAccess(dbName,
tableName));
- OptionalLong version = OptionalLong.empty();
- if (tableSnapshot.isPresent()) {
- TableSnapshot snapshot = tableSnapshot.get();
- if (snapshot.getType() ==
TableSnapshot.VersionType.VERSION) {
- version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
- } else {
- long timestamp =
TimeUtils.timeStringToLong(snapshot.getValue(), TimeUtils.getTimeZone());
- if (timestamp < 0) {
- throw new IllegalArgumentException(
- "Cannot parse Lance FOR TIME AS OF value
'" + snapshot.getValue() + "'");
- }
- try (Dataset latest = openDataset(allocator,
tableAccess, OptionalLong.empty(), metrics)) {
- version =
OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE,
- () ->
LanceSnapshotResolver.getVersionAtOrBefore(latest, timestamp)));
- }
+ OptionalLong direct = directMainVersion(state, metrics);
+ if (direct.isPresent() || isLatestMain(selector)) {
+ state.version = direct;
+ try (Dataset dataset = openDataset(allocator,
state.access, direct, metrics)) {
+ result = reader.read(dataset, state.access, metrics);
+ }
+ } else {
+ try (Dataset main = openDataset(allocator, state.access,
OptionalLong.empty(), metrics)) {
+ result = readFromLatest(main, state, reader, metrics);
}
- }
- try (Dataset dataset = openDataset(allocator, tableAccess,
version, metrics)) {
- result = reader.read(dataset, tableAccess, metrics);
}
}
metrics.succeeded();
return result;
+ } catch (LanceUserFacingException e) {
+ throw new RuntimeException(e.getMessage(), e);
} catch (Exception e) {
- throw LanceErrorMessages.failure("Failed to load Lance table
metadata for " + dbName + "." + tableName, e,
- tableAccess == null ? null : tableAccess.getDatasetUri(),
- tableAccess == null ? namespaceStorageOptions :
tableAccess.getStorageOptions(), catalogSecrets);
+ LanceTableAccess access = state.access;
+ String uri = access == null ? null : access.getDatasetUri();
+ Map<String, String> options = access == null ?
namespaceStorageOptions : access.getStorageOptions();
+ String what = state.displayName();
+ if (state.branch.isPresent() && !state.branchCheckedOut &&
isBranchNotFound(e, state.branch.get())) {
+ throw new RuntimeException("Lance branch '" +
state.branch.get() + "' of " + state.tableName
+ + state.selector.getTag().map(tag -> " (tag '" + tag +
"')").orElse("")
+ + " was not found" + (isNamespaceMiss(e, "table branch
not found") ? " in the namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ if (state.version.isPresent() && isVersionNotFound(e)) {
+ throw new RuntimeException("Lance version " +
state.version.getAsLong() + " of " + what
+ + state.selector.getTag().map(tag -> " (tag '" + tag +
"')").orElse("")
+ + " was not found" + (isNamespaceMiss(e, "table
version not found") ? " in the namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ String hint = access != null && access.isManagedVersioning() &&
isAccessDenied(e)
+ ? " (reading a namespace-managed Lance table may need
write access to finalize a staged manifest)"
+ : "";
+ throw LanceErrorMessages.failure("Failed to load Lance table
metadata for " + what + hint, e, uri, options,
+ catalogSecrets);
} finally {
metrics.close();
}
}
+ /** What a read has resolved so far; the catch block reports errors
against it. */
+ private static final class ReadState {
+ private final LanceRefSelector selector;
+ private final String tableName;
+ private LanceTableAccess access;
+ private Optional<String> branch;
+ /** Set once the branch's latest version was checked out, i.e. the
branch exists. */
+ private boolean branchCheckedOut;
+ private OptionalLong version = OptionalLong.empty();
+ /** The namespace's version list per chain ("" is main), fetched at
most once per read. */
+ private final Map<String, List<TableVersion>> namespaceVersions = new
HashMap<>();
+
+ private ReadState(LanceRefSelector selector, String tableName) {
+ this.selector = selector;
+ this.tableName = tableName;
+ this.branch = selector.getBranch();
+ }
+
+ private String displayName() {
+ return tableName + branch.map(name -> "@" + name).orElse("");
+ }
+ }
+
+ private static boolean isLatestMain(LanceRefSelector selector) {
+ return !selector.getTag().isPresent() &&
!selector.getBranch().isPresent()
+ && !selector.getSnapshot().isPresent();
+ }
+
+ /**
+ * The main-chain version a selector names without looking at the latest
manifest: an explicit
+ * version, or {@code FOR TIME AS OF} on a managed table whose namespace
reports commit times.
+ */
+ private OptionalLong directMainVersion(ReadState state,
LanceMetadataMetrics metrics) {
+ LanceRefSelector selector = state.selector;
+ if (selector.getTag().isPresent() || selector.getBranch().isPresent()
|| !selector.getSnapshot().isPresent()) {
+ return OptionalLong.empty();
+ }
+ TableSnapshot snapshot = selector.getSnapshot().get();
+ if (snapshot.getType() == TableSnapshot.VersionType.VERSION) {
+ return
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
+ }
+ if (!state.access.isManagedVersioning()) {
+ return OptionalLong.empty();
+ }
+ long timestamp = parseTimeTravelTimestamp(snapshot.getValue());
+ return LanceSnapshotResolver.namespaceVersionAtOrBefore(
+ namespaceVersions(state, state.access, metrics), timestamp,
snapshot.getValue());
+ }
+
+ /** Resolves the selector against the open latest main chain and reads the
selected snapshot. */
+ private <T> T readFromLatest(Dataset main, ReadState state,
SnapshotReader<T> reader, LanceMetadataMetrics metrics)
+ throws Exception {
+ LanceRefSelector selector = state.selector;
+ if (selector.getTag().isPresent()) {
+ String tag = selector.getTag().get();
+ Tag target = metrics.measure(Stage.VERSION_RESOLVE, () ->
main.tags().list().stream()
+ .filter(candidate ->
tag.equals(candidate.getName())).findFirst()
+ .orElseThrow(() -> new LanceUserFacingException(
+ "Lance tag '" + tag + "' of " + state.tableName +
" was not found")));
+ state.branch = target.getBranch().filter(name ->
!MAIN_BRANCH.equals(name));
+ state.version = OptionalLong.of(target.getVersion());
+ if (!state.branch.isPresent()) {
+ try (Dataset dataset = checkout(main,
Ref.ofMain(target.getVersion()), metrics)) {
+ return reader.read(dataset, state.access, metrics);
+ }
+ }
+ }
+ if (state.branch.isPresent()) {
+ String branch = state.branch.get();
+ // Check out the branch's latest version first even when a version
is already known, so
+ // a missing branch and a missing version inside an existing
branch are told apart.
+ try (Dataset latest = checkout(main, Ref.ofBranch(branch),
metrics)) {
+ state.branchCheckedOut = true;
+ LanceTableAccess branchAccess = accessOf(latest, state);
+ if (!state.version.isPresent() &&
selector.getSnapshot().isPresent()) {
+ state.version = resolveSnapshotVersion(latest,
branchAccess, selector.getSnapshot().get(), state,
+ metrics);
+ }
+ if (!state.version.isPresent()) {
+ return reader.read(latest, branchAccess, metrics);
+ }
+ try (Dataset dataset = checkout(latest, Ref.ofBranch(branch,
state.version.getAsLong()), metrics)) {
+ return reader.read(dataset, branchAccess, metrics);
+ }
+ }
+ }
+ // FOR TIME AS OF on the main chain, resolved from storage commit
times.
+ state.version = resolveSnapshotVersion(main, state.access,
selector.getSnapshot().get(), state, metrics);
+ try (Dataset dataset = checkout(main,
Ref.ofMain(state.version.getAsLong()), metrics)) {
+ return reader.read(dataset, state.access, metrics);
+ }
+ }
+
+ /**
+ * The access for a dataset checked out from the table: the main chain
keeps the table access,
+ * and a branch takes the directory the SDK checked out, which is what the
BE opens by URI.
+ */
+ private static LanceTableAccess accessOf(Dataset dataset, ReadState state)
{
+ return state.branch.isPresent() ?
state.access.onBranch(state.branch.get(), dataset.uri()) : state.access;
+ }
+
+ /** A selector error whose message is user-facing as is, such as a tag
that does not exist. */
+ private static final class LanceUserFacingException extends
RuntimeException {
+ private LanceUserFacingException(String message) {
+ super(message);
+ }
+ }
+
+ private RuntimeException sanitizedCause(Throwable error, String uri,
Map<String, String> options) {
+ return new RuntimeException(LanceErrorMessages.sanitize(error, uri,
options, catalogSecrets));
+ }
+
+ /**
+ * Checks out a ref of an already open dataset. The SDK resolves the ref
itself, from the
+ * dataset directory or, for a namespace-managed dataset, with its own
namespace client.
+ */
+ private static Dataset checkout(Dataset dataset, Ref ref,
LanceMetadataMetrics metrics) {
+ return metrics.measure(Stage.VERSION_RESOLVE, () ->
dataset.checkout(ref));
+ }
+
+ /**
+ * Resolves a {@code FOR VERSION AS OF} / {@code FOR TIME AS OF} snapshot
against the chain
+ * {@code latest} is checked out on: the main chain, or a branch when
{@code access} is a
+ * branch access.
+ */
+ private OptionalLong resolveSnapshotVersion(Dataset latest,
LanceTableAccess access, TableSnapshot snapshot,
+ ReadState state, LanceMetadataMetrics metrics) {
+ if (snapshot.getType() == TableSnapshot.VersionType.VERSION) {
+ return
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
+ }
+ long timestamp = parseTimeTravelTimestamp(snapshot.getValue());
+ try {
+ return OptionalLong.of(resolveVersionAtOrBefore(latest, access,
timestamp, snapshot.getValue(), state,
+ metrics));
+ } catch (IllegalArgumentException e) {
+ if (!access.getBranch().isPresent()) {
+ throw e;
+ }
+ // A branch's chain starts at the version it was created from and
carries its own
+ // commit times, so an earlier timestamp has nothing to select on
the branch.
+ throw new LanceUserFacingException("Lance branch '" +
access.getBranch().get() + "' of "
+ + state.tableName + " has no version at or before '" +
snapshot.getValue()
+ + "'; a branch only holds the versions from its creation
on");
+ }
+ }
+
+ /**
+ * Whether a failed branch checkout means the branch does not exist. The
SDK reports
+ * "branch <name> does not exist", a namespace "Table branch not found",
or a missing manifest
+ * under the branch directory when nothing was ever committed there.
+ */
+ private static boolean isBranchNotFound(Throwable throwable, String
branch) {
+ if (ExceptionUtils.indexOfType(throwable,
TableBranchNotFoundException.class) >= 0) {
+ return true;
+ }
+ String rootMessage = ExceptionUtils.getRootCauseMessage(throwable);
+ if (rootMessage == null) {
+ return false;
+ }
+ String lower = rootMessage.toLowerCase(Locale.ROOT);
+ String name = branch.toLowerCase(Locale.ROOT);
+ return lower.contains("table branch not found")
+ || lower.contains("branch " + name + " does not exist")
+ || (lower.contains("not found") && lower.contains("tree/" +
name + "/"));
+ }
+
+ /**
+ * Whether a not-found came from the namespace rather than storage. The
SDK surfaces a
+ * namespace error by its display text ("Table version not found: ..."),
and the Java client
+ * by its exception type.
+ */
+ private static boolean isNamespaceMiss(Throwable throwable, String
namespaceText) {
+ if (ExceptionUtils.indexOfType(throwable,
TableVersionNotFoundException.class) >= 0
+ || ExceptionUtils.indexOfType(throwable,
TableBranchNotFoundException.class) >= 0) {
+ return true;
+ }
+ String rootMessage = ExceptionUtils.getRootCauseMessage(throwable);
+ return rootMessage != null &&
rootMessage.toLowerCase(Locale.ROOT).contains(namespaceText);
+ }
+
+ /** An HTTP 403 as the object stores report it, or an explicit
access-denied error. */
+ private static final Pattern ACCESS_DENIED = Pattern.compile(
+ "accessdenied|access denied|permission
denied|forbidden|(status|http|code)\\W{0,3}403\\b");
+
+ private static boolean isAccessDenied(Throwable throwable) {
+ String rootMessage = ExceptionUtils.getRootCauseMessage(throwable);
+ return rootMessage != null &&
ACCESS_DENIED.matcher(rootMessage.toLowerCase(Locale.ROOT)).find();
+ }
+
+ /**
+ * Every version the namespace records for the chain {@code access}
addresses, listed once per
+ * read. The whole list is needed: the storage fallback filters by it, and
neither the order a
+ * namespace returns nor monotonic commit times can be relied on to stop
early.
+ */
+ private List<TableVersion> namespaceVersions(ReadState state,
LanceTableAccess access,
+ LanceMetadataMetrics metrics) {
+ return
state.namespaceVersions.computeIfAbsent(access.getBranch().orElse(""), chain ->
{
+ List<TableVersion> versions =
metrics.measure(Stage.VERSION_RESOLVE,
+ () -> namespaceClient.listManagedVersions(access));
+ if (versions.isEmpty()) {
+ throw new LanceUserFacingException("Lance namespace lists no
versions for "
+ + state.tableName + (chain.isEmpty() ? "" : "@" +
chain));
+ }
+ return versions;
+ });
+ }
+
+ /**
+ * Resolves {@code FOR TIME AS OF} to a version on the chain {@code
latest} is checked out on.
+ * A namespace-managed table is resolved from the commit times the
namespace records, so that
+ * only versions the namespace knows are selected. If the namespace lists
its versions without
+ * commit times, the times come from the manifests present in storage,
restricted to the
+ * versions the namespace lists; a storage-versioned table is resolved
from storage alone.
+ */
+ private long resolveVersionAtOrBefore(Dataset latest, LanceTableAccess
access, long timestamp,
+ String requestedText, ReadState state, LanceMetadataMetrics
metrics) {
+ Set<Long> recordedVersions = null;
+ if (access.isManagedVersioning()) {
+ List<TableVersion> recorded = namespaceVersions(state, access,
metrics);
+ OptionalLong fromNamespace =
LanceSnapshotResolver.namespaceVersionAtOrBefore(
+ recorded, timestamp, requestedText);
+ if (fromNamespace.isPresent()) {
+ LOG.debug("Resolved Lance FOR TIME AS OF '{}' to version {}
from the namespace",
+ requestedText, fromNamespace.getAsLong());
+ return fromNamespace.getAsLong();
+ }
+ recordedVersions =
recorded.stream().map(TableVersion::getVersion).collect(Collectors.toSet());
+ }
+ Set<Long> allowedVersions = recordedVersions;
+ long version = metrics.measure(Stage.VERSION_RESOLVE, () -> {
+ List<Version> versions = latest.listVersions();
Review Comment:
[P1] Do not enumerate storage to resolve an untimed managed history. If even
one namespace version omits `timestamp_millis`, this calls
`latest.listVersions()`, but Lance v12's [namespace manifest
store](https://github.com/lance-format/lance/blob/v12.0.0/rust/lance/src/io/commit/namespace_manifest.rs)
does not implement history listing, so the external handler falls back to
canonical object-store manifests. Opening latest only finalizes that selected
manifest; an older namespace-recorded version can still exist solely at its
staged path and disappear here, making `FOR TIME AS OF` choose an older version
or fail. Resolve/finalize each namespace-recorded version before reading its
commit time (or reject untimed managed histories), and test an older staged
version with missing namespace timestamps.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,68 +250,382 @@ public LanceTableMetadata loadBasicTableMetadata(String
dbName, String tableName
}
public Schema loadTableSchema(String dbName, String tableName) {
- return readTableSnapshot(dbName, tableName, Optional.empty(),
+ return readTableSnapshot(dbName, tableName, LanceRefSelector.latest(),
(dataset, access, metrics) -> metrics.measure(Stage.SCHEMA,
dataset::getSchema));
}
public LanceTableMetadata loadTableMetadata(String dbName, String
tableName,
Optional<TableSnapshot> tableSnapshot) {
- return loadQueryMetadata(dbName, tableName, tableSnapshot,
LanceMetadataLoader.MetadataScope.WITH_INDEXES);
+ return loadTableMetadata(dbName, tableName,
LanceRefSelector.snapshot(tableSnapshot));
+ }
+
+ public LanceTableMetadata loadTableMetadata(String dbName, String
tableName, LanceRefSelector selector) {
+ return loadQueryMetadata(dbName, tableName, selector,
LanceMetadataLoader.MetadataScope.WITH_INDEXES);
}
private LanceTableMetadata loadQueryMetadata(String dbName, String
tableName,
Optional<TableSnapshot> tableSnapshot,
LanceMetadataLoader.MetadataScope mode) {
- return readTableSnapshot(dbName, tableName, tableSnapshot,
+ return loadQueryMetadata(dbName, tableName,
LanceRefSelector.snapshot(tableSnapshot), mode);
+ }
+
+ private LanceTableMetadata loadQueryMetadata(String dbName, String
tableName,
+ LanceRefSelector selector, LanceMetadataLoader.MetadataScope mode)
{
+ return readTableSnapshot(dbName, tableName, selector,
(dataset, access, metrics) ->
LanceMetadataLoader.read(dataset, access, mode, metrics));
}
- /** Pins one resource generation, resolved table access, and the Dataset
version for the whole read. */
- private <T> T readTableSnapshot(String dbName, String tableName,
Optional<TableSnapshot> tableSnapshot,
+ /**
+ * Pins one resource generation, resolved table access, and the Dataset
version for the whole read.
+ *
+ * <p>The latest version of the main chain is opened once and every other
selector is a
+ * checkout from that handle, so the SDK resolves the ref with the same
commit handler
+ * (the namespace's, for a managed table). A tag is resolved first to the
chain and version it
+ * points at, so a tag created on a branch selects that branch. The two
shortcuts that skip the
+ * latest open are an explicit version on the main chain, and {@code FOR
TIME AS OF} on a
+ * managed table whose namespace reports commit times.
+ */
+ private <T> T readTableSnapshot(String dbName, String tableName,
LanceRefSelector selector,
SnapshotReader<T> reader) {
- LanceTableAccess tableAccess = null;
+ ReadState state = new ReadState(selector, dbName + "." + tableName);
LanceMetadataMetrics metrics =
LanceMetadataMetrics.startMetadataRead();
try {
T result;
try (BufferAllocator allocator =
namespaceAllocator.newChildAllocator(
"lance-metadata-read", 0, namespaceAllocator.getLimit())) {
- tableAccess = metrics.measure(Stage.TABLE_ACCESS,
+ state.access = metrics.measure(Stage.TABLE_ACCESS,
() -> namespaceClient.resolveTableAccess(dbName,
tableName));
- OptionalLong version = OptionalLong.empty();
- if (tableSnapshot.isPresent()) {
- TableSnapshot snapshot = tableSnapshot.get();
- if (snapshot.getType() ==
TableSnapshot.VersionType.VERSION) {
- version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
- } else {
- long timestamp =
TimeUtils.timeStringToLong(snapshot.getValue(), TimeUtils.getTimeZone());
- if (timestamp < 0) {
- throw new IllegalArgumentException(
- "Cannot parse Lance FOR TIME AS OF value
'" + snapshot.getValue() + "'");
- }
- try (Dataset latest = openDataset(allocator,
tableAccess, OptionalLong.empty(), metrics)) {
- version =
OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE,
- () ->
LanceSnapshotResolver.getVersionAtOrBefore(latest, timestamp)));
- }
+ OptionalLong direct = directMainVersion(state, metrics);
+ if (direct.isPresent() || isLatestMain(selector)) {
+ state.version = direct;
+ try (Dataset dataset = openDataset(allocator,
state.access, direct, metrics)) {
+ result = reader.read(dataset, state.access, metrics);
Review Comment:
[P1] Keep managed access aligned with the SDK's second describe.
`state.access` may be a cached A location, but `openManagedDataset`
re-describes the table and can open a newly returned B location; this callback
then builds FE metadata from B while preserving A's URI/options for
`LanceScanNode` and the BE. During the configured cache window, a relocated
table can therefore plan B fragments/version and ask lance-c to open A, causing
failures or wrong data if A remains readable. The [Lance v12
builder](https://github.com/lance-format/lance/blob/v12.0.0/java/src/main/java/org/lance/OpenDatasetBuilder.java#L215-L264)
confirms the second describe/open. Resolve one authoritative access for both
sides (or derive the main-chain access from the opened dataset and matching
fresh options), and cover a managed location change while cached access is
still live.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceExternalTable.java:
##########
@@ -94,7 +93,68 @@ public LanceTableMetadata getMetadata(Optional<MvccSnapshot>
snapshot) {
@Override
public MvccSnapshot loadSnapshot(Optional<TableSnapshot> tableSnapshot,
Optional<TableScanParams> scanParams) {
- return new LanceMvccSnapshot(loadMetadata(tableSnapshot));
+ // As for Iceberg and Paimon tables, a non-numeric FOR VERSION AS OF
names a tag.
+ boolean versionIsTag = tableSnapshot.isPresent()
+ && tableSnapshot.get().getType() ==
TableSnapshot.VersionType.VERSION
+ &&
!LanceSnapshotResolver.isVersionNumber(tableSnapshot.get().getValue());
+ LanceRefSelector selector = versionIsTag
+ ? LanceRefSelector.tag(tableSnapshot.get().getValue()) :
LanceRefSelector.snapshot(tableSnapshot);
+ if (scanParams.isPresent()) {
+ TableScanParams params = scanParams.get();
+ if (params.isBranch()) {
+ String branch = refName(params);
+ // Lance calls the main chain "main"; it lives at the table
root, not under tree/.
+ if (LanceCatalogClient.MAIN_BRANCH.equals(branch)) {
+ // selector stays the main-chain one, including a tag
named in FOR VERSION AS OF.
Review Comment:
[P1] Preserve the explicit `main` constraint here. For a nonnumeric `FOR
VERSION AS OF`, `selector` is already a tag selector; leaving it unchanged lets
`readFromLatest` replace the branch with the tag's target. The committed `rel`
tag points to `dev`, so `tbl@branch(main) FOR VERSION AS OF 'rel'` silently
returns dev rows. Reject tag syntax with every explicit `@branch`, or verify
that the resolved tag belongs to main, and add this existing `rel` case to
regression coverage.
##########
docker/thirdparties/docker-compose/iceberg/scripts/lance_build_time_travel.py:
##########
@@ -0,0 +1,174 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+"""Generate the Lance time-travel regression fixture.
+
+time_travel.lance is a root table of the preinstalled Directory catalog next to
+all_types.lance. Every other fixture is compacted to a single version, so this
is the one
+dataset whose history survives: three commits, none of them cleaned up.
+
+ version 1 create row_id 1..3 tag column "v1" Lance tag v1
+ version 2 append row_id 4..6 tag column "v2" Lance tag v2
+ version 3 append row_id 7..9 tag column "v3" Lance tag v3
+
+Each version also carries a Lance tag of the same name (stored under
_refs/tags/), which the
+suites select with tbl@tag(v2). A branch "dev" (stored under tree/dev/,
metadata in
+_refs/branches/) is forked from version 2 and carries one extra commit:
+
+ dev version 3 append row_id 100 tag column "dev"
+
+so the branch's version 3 has rows 1..6 and 100, while main's version 3 has
rows 1..9; the
+suites select it with tbl@branch(dev). The main chain is not touched by the
branch. A tag "rel"
+points at the branch's version 3, so tbl@tag(rel) must read the branch, not
main's version 3.
+
+A Lance branch is a shallow clone: its manifests record the parent's location
as an absolute
+URI (Manifest.base_paths), so a branch only reads where it was created. The
committed
+branch was therefore created against the fixture's final location,
s3://warehouse/lance/
+time_travel.lance, and synced back into this directory; build() writes the
main chain and
+the tags locally, and --create-branch <uri> creates the branch at the uploaded
location:
+
+ python3 lance_build_time_travel.py --create-branch
s3://warehouse/lance/time_travel.lance \
+ --storage-option endpoint=http://127.0.0.1:19000 --storage-option
access_key_id=admin \
+ --storage-option secret_access_key=password --storage-option
region=us-east-1 \
+ --storage-option allow_http=true
+ # then sync tree/ and _refs/branches/ from that location into
preinstalled_data/lance/time_travel.lance/
+
+check() verifies the branch files and that they point at that URI without
opening the branch,
+which is not possible offline.
+
+The commits are spaced apart so that FOR TIME AS OF can land between two of
them. Lance
+stores the commit time of each version in its manifest, and that time is
whatever wall
+clock this script ran at. Doris has no SQL to read those times back, so the
regression
+suites hard-code them (test_lance_time_travel and test_lance_rest_time_travel,
together
+with their .out files). Regenerating this dataset therefore means updating
those suites
+from the times this script prints; that is why
lance_build_preinstalled_catalog.py carries
+the committed directory over as-is, like all_types.lance, instead of
rebuilding it, and only
+runs check() on it.
+
+The same directory serves three catalogs in the regression suites: the
filesystem catalog,
+the REST catalog with storage-native versions, and the REST catalog with
namespace-managed
+versions, where lance_rest_server.py answers the version endpoints from a
static list
+matching the versions written here.
+
+Run it with the writer pinned in lance_fixture_requirements.txt; the main
script's
+check_pinned_writer() enforces the pin for the whole catalog, this script
alone does not.
+
+Usage:
+ python3 lance_build_time_travel.py preinstalled_data/lance/time_travel.lance
+ python3 lance_build_time_travel.py --check
preinstalled_data/lance/time_travel.lance
+"""
+import argparse
+import shutil
+import time
+from pathlib import Path
+
+import lance
+import pyarrow as pa
+
+COMMITS = (("create", 1, 3, "v1"), ("append", 4, 6, "v2"), ("append", 7, 9,
"v3"))
+BRANCH = "dev"
+BRANCH_ROW_ID = 100
+BRANCH_TAG = "rel"
+BRANCH_PARENT_URI = "s3://warehouse/lance/time_travel.lance"
+COMMIT_GAP_SECONDS = 1.5
+
+
+def rows_of(low: int, high: int, tag: str) -> pa.Table:
+ return pa.table({
+ "row_id": pa.array(range(low, high + 1), pa.int32()),
+ "tag": pa.array([tag] * (high - low + 1), pa.string()),
+ })
+
+
+def build(output: Path) -> None:
+ if output.exists():
+ shutil.rmtree(output)
+ for index, (mode, low, high, tag) in enumerate(COMMITS):
+ if index > 0:
+ time.sleep(COMMIT_GAP_SECONDS)
+ # Match all_types.lance (data storage version 2.2) so every committed
Lance data file
+ # shares one on-disk format with the rest of the fixture.
+ lance.write_dataset(rows_of(low, high, tag), str(output), mode=mode,
+ data_storage_version="2.2")
+ dataset = lance.dataset(str(output))
+ for version, (_, _, _, tag) in zip((1, 2, 3), COMMITS):
+ dataset.tags.create(tag, version)
+ print(f"main chain and tags written; create the branch with
--create-branch {BRANCH_PARENT_URI}")
+
+
+def create_branch(uri: str, storage_options: dict) -> None:
+ """Forks the branch at the dataset's final location and appends its extra
row there."""
+ dataset = lance.dataset(uri, storage_options=storage_options)
+ assert [v["version"] for v in dataset.versions()] == [1, 2, 3], "upload
the main chain first"
+ dataset.create_branch(BRANCH, 2)
+ lance.write_dataset(rows_of(BRANCH_ROW_ID, BRANCH_ROW_ID, BRANCH),
f"{uri}/tree/{BRANCH}",
+ mode="append", data_storage_version="2.2",
storage_options=storage_options)
+ branch = lance.dataset(f"{uri}/tree/{BRANCH}",
storage_options=storage_options)
+ assert branch.version == 3 and
sorted(branch.to_table()["row_id"].to_pylist()) == [1, 2, 3, 4, 5, 6,
BRANCH_ROW_ID]
+ dataset.tags.create(BRANCH_TAG, (BRANCH, 3))
+ print(f"branch {BRANCH} and tag {BRANCH_TAG} created at {uri}; sync tree/,
_refs/branches/ and _refs/tags/ back")
+
+
+def check(output: Path) -> None:
+ dataset = lance.dataset(str(output))
+ versions = dataset.versions()
+ assert [v["version"] for v in versions] == [1, 2, 3], (
+ f"time-travel fixture must keep exactly versions 1..3: {versions}")
+ timestamps = [v["timestamp"] for v in versions]
+ assert timestamps == sorted(timestamps) and len(set(timestamps)) == 3, (
+ f"time-travel fixture commit times must be distinct and increasing:
{timestamps}")
+ assert all((b - a).total_seconds() >= 1 for a, b in zip(timestamps,
timestamps[1:])), (
+ f"time-travel fixture commits must be at least one second apart:
{timestamps}")
+ for version, (_, _, high, tag) in zip((1, 2, 3), COMMITS):
+ table = dataset.checkout_version(version).to_table().sort_by("row_id")
+ assert table["row_id"].to_pylist() == list(range(1, high + 1)), (
+ f"version {version} rows differ from expected: {table}")
+ assert table["tag"].to_pylist()[-1] == tag, f"version {version} tag
differs: {table}"
+ tags = {name: (ref["branch"], ref["version"]) for name, ref in
dataset.tags.list().items()}
+ assert tags == {"v1": (None, 1), "v2": (None, 2), "v3": (None, 3),
BRANCH_TAG: (BRANCH, 3)}, (
+ f"time-travel fixture tags differ: {tags}")
+ assert list(dataset.branches.list()) == [BRANCH], f"time-travel fixture
branches differ: {dataset.branches.list()}"
+ branch_manifests = sorted((output / "tree" / BRANCH /
"_versions").glob("*.manifest"))
+ assert len(branch_manifests) == 2, f"branch {BRANCH} must carry versions 2
and 3: {branch_manifests}"
+ for manifest in branch_manifests:
+ assert BRANCH_PARENT_URI.encode() in manifest.read_bytes(), (
+ f"{manifest} must reference the parent at {BRANCH_PARENT_URI}; a
branch created elsewhere is unreadable there")
+ for version in versions:
+ print(f"time_travel.lance version {version['version']} committed at "
+ f"{version['timestamp'].isoformat()}")
+
+
+def main() -> int:
+ parser = argparse.ArgumentParser(description=__doc__.splitlines()[0])
+ parser.add_argument("output", type=Path, help="path of time_travel.lance")
Review Comment:
[P2] Make `output` optional in `--create-branch` mode. The documented
regeneration command supplies only `--create-branch <URI>` and storage options,
but argparse requires this positional before dispatch, so the mandatory
branch/tag creation step exits without running. Use `nargs='?'` and explicitly
require `output` only for build/check modes, or update every documented caller
with a meaningful positional argument.
--
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]