zy-kkk commented on code in PR #68453:
URL: https://github.com/apache/doris/pull/68453#discussion_r4129580932
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -227,22 +241,106 @@ 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.
+ if (StringUtils.isBlank(table.getLocation())) {
+ throw new RuntimeException("Lance namespace returned no
location for managed table " + tableId);
+ }
+ if (StringUtils.isNotBlank(table.getTableUri()) &&
!StringUtils.removeEnd(table.getTableUri(), "/")
Review Comment:
Fixed in 5a682a80bbc. The check now compares the two URIs without their
query and trailing slash, so a `table_uri` that adds presigned credentials to
the location is accepted, and Doris opens that `table_uri`. A `table_uri` that
names another path, or adds a fragment (Lance ignores it, but would join a
branch directory after it), still fails; so do other spellings of the same
place, which Lance's own namespaces do not produce. An `s3+ddb` URI is now
rejected for a managed table: it commits through DynamoDB, whose handler
records and finalizes versions even on read, while the namespace already
manages the versions. Tests:
`LanceManagedAccessTest.testTableUriMayAddAQueryToTheLocation` and
`testManagedTableRejectsADynamoDbCommitUri` (run on CI).
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -227,22 +241,106 @@ 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.
+ if (StringUtils.isBlank(table.getLocation())) {
+ throw new RuntimeException("Lance namespace returned no
location for managed table " + tableId);
+ }
+ if (StringUtils.isNotBlank(table.getTableUri()) &&
!StringUtils.removeEnd(table.getTableUri(), "/")
+ .equals(StringUtils.removeEnd(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:
Agreed, and it is more than throughput: Lance's REST client sets no request
timeout, so a single stalled describe held the lock for every table of the
catalog.
Fixed in 5a682a80bbc: table describes now run outside `namespaceLock`, as
the version requests already did (this also supersedes what the PR comment said
about the lock). The lock keeps only the namespace and table listings and the
existence checks. The access cache still merges concurrent misses on one table.
Test: `LanceManagedAccessTest.testManagedDescribesRunConcurrently` (runs on CI)
has two describes that each wait for the other to start, which cannot finish
under the old lock.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,64 +248,520 @@ 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 every other
selector is a checkout
+ * from it. A tag is resolved first to the chain and version it points at,
so a tag created on a
+ * branch selects that branch. An explicit version on the main chain skips
the handle.
+ *
+ * <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}). The
+ * handle only supplies what storage holds: tag files, branch locations,
and the manifest
+ * listing that FOR TIME AS OF 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 (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; for a branch, {@link #accessOf} derives the
branch's from it. */
+ 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() && state.access.isManagedVersioning() &&
selector.getSnapshot().isPresent()
+ && selector.getSnapshot().get().getType() ==
TableSnapshot.VersionType.VERSION) {
+ // The namespace tells a missing branch from a missing version, so
the version is
+ // checked out directly, whatever state the branch's newest
version is in.
+ String branch = state.branch.get();
+ state.version = OptionalLong.of(
+
LanceSnapshotResolver.parseVersion(selector.getSnapshot().get().getValue()));
+ requireRecorded(state, state.branch, state.version.getAsLong(),
metrics);
+ state.branchExists = true;
+ try (Dataset dataset = checkout(main, Ref.ofBranch(branch,
state.version.getAsLong()), metrics)) {
+ return reader.read(dataset, accessOf(dataset, 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.
+ Ref branchHead = Ref.ofBranch(branch);
+ if (state.access.isManagedVersioning()) {
+ if (selector.getSnapshot().isPresent()) {
+ // FOR TIME AS OF selects among the versions the namespace
records and checks
+ // the one it selects, as on main, so the branch's newest
version in storage
+ // only serves to list the branch's manifests.
+ namespaceVersions(state, state.branch, metrics);
+ } else {
+ branchHead = Ref.ofBranch(branch, recordedHead(state,
Optional.of(branch), metrics));
+ }
+ state.branchExists = true;
+ }
+ try (Dataset latest = checkout(main, branchHead, 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);
+ }
+ }
+
+ 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) {
+ 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:
Confirmed: the storage manifests' commit times decide the selection, so a
version recorded at a manifest Doris does not read can move it even when that
version is not the one selected.
Fixed in 5a682a80bbc. Before selecting by time, every version the namespace
records on the chain must be at its canonical path or at a staged manifest
beside it; otherwise the read fails with the recorded and the canonical path.
Test:
`LanceManagedVersioningTest.testFinalizedManifestOutsideItsCanonicalPathFailsTheRead`
now also asks for a time before the misrecorded version's commit, on `main`
and on a branch, and expects the failure.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceManifestPaths.java:
##########
@@ -0,0 +1,141 @@
+// 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.
+
+package org.apache.doris.datasource.lance;
+
+import org.apache.commons.lang3.StringUtils;
+
+import java.io.ByteArrayOutputStream;
+import java.math.BigInteger;
+import java.nio.charset.StandardCharsets;
+import java.util.Optional;
+import java.util.regex.Pattern;
+
+/**
+ * Where a namespace-managed version's manifest must be for Doris to read that
version.
+ *
+ * <p>Doris reads a managed version as it reads any other, by the dataset URI
and the version
+ * number, and so does the BE: Lance then opens {@code
<chain>/_versions/<u64::MAX - v>.manifest},
+ * or {@code <v>.manifest} in the V1 naming scheme, where the chain is the
table root or
+ * {@code <root>/tree/<branch>}. The namespace records a manifest path for
each version. Lance's
+ * own namespaces finalize a commit in CreateTableVersion and record that
canonical path. A
+ * namespace that records the path Lance's client sends records the staged
manifest beside it,
+ * named {@code <canonical>-<id>} ({@code make_staging_manifest_path}), and
keeps it after the
+ * commit is finalized, since Lance's namespace store cannot update a record.
A recorded path
+ * anywhere else names a manifest Doris would not read. It is also what a
namespace answers once
+ * it has moved the table away from the location this read described, if the
move changed the
+ * path inside the bucket; a path is relative to its bucket or container, so a
move to another one
+ * under the same path is not seen here.
+ */
+final class LanceManifestPaths {
+
+ /** How the namespace records a version whose manifest is where Doris
reads it. */
+ enum Recorded {
+ /** At its canonical path. */
+ CANONICAL,
+ /**
+ * At a staged manifest beside the canonical path. The version may not
have been finalized
+ * yet, or was finalized after the namespace recorded the staged path.
+ */
+ STAGED
+ }
+
+ private static final String MANIFEST_EXTENSION = ".manifest";
+
+ private static final BigInteger U64_MAX =
BigInteger.ONE.shiftLeft(64).subtract(BigInteger.ONE);
+
+ /** A URL scheme; lance-io takes a single letter before the colon for a
Windows drive instead. */
+ private static final Pattern URL_SCHEME =
Pattern.compile("^[A-Za-z][A-Za-z0-9+.-]+:");
+
+ private LanceManifestPaths() {
+ }
+
+ /**
+ * How the namespace records version {@code version} of the chain of
{@code tableUri} (the
+ * table root) on {@code branch}.
+ *
+ * @throws IllegalStateException if the recorded path is neither the
canonical path nor a
+ * staged manifest beside it
+ */
+ static Recorded check(String tableUri, Optional<String> branch, long
version, String manifestPath,
+ String tableName) {
+ String chain = objectStorePath(tableUri);
+ if (branch.isPresent()) {
+ chain = (chain.isEmpty() ? "" : chain + "/") + "tree/" +
branch.get();
+ }
+ String versions = (chain.isEmpty() ? "" : chain + "/") + "_versions/";
+ // Padded by hand: String.format would use the FE's locale digits, and
Lance writes ASCII.
+ String canonical = versions +
StringUtils.leftPad(U64_MAX.subtract(BigInteger.valueOf(version)).toString(),
+ 20, '0') + MANIFEST_EXTENSION;
+ // Lance parses the recorded path first, which drops surrounding
slashes.
+ String recorded = manifestPath == null ? "" :
StringUtils.strip(manifestPath, "/");
+ for (String name : new String[] {canonical, versions + version +
MANIFEST_EXTENSION}) {
Review Comment:
Two different manifests for one version under the two naming schemes are not
something Lance produces. A dataset commits in the scheme of its latest
manifest, a namespace like Lance's Directory namespace accepts only the next
version number, and `migrate_scheme_to_v2` (lance-table `io/commit.rs`) renames
each V1 manifest to its V2 name. Both names can exist at once (an interrupted
rename, or a namespace-mode reader finalizing a staged V1 record again after
the migration), but then they hold the same manifest. Accepting the V1 name
keeps a migrated table readable when its namespace still records the old names:
the numeric open reads the V2 file, which is the same manifest.
Failing closed here would take a storage probe on every read, for a state
only a writer that bypasses Lance can create. The BE resolves the version the
same way (V2, then V1), so the FE and the BE read the same file either way. I'd
keep it as is.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -235,64 +248,520 @@ 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 every other
selector is a checkout
+ * from it. A tag is resolved first to the chain and version it points at,
so a tag created on a
+ * branch selects that branch. An explicit version on the main chain skips
the handle.
+ *
+ * <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}). The
+ * handle only supplies what storage holds: tag files, branch locations,
and the manifest
+ * listing that FOR TIME AS OF 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 (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)) {
Review Comment:
Confirmed for a namespace that keeps its own version records: it can answer
for a branch whose `main` has no readable manifest left, and the read failed
while opening `main`. (Lance's Directory namespace itself opens a branch
through the table root, so in that state it reports the branch as missing, and
Doris reports "Lance branch 'dev' ... was not found in the namespace".)
Fixed in 5a682a80bbc. A managed table's branch is now read from its own
directory, `<table>/tree/<branch>` joined as Lance's `BranchLocation` joins it,
by URI and version as the BE reads it. The namespace answers for the branch's
head, version or history first, and `main` is not opened. This also supersedes
my earlier update on the empty-main thread, which still described branch reads
as checkouts from a `main` handle. Tests:
`LanceManagedVersioningTest.testBranchReadsWithoutAReadableMainChain` deletes
every `main` manifest and reads the branch's head, an explicit version and a
time; `LanceManagedAccessTest.testBranchUriFollowsLance` (runs on CI) covers
the join. A tag still needs the table root to open, since Lance reads tag files
through an open dataset.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -227,22 +241,106 @@ 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.
+ if (StringUtils.isBlank(table.getLocation())) {
+ throw new RuntimeException("Lance namespace returned no
location for managed table " + tableId);
+ }
+ if (StringUtils.isNotBlank(table.getTableUri()) &&
!StringUtils.removeEnd(table.getTableUri(), "/")
+ .equals(StringUtils.removeEnd(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()));
+ }
+
+ /**
+ * 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.
+ */
+ private Map<String, String> storageOptions(String datasetUri, Map<String,
String> vendedOptions) {
+ return
LanceStorageOptions.fromDorisAndVendedStorageOptions(datasetUri,
storageProperties, vendedOptions);
+ }
+
+ /**
+ * Every version the namespace records for a managed chain. No page size
is requested: Lance's
+ * Directory namespace applies a limit without returning a page token,
which would silently
+ * truncate the history, while a namespace that pages on its own still
returns one.
+ */
+ List<TableVersion> listManagedVersions(LanceTableAccess access,
Optional<String> branch) {
+ List<TableVersion> result = new ArrayList<>();
+ String pageToken = null;
+ Set<String> consumedTokens = new HashSet<>();
+ do {
+ ListTableVersionsRequest request = new
ListTableVersionsRequest().id(access.getNamespaceTableId());
+ branch.ifPresent(request::branch);
+ if (pageToken != null) {
+ request.pageToken(pageToken);
+ }
+ ListTableVersionsResponse response =
namespace.listTableVersions(request);
+ if (response.getVersions() != null) {
+ result.addAll(response.getVersions());
+ }
+ pageToken = response.getPageToken();
+ if (StringUtils.isNotEmpty(pageToken) &&
!consumedTokens.add(pageToken)) {
+ throw new IllegalStateException("Lance namespace repeated a
pagination token");
+ }
+ } while (StringUtils.isNotEmpty(pageToken));
+ return result;
+ }
+
+ /**
+ * The newest version the namespace records for a managed chain, asked for
the way the Lance
+ * SDK asks when it opens the latest version: newest first, one entry.
Empty when the chain
+ * records no version.
+ */
+ Optional<TableVersion> latestManagedVersion(LanceTableAccess access,
Optional<String> branch) {
+ ListTableVersionsRequest request = new
ListTableVersionsRequest().id(access.getNamespaceTableId())
+ .descending(true).limit(1);
+ branch.ifPresent(request::branch);
+ ListTableVersionsResponse response =
namespace.listTableVersions(request);
Review Comment:
Fixed in 5a682a80bbc. The newest-version lookup now follows page tokens
until a page holds a version or the listing ends, and rejects a repeated token,
as the full listing already did. Test:
`LanceManagedAccessTest.testNewestVersionFollowsPageTokens` (runs on CI).
--
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]