github-actions[bot] commented on code in PR #68453:
URL: https://github.com/apache/doris/pull/68453#discussion_r4129705352
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,64 +248,545 @@ 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>Every dataset is read by its URI and a version, as the BE reads it.
The latest version of
+ * the main chain in storage is opened once as a handle, and the other
selectors are checkouts
+ * from it, except an explicit version on the main chain and a managed
table's branch (below).
+ * A tag is resolved first to the chain and version it points at, so a tag
created on a branch
+ * selects that branch.
+ *
+ * <p>For a managed table the namespace decides which versions exist.
"Latest" is the newest
+ * version it records, never the newest manifest in storage, and every
version a read selects
+ * must be one it records, at the manifest path Doris reads ({@link
LanceManifestPaths}). A
+ * branch is read from its own directory, as the BE reads it, so it does
not depend on the main
+ * chain; the main handle only supplies tag files and the manifest listing
that FOR TIME AS OF
+ * on main takes commit times from.
+ */
+ 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 (state.access.isManagedVersioning() &&
state.branch.isPresent() && !selector.getTag().isPresent()) {
+ result = readManagedBranch(allocator, state, reader,
metrics);
+ } else 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();
+ // Lance's Directory namespace reports a branch it lacks as a
missing table. The table
+ // was described in this read, so a missing table from a branch's
version request
+ // means the branch.
+ boolean namespaceLacksBranch = access != null &&
access.isManagedVersioning()
+ && ExceptionUtils.indexOfType(e,
TableNotFoundException.class) >= 0;
+ if (state.branch.isPresent() && !state.branchExists
+ && (namespaceLacksBranch || 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" + (namespaceLacksBranch ||
isNamespaceMiss(e) ? " in the namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ if (isVersionNotFound(e) && state.pinned != null
+ && state.pinned.manifest ==
LanceManifestPaths.Recorded.STAGED) {
+ throw new RuntimeException(unreadableStaged(state.pinned,
state), 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) ? " in the
namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ throw LanceErrorMessages.failure("Failed to load Lance table
metadata for " + what, 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;
+ /** The table's access; a branch's access is derived from it with
{@code onBranch}. */
+ private LanceTableAccess access;
+ private Optional<String> branch;
+ /**
+ * Set once the branch is known to exist: the namespace recorded
versions for it, or its
+ * latest version was checked out. Later failures are not reported as
a missing branch.
+ */
+ private boolean branchExists;
+ private OptionalLong version = OptionalLong.empty();
+ /** The managed version this read opens next, as the namespace records
it. */
+ private Recorded pinned;
+ /** The namespace's version list of the chain a FOR TIME AS OF reads
("" is main), fetched once. */
+ 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("");
+ }
+ }
+
+ /** A version of a managed chain the namespace records, and how it records
its manifest. */
+ private static final class Recorded {
+ private final Optional<String> branch;
+ private final long version;
+ private final LanceManifestPaths.Recorded manifest;
+
+ private Recorded(Optional<String> branch, long version,
LanceManifestPaths.Recorded manifest) {
+ this.branch = branch;
+ this.version = version;
+ this.manifest = manifest;
+ }
+ }
+
+ 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 the latest version of a managed table, which the namespace
records.
+ */
+ private OptionalLong directMainVersion(ReadState state,
LanceMetadataMetrics metrics) {
+ LanceRefSelector selector = state.selector;
+ if (selector.getTag().isPresent() || selector.getBranch().isPresent())
{
+ return OptionalLong.empty();
+ }
+ if (!selector.getSnapshot().isPresent()) {
+ return state.access.isManagedVersioning()
+ ? OptionalLong.of(recordedHead(state, Optional.empty(),
metrics))
+ : OptionalLong.empty();
+ }
+ TableSnapshot snapshot = selector.getSnapshot().get();
+ if (snapshot.getType() != TableSnapshot.VersionType.VERSION) {
+ return OptionalLong.empty();
+ }
+ state.version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
+ requireRecorded(state, Optional.empty(), state.version.getAsLong(),
metrics);
+ return state.version;
+ }
+
+ /** 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()) {
+ // Only this tag's file is read, however many tags the table has.
The SDK checks the tag
+ // out on the branch of the version it points at.
+ String tag = selector.getTag().get();
+ state.version =
OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE, () -> tagVersion(main,
tag, state)));
+ try (Dataset target = checkout(main, Ref.ofTag(tag), metrics)) {
+ // The checkout reads the tag file again; the version it read
is the one to check.
+ state.version = OptionalLong.of(target.version());
+ state.branch = branchOf(target.uri(),
state.access.getDatasetUri());
+ requireRecorded(state, state.branch, target.version(),
metrics);
+ return reader.read(target, accessOf(target, state), 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.branchExists = 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 manifest 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);
+ }
+ }
+
+ /**
+ * Reads a branch of a managed table from the branch's directory, by URI
and version as the BE
+ * reads it. The namespace answers for the branch first, so a branch it
lacks is reported as
+ * missing, and the main chain need not be readable. FOR TIME AS OF lists
the branch's
+ * manifests from its newest version in storage and, as on main, selects
among the versions the
+ * namespace records.
+ */
+ private <T> T readManagedBranch(BufferAllocator allocator, ReadState
state, SnapshotReader<T> reader,
+ LanceMetadataMetrics metrics) throws Exception {
+ String branch = state.branch.get();
+ LanceTableAccess branchAccess = state.access.onBranch(branch,
+ branchUri(state.access.getDatasetUri(), branch));
+ Optional<TableSnapshot> snapshot = state.selector.getSnapshot();
+ if (!snapshot.isPresent()) {
+ state.version = OptionalLong.of(recordedHead(state, state.branch,
metrics));
+ } else if (snapshot.get().getType() ==
TableSnapshot.VersionType.VERSION) {
+ state.version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.get().getValue()));
+ requireRecorded(state, state.branch, state.version.getAsLong(),
metrics);
+ } else {
+ namespaceVersions(state, state.branch, metrics);
+ state.branchExists = true;
+ try (Dataset latest = openDataset(allocator, branchAccess,
OptionalLong.empty(), metrics)) {
Review Comment:
[P2] Pin the managed branch before resolving time. This new branch path
opens its unversioned physical latest before selecting a namespace-published
historical version. If the namespace publishes readable dev@2 but an
unpublished dev@3 manifest is malformed, `@branch(dev) FOR VERSION AS OF 2`
works while `FOR TIME AS OF` targeting dev@2 fails opening dev@3. Open a
namespace-recorded branch head by version first; the main time path at line 310
has the same physical-head prerequisite. Cover an unreadable unpublished latest
on each chain.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,64 +248,545 @@ 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>Every dataset is read by its URI and a version, as the BE reads it.
The latest version of
+ * the main chain in storage is opened once as a handle, and the other
selectors are checkouts
+ * from it, except an explicit version on the main chain and a managed
table's branch (below).
+ * A tag is resolved first to the chain and version it points at, so a tag
created on a branch
+ * selects that branch.
+ *
+ * <p>For a managed table the namespace decides which versions exist.
"Latest" is the newest
+ * version it records, never the newest manifest in storage, and every
version a read selects
+ * must be one it records, at the manifest path Doris reads ({@link
LanceManifestPaths}). A
+ * branch is read from its own directory, as the BE reads it, so it does
not depend on the main
+ * chain; the main handle only supplies tag files and the manifest listing
that FOR TIME AS OF
+ * on main takes commit times from.
+ */
+ 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 (state.access.isManagedVersioning() &&
state.branch.isPresent() && !selector.getTag().isPresent()) {
+ result = readManagedBranch(allocator, state, reader,
metrics);
+ } else 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();
+ // Lance's Directory namespace reports a branch it lacks as a
missing table. The table
+ // was described in this read, so a missing table from a branch's
version request
+ // means the branch.
+ boolean namespaceLacksBranch = access != null &&
access.isManagedVersioning()
+ && ExceptionUtils.indexOfType(e,
TableNotFoundException.class) >= 0;
+ if (state.branch.isPresent() && !state.branchExists
+ && (namespaceLacksBranch || 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" + (namespaceLacksBranch ||
isNamespaceMiss(e) ? " in the namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ if (isVersionNotFound(e) && state.pinned != null
+ && state.pinned.manifest ==
LanceManifestPaths.Recorded.STAGED) {
+ throw new RuntimeException(unreadableStaged(state.pinned,
state), 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) ? " in the
namespace" : ""),
+ sanitizedCause(e, uri, options));
+ }
+ throw LanceErrorMessages.failure("Failed to load Lance table
metadata for " + what, 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;
+ /** The table's access; a branch's access is derived from it with
{@code onBranch}. */
+ private LanceTableAccess access;
+ private Optional<String> branch;
+ /**
+ * Set once the branch is known to exist: the namespace recorded
versions for it, or its
+ * latest version was checked out. Later failures are not reported as
a missing branch.
+ */
+ private boolean branchExists;
+ private OptionalLong version = OptionalLong.empty();
+ /** The managed version this read opens next, as the namespace records
it. */
+ private Recorded pinned;
+ /** The namespace's version list of the chain a FOR TIME AS OF reads
("" is main), fetched once. */
+ 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("");
+ }
+ }
+
+ /** A version of a managed chain the namespace records, and how it records
its manifest. */
+ private static final class Recorded {
+ private final Optional<String> branch;
+ private final long version;
+ private final LanceManifestPaths.Recorded manifest;
+
+ private Recorded(Optional<String> branch, long version,
LanceManifestPaths.Recorded manifest) {
+ this.branch = branch;
+ this.version = version;
+ this.manifest = manifest;
+ }
+ }
+
+ 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 the latest version of a managed table, which the namespace
records.
+ */
+ private OptionalLong directMainVersion(ReadState state,
LanceMetadataMetrics metrics) {
+ LanceRefSelector selector = state.selector;
+ if (selector.getTag().isPresent() || selector.getBranch().isPresent())
{
+ return OptionalLong.empty();
+ }
+ if (!selector.getSnapshot().isPresent()) {
+ return state.access.isManagedVersioning()
+ ? OptionalLong.of(recordedHead(state, Optional.empty(),
metrics))
+ : OptionalLong.empty();
+ }
+ TableSnapshot snapshot = selector.getSnapshot().get();
+ if (snapshot.getType() != TableSnapshot.VersionType.VERSION) {
+ return OptionalLong.empty();
+ }
+ state.version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
+ requireRecorded(state, Optional.empty(), state.version.getAsLong(),
metrics);
+ return state.version;
+ }
+
+ /** 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()) {
+ // Only this tag's file is read, however many tags the table has.
The SDK checks the tag
+ // out on the branch of the version it points at.
+ String tag = selector.getTag().get();
+ state.version =
OptionalLong.of(metrics.measure(Stage.VERSION_RESOLVE, () -> tagVersion(main,
tag, state)));
+ try (Dataset target = checkout(main, Ref.ofTag(tag), metrics)) {
+ // The checkout reads the tag file again; the version it read
is the one to check.
+ state.version = OptionalLong.of(target.version());
+ state.branch = branchOf(target.uri(),
state.access.getDatasetUri());
+ requireRecorded(state, state.branch, target.version(),
metrics);
+ return reader.read(target, accessOf(target, state), 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.branchExists = 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 manifest 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);
+ }
+ }
+
+ /**
+ * Reads a branch of a managed table from the branch's directory, by URI
and version as the BE
+ * reads it. The namespace answers for the branch first, so a branch it
lacks is reported as
+ * missing, and the main chain need not be readable. FOR TIME AS OF lists
the branch's
+ * manifests from its newest version in storage and, as on main, selects
among the versions the
+ * namespace records.
+ */
+ private <T> T readManagedBranch(BufferAllocator allocator, ReadState
state, SnapshotReader<T> reader,
+ LanceMetadataMetrics metrics) throws Exception {
+ String branch = state.branch.get();
+ LanceTableAccess branchAccess = state.access.onBranch(branch,
+ branchUri(state.access.getDatasetUri(), branch));
+ Optional<TableSnapshot> snapshot = state.selector.getSnapshot();
+ if (!snapshot.isPresent()) {
+ state.version = OptionalLong.of(recordedHead(state, state.branch,
metrics));
+ } else if (snapshot.get().getType() ==
TableSnapshot.VersionType.VERSION) {
+ state.version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.get().getValue()));
+ requireRecorded(state, state.branch, state.version.getAsLong(),
metrics);
+ } else {
+ namespaceVersions(state, state.branch, metrics);
+ state.branchExists = true;
+ try (Dataset latest = openDataset(allocator, branchAccess,
OptionalLong.empty(), metrics)) {
+ state.version = resolveSnapshotVersion(latest, branchAccess,
snapshot.get(), state, metrics);
+ }
+ }
+ state.branchExists = true;
+ try (Dataset dataset = openDataset(allocator, branchAccess,
state.version, metrics)) {
+ return reader.read(dataset, branchAccess, metrics);
+ }
+ }
+
+ /**
+ * The URI of a branch's directory, joined as Lance's {@code
BranchLocation} joins it:
+ * {@code tree/<branch>} under the table root, before the URI's query
string.
+ */
+ static String branchUri(String tableUri, String branch) {
+ int query = tableUri.indexOf('?');
+ String path = query < 0 ? tableUri : tableUri.substring(0, query);
+ String joined = path + (path.endsWith("/") ? "" : "/") + "tree/" +
StringUtils.stripStart(branch, "/");
+ return query < 0 ? joined : joined + tableUri.substring(query);
+ }
+
+ private static long tagVersion(Dataset main, String tag, ReadState state) {
+ try {
+ return main.tags().getVersion(tag);
+ } catch (RuntimeException e) {
+ String rootMessage = ExceptionUtils.getRootCauseMessage(e);
+ if (rootMessage != null && rootMessage.contains("tag " + tag + "
does not exist")) {
+ throw new LanceUserFacingException("Lance tag '" + tag + "' of
" + state.tableName + " was not found");
+ }
+ throw e;
+ }
+ }
+
+ /**
+ * A dataset URI without its query and trailing slash. The query may carry
credentials, which a
+ * namespace can vend anew on every describe.
+ */
+ private static String location(String uri) {
+ return StringUtils.removeEnd(StringUtils.substringBefore(uri, "?"),
"/");
+ }
+
+ /**
+ * The branch a dataset checked out from the table root is on, from its
root directory: the
+ * table root for main, {@code <root>/tree/<branch>} otherwise. Lance
inserts the branch path
+ * before a URI's query string, so the query is compared apart. A URI that
is neither is an
+ * error rather than main, which would hand the BE the wrong chain.
+ */
+ static Optional<String> branchOf(String checkedOutUri, String tableUri) {
+ String root = location(tableUri);
+ String uri = location(checkedOutUri);
+ if (uri.equals(root)) {
+ return Optional.empty();
+ }
+ String branchRoot = root + "/tree/";
+ if (!uri.startsWith(branchRoot) || uri.length() ==
branchRoot.length()) {
+ // The URIs may carry credentials in their query, so they stay out
of the message.
+ throw new IllegalStateException("Cannot tell which branch a Lance
tag was checked out on");
+ }
+ return Optional.of(uri.substring(branchRoot.length()));
+ }
+
+ /**
+ * 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);
+ }
+ }
+
+ /**
+ * The newest version the namespace records for a managed chain, which the
read then opens.
+ * Doris asks for it itself: opening "latest" by URI would read the newest
manifest in storage,
+ * which the namespace may not have published.
+ */
+ private long recordedHead(ReadState state, Optional<String> branch,
LanceMetadataMetrics metrics) {
+ state.pinned = null;
+ Optional<TableVersion> head = metrics.measure(Stage.VERSION_RESOLVE,
+ () -> namespaceClient.latestManagedVersion(state.access,
branch));
+ if (!head.isPresent()) {
+ throw new LanceUserFacingException("Lance namespace lists no
versions for " + state.tableName
+ + branch.map(name -> "@" + name).orElse(""));
+ }
+ long version = head.get().getVersion();
+ state.pinned = new Recorded(branch, version,
LanceManifestPaths.check(state.access.getDatasetUri(), branch,
+ version, head.get().getManifestPath(), state.tableName));
+ return version;
+ }
+
+ /**
+ * Requires the namespace of a managed table to record {@code version} of
the chain on
+ * {@code branch}, at the manifest path Doris reads; nothing for a
storage-versioned table.
+ */
+ private void requireRecorded(ReadState state, Optional<String> branch,
long version,
+ LanceMetadataMetrics metrics) {
+ if (!state.access.isManagedVersioning()) {
+ return;
+ }
+ state.pinned = null;
+ TableVersion recorded = metrics.measure(Stage.VERSION_RESOLVE,
+ () -> namespaceClient.describeManagedVersion(state.access,
branch, version));
+ state.pinned = new Recorded(branch, version,
LanceManifestPaths.check(state.access.getDatasetUri(), branch,
+ version, recorded.getManifestPath(), state.tableName));
+ }
+
+ /**
+ * The error for a version the namespace records at a staged manifest
while its canonical
+ * manifest, which Doris reads, does not exist. Either the commit reserved
the version and was
+ * not finalized, which a reader that uses the namespace would finish and
Doris does not, or
+ * the version was finalized and cleanup later removed it; the namespace's
record does not tell
+ * the two apart.
+ */
+ private static String unreadableStaged(Recorded pinned, ReadState state) {
+ return "Lance version " + pinned.version + " of " + state.tableName
+ + pinned.branch.map(name -> "@" + name).orElse("") + " cannot
be read: " + stagedOnly();
+ }
+
+ private static String stagedOnly() {
+ return "the namespace records it at a staged manifest, and its
canonical manifest, which Doris reads,"
+ + " does not exist (its commit was not finalized, or cleanup
removed it)";
+ }
+
+ 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
from the dataset directory. */
+ 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) {
+ state.version =
OptionalLong.of(LanceSnapshotResolver.parseVersion(snapshot.getValue()));
+ requireRecorded(state, access.getBranch(),
state.version.getAsLong(), metrics);
+ return state.version;
+ }
+ long timestamp = parseTimeTravelTimestamp(snapshot.getValue());
+ try {
+ return OptionalLong.of(resolveVersionAtOrBefore(latest, access,
timestamp, snapshot.getValue(), state,
+ metrics));
+ } catch (LanceSnapshotResolver.NoVersionAtOrBeforeException e) {
+ if (!access.getBranch().isPresent()) {
+ throw new LanceUserFacingException("Lance table " +
state.tableName + " has no version at or before '"
+ + snapshot.getValue() + "'");
+ }
+ // 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 failure means the branch does not exist. A checkout reports
"branch <name> does
+ * not exist", or a missing manifest under the branch directory when
nothing was ever
+ * committed there; a namespace that reports the branch itself throws
+ * {@link TableBranchNotFoundException}.
+ */
+ 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("branch " + name + " does not exist")
+ || (lower.contains("not found") && lower.contains("tree/" +
name + "/"));
+ }
+
+ /** Whether a not-found came from the namespace client rather than
storage. */
+ private static boolean isNamespaceMiss(Throwable throwable) {
+ return ExceptionUtils.indexOfType(throwable,
TableVersionNotFoundException.class) >= 0
+ || ExceptionUtils.indexOfType(throwable,
TableBranchNotFoundException.class) >= 0;
+ }
+
+ /**
+ * Every version the namespace records for the chain on {@code branch},
listed once per read.
+ * The whole list is needed: FOR TIME AS OF selects among 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,
Optional<String> branch,
+ LanceMetadataMetrics metrics) {
+ return state.namespaceVersions.computeIfAbsent(branch.orElse(""),
chain -> {
+ List<TableVersion> versions =
metrics.measure(Stage.VERSION_RESOLVE,
+ () -> namespaceClient.listManagedVersions(state.access,
branch));
+ 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,
+ * from the commit times the manifests in storage record, over the history
+ * {@link LanceSnapshotResolver} describes. A managed table only selects
among the versions its
+ * namespace records. A version it no longer records between recorded ones
cuts the history
+ * like a removed one, and so does a recorded version storage lacks, since
its commit time is
+ * unknown.
+ */
+ private long resolveVersionAtOrBefore(Dataset latest, LanceTableAccess
access, long timestamp,
+ String requestedText, ReadState state, LanceMetadataMetrics
metrics) {
+ state.pinned = null;
+ Map<Long, TableVersion> records = null;
+ if (access.isManagedVersioning()) {
+ records = new HashMap<>();
+ for (TableVersion recorded : namespaceVersions(state,
access.getBranch(), metrics)) {
+ if (recorded.getVersion() != null) {
+ // Selection compares the commit times of the manifests
Doris reads, so each
+ // recorded version must be at one of them, not only the
selected one.
+ LanceManifestPaths.check(state.access.getDatasetUri(),
access.getBranch(), recorded.getVersion(),
+ recorded.getManifestPath(), state.tableName);
+ records.put(recorded.getVersion(), recorded);
+ }
+ }
+ }
+ Map<Long, TableVersion> recordedById = records;
+ NavigableSet<Long> recorded = records == null ? null : new
TreeSet<>(records.keySet());
+ long version = metrics.measure(Stage.VERSION_RESOLVE, () -> {
+ try {
+ return
LanceSnapshotResolver.versionAtOrBefore(latest.listVersions(), recorded,
timestamp,
Review Comment:
[P2] Skip unrecorded manifests when resolving managed time travel. Even
after the new recorded-path checks pass, `latest.listVersions()` deserializes
every physical manifest before `versionAtOrBefore` filters by namespace IDs. If
the namespace records v1, v2 and v4, while storage has a malformed leftover v3
and a readable latest v4, `FOR VERSION AS OF 4` works but `FOR TIME AS OF`
after v4 fails while reading the unrecorded v3. Lance v12 `Dataset::versions()`
propagates that read error. Resolve times from only namespace-recorded
manifests, or list IDs first and read only those; cover an unreadable
unrecorded intermediate version.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -227,22 +242,137 @@ private List<String> tableAccessKey(String dbName,
String tableName) {
private CachedTableAccess loadTableAccess(List<String> tableId) {
DescribeTableResponse table = describeTable(tableId);
- if (Boolean.TRUE.equals(table.getManagedVersioning())) {
- throw new UnsupportedOperationException(
- "Lance managed versioning is not supported by the current
BE reader");
+ if (Boolean.TRUE.equals(table.getIsOnlyDeclared())) {
+ throw new RuntimeException("Lance table is declared in the
namespace but has no data yet");
}
String datasetUri = StringUtils.firstNonBlank(table.getTableUri(),
table.getLocation());
if (datasetUri == null) {
throw new RuntimeException("Lance namespace returned no table URI
for " + tableId);
}
- // One option map serves both readers: the FE opens the dataset
through the Lance Java SDK
- // and the BE through lance-c, so neither can end up with credentials
the other lacks. The
- // dataset URL picks the option vocabulary, the same way Lance picks a
provider from it.
- Map<String, String> storageOptions =
LanceStorageOptions.fromDorisAndVendedStorageOptions(datasetUri,
- storageProperties, table.getStorageOptions());
- return new CachedTableAccess(new LanceTableAccess(datasetUri,
storageOptions),
- tableAccessTtlNanos(datasetUri, table.getStorageOptions()));
+ LanceTableAccess access;
+ boolean managed = Boolean.TRUE.equals(table.getManagedVersioning());
+ if (managed) {
+ // The namespace decides which versions exist; the FE and the BE
both read one of them
+ // by URI, which is all lance-c supports. The manifest paths the
namespace records are
+ // object-store paths under `location`, so `table_uri` must name
the same place; it may
+ // add a query, such as presigned credentials, which Doris then
opens it with.
+ if (StringUtils.isBlank(table.getLocation())) {
+ throw new RuntimeException("Lance namespace returned no
location for managed table " + tableId);
+ }
+ // An s3+ddb URI commits through DynamoDB, whose handler records
and finalizes versions
+ // itself, even on read; the namespace already decides the
versions of a managed table.
+ if (StringUtils.startsWithIgnoreCase(datasetUri, "s3+ddb:")) {
+ throw new RuntimeException("Lance namespace returned an s3+ddb
URI for managed table " + tableId
+ + ", whose versions the namespace manages");
+ }
+ if (StringUtils.isNotBlank(table.getTableUri())
+ &&
!withoutQuery(table.getTableUri()).equals(withoutQuery(table.getLocation()))) {
+ throw new RuntimeException("Lance namespace returned a
table_uri that differs from location for "
+ + "managed table " + tableId);
+ }
+ access = LanceTableAccess.managedByNamespace(datasetUri,
+ storageOptions(datasetUri, table.getStorageOptions()),
tableId);
+ } else {
+ access = new LanceTableAccess(datasetUri,
storageOptions(datasetUri, table.getStorageOptions()));
+ }
+ // A managed access is not cached: the version list the read asks for
next is the
+ // namespace's current one, and must be checked against the location
the namespace
+ // reports now, not against one it reported before moving the table.
+ return new CachedTableAccess(access, managed ? 0 :
tableAccessTtlNanos(datasetUri, table.getStorageOptions()));
Review Comment:
[P2] Bypass `Cache.get` for managed accesses that need a fresh describe per
read. Caffeine 2.9.3 captures the lookup time before coalescing same-key loads:
if Q1 is describing access A, the namespace rotates to B, and Q2 starts while
Q1 is still loading, Q2 can return A after Q1 inserts it because Q2 compares
A's zero-TTL expiry against its earlier lookup time. Q2 then uses revoked
credentials or the old location despite starting after the rotation. Keep
same-table overlapping managed reads from sharing a describe result; a
latch-controlled A-to-B rotation test can cover this.
--
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]