github-actions[bot] commented on code in PR #68349:
URL: https://github.com/apache/doris/pull/68349#discussion_r4119626280
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -124,6 +152,55 @@ List<String> listDatabaseNames() {
return new ArrayList<>(databases);
}
+ boolean isRootDatabase(String dbName) {
+ return rootDatabase.equals(dbName);
+ }
+
+ boolean databaseExists(String dbName) {
+ if (isRootDatabase(dbName)) {
+ return true;
+ }
+ try {
+ NamespaceExistsRequest request = new
NamespaceExistsRequest().id(buildNamespaceId(dbName));
+ synchronized (namespaceLock) {
+ namespace.namespaceExists(request);
+ }
+ return true;
+ } catch (NamespaceNotFoundException e) {
+ return false;
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void createDatabase(String dbName, Map<String, String> properties) {
+ try {
+ CreateNamespaceRequest request = new CreateNamespaceRequest()
+ .id(buildNamespaceId(dbName))
+ .mode("Create")
+ .properties(properties == null ? Collections.emptyMap() :
properties);
+ synchronized (namespaceLock) {
+ namespace.createNamespace(request);
+ }
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void dropDatabase(String dbName, boolean ifExists, boolean force) {
+ try {
+ DropNamespaceRequest request = new DropNamespaceRequest()
Review Comment:
[P1] Implement or reject FORCE for populated filesystem namespaces. This
sends `behavior("Cascade")`, but the pinned Lance 11.0.0 DirectoryNamespace
manifest `drop_namespace` ignores `behavior` and returns `NamespaceNotEmpty`
whenever the namespace has child tables or namespaces
([source](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir/manifest.rs#L3360-L3416)).
Thus `DROP DATABASE ... FORCE` still fails for the populated Lance databases
this PR advertises. Provide an explicit supported cascade or reject the
unsupported promise and cover it with a real namespace test.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -198,6 +275,72 @@ boolean tableExists(String dbName, String tblName) {
}
}
+ void createTable(String dbName, String tableName, Map<String, String>
properties,
+ byte[] arrowStream) {
+ try {
+ CreateTableRequest request = new CreateTableRequest()
+ .id(buildTableId(dbName, tableName))
+ .mode("Create")
+ .properties(properties == null ? Collections.emptyMap() :
properties)
+ .storageOptions(namespaceStorageOptions);
+ synchronized (namespaceLock) {
+ namespace.createTable(request, arrowStream);
+ }
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void dropTable(String dbName, String tableName) {
Review Comment:
[P1] Handle visible root tables that have no namespace-manifest entry. With
Lance 11.0.0 defaults, root `list_tables`/`table_exists` discover an existing
`name.lance` directory, but DROP TABLE and column mutations delegate to
ManifestNamespace, whose lookup returns `TableNotFound` when that table was not
registered ([directory
path](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir.rs#L3330-L3500),
[column
paths](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir.rs#L3745-L3859),
[manifest drop
path](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir/manifest.rs#L3159-L3208)).
Doris therefore cannot mutate a table it lists; `DROP TABLE IF EXISTS`
suppresses the error and can leave the physical table intact. Use a supported
directory deletion or migration path and cover an externally created root
dataset across these mutations.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,550 @@
+// 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.doris.analysis.ColumnPosition;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.common.ErrorReport;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.datasource.lance.metadata.LanceTypeConverter;
+import org.apache.doris.datasource.operations.ExternalMetadataOps;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo;
+
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.lance.namespace.errors.NamespaceAlreadyExistsException;
+import org.lance.namespace.errors.NamespaceNotFoundException;
+import org.lance.namespace.errors.TableAlreadyExistsException;
+import org.lance.namespace.errors.TableNotFoundException;
+import org.lance.namespace.model.AddColumnsEntry;
+import org.lance.namespace.model.AlterColumnsEntry;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Doris external metadata operations backed by the Lance Namespace API. */
+public class LanceMetadataOps implements ExternalMetadataOps {
+ private static final Logger LOG =
LogManager.getLogger(LanceMetadataOps.class);
+ private static final String TABLE_COMMENT_PROPERTY = "comment";
+
+ private final LanceExternalCatalog catalog;
+
+ public LanceMetadataOps(LanceExternalCatalog catalog) {
+ this.catalog = catalog;
+ }
+
+ @Override
+ public boolean createDbImpl(String dbName, boolean ifNotExists,
Map<String, String> properties)
+ throws DdlException {
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ if (catalog.getDbNullable(dbName) != null) {
+ if (ifNotExists) {
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ if (client.databaseExists(dbName)) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ try {
+ client.createDatabase(dbName, new HashMap<>(
+
Optional.ofNullable(properties).orElse(Collections.emptyMap())));
+ return false;
+ } catch (NamespaceAlreadyExistsException e) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ @Override
+ public void afterCreateDb() {
+ catalog.resetMetaCacheNames();
+ }
+
+ @Override
+ public boolean dropDbImpl(String dbName, boolean ifExists, boolean force)
throws DdlException {
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ return execute("Failed to drop Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (db == null) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ return false;
+ }
+ String remoteDbName = db.getRemoteName();
+ if (client.isRootDatabase(remoteDbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (!client.databaseExists(remoteDbName)) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ try {
+ client.dropDatabase(remoteDbName, ifExists, force);
+ } catch (NamespaceNotFoundException e) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ return true;
+ });
+ }
+
+ @Override
+ public void afterDropDb(String dbName) {
+ Optional<ExternalDatabase<? extends ExternalTable>> db =
catalog.getDbForReplay(dbName);
+ if (db.isPresent()) {
+ catalog.unregisterDatabase(db.get().getFullName());
+ return;
+ }
+ catalog.unregisterDatabase(dbName);
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public void afterDropDbNoOp(String dbName) {
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ if (db == null) {
+ throw new DdlException("Failed to get database: '" + dbName
+ + "' in catalog: " + catalog.getName());
+ }
+ List<Column> columns = createTableInfo.getColumns();
+ validateCreateColumns(columns);
+ Schema schema = LanceTypeConverter.toArrowSchema(columns);
+ Map<String, String> properties = new HashMap<>(
+
Optional.ofNullable(createTableInfo.getProperties()).orElse(Collections.emptyMap()));
+ if (StringUtils.isNotBlank(createTableInfo.getComment())) {
+ properties.put(TABLE_COMMENT_PROPERTY,
createTableInfo.getComment());
+ }
+
+ return execute("Failed to create Lance table " + dbName + "." +
tableName, client -> {
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ if (db.getTableNullable(tableName) != null) {
+ resetTableNameCache(dbName);
+ if (db.getTableNullable(tableName) != null) {
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ }
+ try {
+ client.createTable(db.getRemoteName(), tableName, schema,
properties);
+ return false;
+ } catch (TableAlreadyExistsException e) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ private static void validateCreateColumns(List<Column> columns) throws
UserException {
+ for (Column column : columns) {
+ if (column.isAggregated()) {
+ throw new UserException("Lance columns do not support
aggregation: " + column.getName());
+ }
+ if (column.isAutoInc()) {
+ throw new UserException("Lance columns do not support
AUTO_INCREMENT: " + column.getName());
+ }
+ if (column.isGeneratedColumn()) {
+ throw new UserException("Lance columns do not support
generated columns: " + column.getName());
+ }
+ if (column.getDefaultValue() != null) {
+ throw new UserException("Lance table creation does not support
column defaults: "
+ + column.getName());
+ }
+ }
+ }
+
+ @Override
+ public void afterCreateTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ resetTableNameCache(dbName);
+ }
+
+ @Override
+ public void dropTableImpl(ExternalTable dorisTable, boolean ifExists)
throws DdlException {
+ String dbName = dorisTable.getRemoteDbName();
+ String tableName = dorisTable.getRemoteName();
+ execute("Failed to drop Lance table " + dbName + "." + tableName,
client -> {
+ try {
+ client.dropTable(dbName, tableName);
+ } catch (TableNotFoundException e) {
+ if (!ifExists) {
+
ErrorReport.reportDdlException(ErrorCode.ERR_UNKNOWN_TABLE, tableName, dbName);
+ }
+ }
+ return null;
+ });
+ }
+
+ @Override
+ public void afterDropTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ Optional<ExternalDatabase<?>> db = catalog.getDbForReplay(dbName);
+ db.ifPresent(externalDatabase ->
externalDatabase.unregisterTable(tblName));
Review Comment:
[P2] Retire cached table objects when replay cannot resolve the dropped
table name. With `lower_case_table_names=2`, a names refresh after the remote
drop can remove the case map entry before a follower replays DROP TABLE.
`unregisterTable(tblName)` then finds no table and leaves the old object in
`metaObjCache`; recreating the same remote name can reuse its stale schema.
Handle the unresolved mapping as `RefreshManager.replayRefreshTable` does, and
cover drop/recreate after a names refresh.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java:
##########
@@ -4896,6 +4899,18 @@ public static void getDdlStmt(Command command, String
dbName, TableIf table, Lis
}
}
sb.append("\n)");
+ } else if (table.getType() == TableType.LANCE_EXTERNAL_TABLE) {
+ Map<String, String> properties = new TreeMap<>(
+ ((LanceExternalTable) table).getTableProperties());
+ String tableComment = properties.remove("comment");
+ if (StringUtils.isNotBlank(tableComment)) {
+ sb.append("\nCOMMENT
").append(SqlLiteralUtils.quoteStringLiteral(tableComment));
+ }
+ if (!properties.isEmpty()) {
Review Comment:
[P2] Escape Lance property values in SHOW CREATE TABLE. For a child
namespace or REST table whose describe response includes a valid `PROPERTIES
('owner' = 'a"b')` value, `PrintableMap` wraps it without escaping and emits
`"owner" = "a"b"`, which cannot be parsed as a CREATE TABLE statement. Render
the key/value as escaped SQL literals and add a SHOW CREATE parse or replay
case with a quote in the value.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,550 @@
+// 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.doris.analysis.ColumnPosition;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.common.ErrorReport;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.datasource.lance.metadata.LanceTypeConverter;
+import org.apache.doris.datasource.operations.ExternalMetadataOps;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo;
+
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.lance.namespace.errors.NamespaceAlreadyExistsException;
+import org.lance.namespace.errors.NamespaceNotFoundException;
+import org.lance.namespace.errors.TableAlreadyExistsException;
+import org.lance.namespace.errors.TableNotFoundException;
+import org.lance.namespace.model.AddColumnsEntry;
+import org.lance.namespace.model.AlterColumnsEntry;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Doris external metadata operations backed by the Lance Namespace API. */
+public class LanceMetadataOps implements ExternalMetadataOps {
+ private static final Logger LOG =
LogManager.getLogger(LanceMetadataOps.class);
+ private static final String TABLE_COMMENT_PROPERTY = "comment";
+
+ private final LanceExternalCatalog catalog;
+
+ public LanceMetadataOps(LanceExternalCatalog catalog) {
+ this.catalog = catalog;
+ }
+
+ @Override
+ public boolean createDbImpl(String dbName, boolean ifNotExists,
Map<String, String> properties)
+ throws DdlException {
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ if (catalog.getDbNullable(dbName) != null) {
+ if (ifNotExists) {
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ if (client.databaseExists(dbName)) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ try {
+ client.createDatabase(dbName, new HashMap<>(
+
Optional.ofNullable(properties).orElse(Collections.emptyMap())));
+ return false;
+ } catch (NamespaceAlreadyExistsException e) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ @Override
+ public void afterCreateDb() {
+ catalog.resetMetaCacheNames();
+ }
+
+ @Override
+ public boolean dropDbImpl(String dbName, boolean ifExists, boolean force)
throws DdlException {
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ return execute("Failed to drop Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (db == null) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ return false;
+ }
+ String remoteDbName = db.getRemoteName();
+ if (client.isRootDatabase(remoteDbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (!client.databaseExists(remoteDbName)) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ try {
+ client.dropDatabase(remoteDbName, ifExists, force);
+ } catch (NamespaceNotFoundException e) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ return true;
+ });
+ }
+
+ @Override
+ public void afterDropDb(String dbName) {
+ Optional<ExternalDatabase<? extends ExternalTable>> db =
catalog.getDbForReplay(dbName);
+ if (db.isPresent()) {
+ catalog.unregisterDatabase(db.get().getFullName());
+ return;
+ }
+ catalog.unregisterDatabase(dbName);
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public void afterDropDbNoOp(String dbName) {
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ if (db == null) {
+ throw new DdlException("Failed to get database: '" + dbName
+ + "' in catalog: " + catalog.getName());
+ }
+ List<Column> columns = createTableInfo.getColumns();
+ validateCreateColumns(columns);
+ Schema schema = LanceTypeConverter.toArrowSchema(columns);
+ Map<String, String> properties = new HashMap<>(
+
Optional.ofNullable(createTableInfo.getProperties()).orElse(Collections.emptyMap()));
+ if (StringUtils.isNotBlank(createTableInfo.getComment())) {
+ properties.put(TABLE_COMMENT_PROPERTY,
createTableInfo.getComment());
+ }
+
+ return execute("Failed to create Lance table " + dbName + "." +
tableName, client -> {
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ if (db.getTableNullable(tableName) != null) {
+ resetTableNameCache(dbName);
+ if (db.getTableNullable(tableName) != null) {
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ }
+ try {
+ client.createTable(db.getRemoteName(), tableName, schema,
properties);
+ return false;
+ } catch (TableAlreadyExistsException e) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ private static void validateCreateColumns(List<Column> columns) throws
UserException {
+ for (Column column : columns) {
+ if (column.isAggregated()) {
+ throw new UserException("Lance columns do not support
aggregation: " + column.getName());
+ }
+ if (column.isAutoInc()) {
+ throw new UserException("Lance columns do not support
AUTO_INCREMENT: " + column.getName());
+ }
+ if (column.isGeneratedColumn()) {
+ throw new UserException("Lance columns do not support
generated columns: " + column.getName());
+ }
+ if (column.getDefaultValue() != null) {
+ throw new UserException("Lance table creation does not support
column defaults: "
+ + column.getName());
+ }
+ }
+ }
+
+ @Override
+ public void afterCreateTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ resetTableNameCache(dbName);
+ }
+
+ @Override
+ public void dropTableImpl(ExternalTable dorisTable, boolean ifExists)
throws DdlException {
+ String dbName = dorisTable.getRemoteDbName();
+ String tableName = dorisTable.getRemoteName();
+ execute("Failed to drop Lance table " + dbName + "." + tableName,
client -> {
+ try {
+ client.dropTable(dbName, tableName);
+ } catch (TableNotFoundException e) {
+ if (!ifExists) {
+
ErrorReport.reportDdlException(ErrorCode.ERR_UNKNOWN_TABLE, tableName, dbName);
+ }
+ }
+ return null;
+ });
+ }
+
+ @Override
+ public void afterDropTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ Optional<ExternalDatabase<?>> db = catalog.getDbForReplay(dbName);
+ db.ifPresent(externalDatabase ->
externalDatabase.unregisterTable(tblName));
+ }
+
+ @Override
+ public void renameTableImpl(String dbName, String oldName, String newName)
throws DdlException {
+ throw new DdlException("Lance table rename is not supported by the
pinned Lance SDK");
+ }
+
+ @Override
+ public void afterRenameTable(String dbName, String oldName, String
newName) {
+ catalog.invalidateTableAccessCache();
+ Optional<ExternalDatabase<?>> db = catalog.getDbForReplay(dbName);
+ if (db.isPresent()) {
+ db.get().unregisterTable(oldName);
+ db.get().resetMetaCacheNames();
+ }
+ }
+
+ @Override
+ public void addColumn(ExternalTable dorisTable, Column column,
ColumnPosition position, long updateTime)
+ throws UserException {
+ validateAddColumn(column, position);
+ ensureNewColumnNames(dorisTable, Collections.singletonList(column));
+ AddColumnsEntry entry = new AddColumnsEntry()
+ .name(column.getName())
+
.expression(LanceTypeConverter.toAddColumnExpression(column.getType()));
+ execute("Failed to add column " + column.getName() + " to Lance table "
+ + tableName(dorisTable),
+ client -> {
+ client.addColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+ Collections.singletonList(entry));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void addColumns(ExternalTable dorisTable, List<Column> columns,
long updateTime)
+ throws UserException {
+ if (columns.isEmpty()) {
+ return;
+ }
+ List<AddColumnsEntry> entries = new ArrayList<>(columns.size());
+ for (Column column : columns) {
+ validateAddColumn(column, null);
+ entries.add(new AddColumnsEntry()
+ .name(column.getName())
+
.expression(LanceTypeConverter.toAddColumnExpression(column.getType())));
+ }
+ ensureNewColumnNames(dorisTable, columns);
+ execute("Failed to add columns to Lance table " +
tableName(dorisTable), client -> {
+ client.addColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(), entries);
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void dropColumn(ExternalTable dorisTable, String columnName, long
updateTime)
+ throws UserException {
+ Column currentColumn = requireColumn(dorisTable, columnName);
+ execute("Failed to drop column " + currentColumn.getName() + " from
Lance table "
+ + tableName(dorisTable),
+ client -> {
+ client.dropColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+
Collections.singletonList(currentColumn.getName()));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void renameColumn(ExternalTable dorisTable, String oldName, String
newName, long updateTime)
+ throws UserException {
+ Column currentColumn = requireColumn(dorisTable, oldName);
+ Column conflictingColumn = dorisTable.getColumn(newName);
+ if (conflictingColumn != null && conflictingColumn != currentColumn) {
+ throw new UserException("Column " + newName
+ + " conflicts with an existing Lance column
(case-insensitive)");
+ }
+ if (currentColumn.getName().equals(newName)) {
+ return;
+ }
+ AlterColumnsEntry alteration = new AlterColumnsEntry()
+ .path(currentColumn.getName())
+ .rename(newName);
+ execute("Failed to rename column " + currentColumn.getName() + " to "
+ newName
+ + " in Lance table " + tableName(dorisTable),
+ client -> {
+ client.alterColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+ Collections.singletonList(alteration));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void modifyColumn(ExternalTable dorisTable, Column column,
ColumnPosition position, long updateTime)
+ throws UserException {
+ validateModifyColumn(column, position);
+ Column currentColumn = requireColumn(dorisTable, column.getName());
+ AlterColumnsEntry alteration = new
AlterColumnsEntry().path(currentColumn.getName());
+ boolean changed = false;
+ if (!currentColumn.getType().equals(column.getType())) {
+
alteration.dataType(LanceTypeConverter.toAlterColumnType(column.getType()));
Review Comment:
[P2] Preserve column metadata across a type change. Doris writes a column
comment into Arrow field metadata and reads it back for DESC/SHOW CREATE. On a
`dataType` alteration, Lance 11.0.0 replaces the field from a bare
`ArrowField::new`
([source](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance/src/dataset/schema_evolution.rs#L768-L786)),
discarding that metadata. Thus `CREATE TABLE t (c INT COMMENT 'keep')`
followed by `MODIFY COLUMN c BIGINT` silently loses `keep` after reload.
Preserve the metadata or reject commented-column casts until it can be
retained, with an SDK-backed reload case.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,550 @@
+// 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.doris.analysis.ColumnPosition;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.ErrorCode;
+import org.apache.doris.common.ErrorReport;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.datasource.lance.metadata.LanceTypeConverter;
+import org.apache.doris.datasource.operations.ExternalMetadataOps;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo;
+import
org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo;
+import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo;
+
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.lance.namespace.errors.NamespaceAlreadyExistsException;
+import org.lance.namespace.errors.NamespaceNotFoundException;
+import org.lance.namespace.errors.TableAlreadyExistsException;
+import org.lance.namespace.errors.TableNotFoundException;
+import org.lance.namespace.model.AddColumnsEntry;
+import org.lance.namespace.model.AlterColumnsEntry;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.TreeSet;
+
+/** Doris external metadata operations backed by the Lance Namespace API. */
+public class LanceMetadataOps implements ExternalMetadataOps {
+ private static final Logger LOG =
LogManager.getLogger(LanceMetadataOps.class);
+ private static final String TABLE_COMMENT_PROPERTY = "comment";
+
+ private final LanceExternalCatalog catalog;
+
+ public LanceMetadataOps(LanceExternalCatalog catalog) {
+ this.catalog = catalog;
+ }
+
+ @Override
+ public boolean createDbImpl(String dbName, boolean ifNotExists,
Map<String, String> properties)
+ throws DdlException {
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ if (catalog.getDbNullable(dbName) != null) {
+ if (ifNotExists) {
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ if (client.databaseExists(dbName)) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ }
+ try {
+ client.createDatabase(dbName, new HashMap<>(
+
Optional.ofNullable(properties).orElse(Collections.emptyMap())));
+ return false;
+ } catch (NamespaceAlreadyExistsException e) {
+ if (ifNotExists) {
+ catalog.resetMetaCacheNames();
+ return true;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS,
dbName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ @Override
+ public void afterCreateDb() {
+ catalog.resetMetaCacheNames();
+ }
+
+ @Override
+ public boolean dropDbImpl(String dbName, boolean ifExists, boolean force)
throws DdlException {
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ return execute("Failed to drop Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (db == null) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ return false;
+ }
+ String remoteDbName = db.getRemoteName();
+ if (client.isRootDatabase(remoteDbName)) {
+ throw new DdlException("Cannot drop the configured Lance root
database: " + dbName);
+ }
+ if (!client.databaseExists(remoteDbName)) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ try {
+ client.dropDatabase(remoteDbName, ifExists, force);
+ } catch (NamespaceNotFoundException e) {
+ if (ifExists) {
+ return false;
+ }
+ ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS,
dbName);
+ }
+ return true;
+ });
+ }
+
+ @Override
+ public void afterDropDb(String dbName) {
+ Optional<ExternalDatabase<? extends ExternalTable>> db =
catalog.getDbForReplay(dbName);
+ if (db.isPresent()) {
+ catalog.unregisterDatabase(db.get().getFullName());
+ return;
+ }
+ catalog.unregisterDatabase(dbName);
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public void afterDropDbNoOp(String dbName) {
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ if (db == null) {
+ throw new DdlException("Failed to get database: '" + dbName
+ + "' in catalog: " + catalog.getName());
+ }
+ List<Column> columns = createTableInfo.getColumns();
+ validateCreateColumns(columns);
+ Schema schema = LanceTypeConverter.toArrowSchema(columns);
+ Map<String, String> properties = new HashMap<>(
+
Optional.ofNullable(createTableInfo.getProperties()).orElse(Collections.emptyMap()));
+ if (StringUtils.isNotBlank(createTableInfo.getComment())) {
+ properties.put(TABLE_COMMENT_PROPERTY,
createTableInfo.getComment());
+ }
+
+ return execute("Failed to create Lance table " + dbName + "." +
tableName, client -> {
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ if (db.getTableNullable(tableName) != null) {
+ resetTableNameCache(dbName);
+ if (db.getTableNullable(tableName) != null) {
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ }
+ try {
+ client.createTable(db.getRemoteName(), tableName, schema,
properties);
+ return false;
+ } catch (TableAlreadyExistsException e) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ throw new IllegalStateException("unreachable");
+ }
+ });
+ }
+
+ private static void validateCreateColumns(List<Column> columns) throws
UserException {
+ for (Column column : columns) {
+ if (column.isAggregated()) {
+ throw new UserException("Lance columns do not support
aggregation: " + column.getName());
+ }
+ if (column.isAutoInc()) {
+ throw new UserException("Lance columns do not support
AUTO_INCREMENT: " + column.getName());
+ }
+ if (column.isGeneratedColumn()) {
+ throw new UserException("Lance columns do not support
generated columns: " + column.getName());
+ }
+ if (column.getDefaultValue() != null) {
+ throw new UserException("Lance table creation does not support
column defaults: "
+ + column.getName());
+ }
+ }
+ }
+
+ @Override
+ public void afterCreateTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ resetTableNameCache(dbName);
+ }
+
+ @Override
+ public void dropTableImpl(ExternalTable dorisTable, boolean ifExists)
throws DdlException {
+ String dbName = dorisTable.getRemoteDbName();
+ String tableName = dorisTable.getRemoteName();
+ execute("Failed to drop Lance table " + dbName + "." + tableName,
client -> {
+ try {
+ client.dropTable(dbName, tableName);
+ } catch (TableNotFoundException e) {
+ if (!ifExists) {
+
ErrorReport.reportDdlException(ErrorCode.ERR_UNKNOWN_TABLE, tableName, dbName);
+ }
+ }
+ return null;
+ });
+ }
+
+ @Override
+ public void afterDropTable(String dbName, String tblName) {
+ catalog.invalidateTableAccessCache();
+ Optional<ExternalDatabase<?>> db = catalog.getDbForReplay(dbName);
+ db.ifPresent(externalDatabase ->
externalDatabase.unregisterTable(tblName));
+ }
+
+ @Override
+ public void renameTableImpl(String dbName, String oldName, String newName)
throws DdlException {
+ throw new DdlException("Lance table rename is not supported by the
pinned Lance SDK");
+ }
+
+ @Override
+ public void afterRenameTable(String dbName, String oldName, String
newName) {
+ catalog.invalidateTableAccessCache();
+ Optional<ExternalDatabase<?>> db = catalog.getDbForReplay(dbName);
+ if (db.isPresent()) {
+ db.get().unregisterTable(oldName);
+ db.get().resetMetaCacheNames();
+ }
+ }
+
+ @Override
+ public void addColumn(ExternalTable dorisTable, Column column,
ColumnPosition position, long updateTime)
+ throws UserException {
+ validateAddColumn(column, position);
+ ensureNewColumnNames(dorisTable, Collections.singletonList(column));
+ AddColumnsEntry entry = new AddColumnsEntry()
+ .name(column.getName())
+
.expression(LanceTypeConverter.toAddColumnExpression(column.getType()));
+ execute("Failed to add column " + column.getName() + " to Lance table "
+ + tableName(dorisTable),
+ client -> {
+ client.addColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+ Collections.singletonList(entry));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void addColumns(ExternalTable dorisTable, List<Column> columns,
long updateTime)
+ throws UserException {
+ if (columns.isEmpty()) {
+ return;
+ }
+ List<AddColumnsEntry> entries = new ArrayList<>(columns.size());
+ for (Column column : columns) {
+ validateAddColumn(column, null);
+ entries.add(new AddColumnsEntry()
+ .name(column.getName())
+
.expression(LanceTypeConverter.toAddColumnExpression(column.getType())));
+ }
+ ensureNewColumnNames(dorisTable, columns);
+ execute("Failed to add columns to Lance table " +
tableName(dorisTable), client -> {
+ client.addColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(), entries);
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void dropColumn(ExternalTable dorisTable, String columnName, long
updateTime)
+ throws UserException {
+ Column currentColumn = requireColumn(dorisTable, columnName);
+ execute("Failed to drop column " + currentColumn.getName() + " from
Lance table "
+ + tableName(dorisTable),
+ client -> {
+ client.dropColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+
Collections.singletonList(currentColumn.getName()));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void renameColumn(ExternalTable dorisTable, String oldName, String
newName, long updateTime)
+ throws UserException {
+ Column currentColumn = requireColumn(dorisTable, oldName);
+ Column conflictingColumn = dorisTable.getColumn(newName);
+ if (conflictingColumn != null && conflictingColumn != currentColumn) {
+ throw new UserException("Column " + newName
+ + " conflicts with an existing Lance column
(case-insensitive)");
+ }
+ if (currentColumn.getName().equals(newName)) {
+ return;
+ }
+ AlterColumnsEntry alteration = new AlterColumnsEntry()
+ .path(currentColumn.getName())
+ .rename(newName);
+ execute("Failed to rename column " + currentColumn.getName() + " to "
+ newName
+ + " in Lance table " + tableName(dorisTable),
+ client -> {
+ client.alterColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+ Collections.singletonList(alteration));
+ return null;
+ });
+ refreshTable(dorisTable, updateTime);
+ }
+
+ @Override
+ public void modifyColumn(ExternalTable dorisTable, Column column,
ColumnPosition position, long updateTime)
+ throws UserException {
+ validateModifyColumn(column, position);
+ Column currentColumn = requireColumn(dorisTable, column.getName());
+ AlterColumnsEntry alteration = new
AlterColumnsEntry().path(currentColumn.getName());
+ boolean changed = false;
+ if (!currentColumn.getType().equals(column.getType())) {
+
alteration.dataType(LanceTypeConverter.toAlterColumnType(column.getType()));
+ changed = true;
+ }
+ if (column.isNullableSpecified()
+ && currentColumn.isAllowNull() != column.isAllowNull()) {
+ alteration.nullable(column.isAllowNull());
Review Comment:
[P1] Split or reject a combined type cast and NOT NULL change. The new
regression adds nullable `score INT` and then runs `MODIFY COLUMN score BIGINT
NOT NULL`; this code puts both `dataType` and `nullable(false)` in one Lance
alteration. Lance 11.0.0 explicitly rejects that combination and asks callers
to cast first, then change nullability
([source](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance/src/dataset/schema_evolution.rs#L829-L835)).
The advertised ALTER and regression cannot succeed through the real SDK;
handle the two steps and their failure semantics, or reject the combination
explicitly.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceExternalTable.java:
##########
@@ -56,6 +59,14 @@ public Optional<SchemaCacheValue> initSchema() {
return Optional.of(new
SchemaCacheValue(LanceSchemaHelper.toDorisColumns(schema)));
}
+ public Map<String, String> getTableProperties() {
+ DescribeTableResponse table = ((LanceExternalCatalog)
catalog).describeTable(
+ db.getRemoteName(), remoteName);
Review Comment:
[P2] Recover properties for filesystem root tables before rendering SHOW
CREATE. With the pinned Lance 11.0.0 defaults,
`DirectoryNamespace::describe_table` skips the manifest for a root table and
its directory response has no `properties` field
([source](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir.rs#L1988-L2113)).
This turns that response into an empty map, so `SHOW CREATE TABLE` silently
drops the table comment and user properties even when CREATE stored them in the
manifest. Read those properties through a supported path or make the root
describe mode return them; cover the configured root database as well as the
child database in the regression.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceCatalogClient.java:
##########
@@ -217,6 +224,71 @@ boolean tableExists(String dbName, String tableName) {
return namespaceClient.tableExists(dbName, tableName);
}
+ boolean isRootDatabase(String dbName) {
+ return namespaceClient.isRootDatabase(dbName);
+ }
+
+ boolean databaseExists(String dbName) {
+ return namespaceClient.databaseExists(dbName);
+ }
+
+ void createDatabase(String dbName, Map<String, String> properties) {
+ namespaceClient.createDatabase(dbName, properties);
+ }
+
+ void dropDatabase(String dbName, boolean ifExists, boolean force) {
+ namespaceClient.dropDatabase(dbName, ifExists, force);
+ }
+
+ void createTable(String dbName, String tableName, Schema schema,
Map<String, String> properties) {
+ namespaceClient.createTable(dbName, tableName, properties,
emptyArrowStream(schema));
+ }
+
+ void dropTable(String dbName, String tableName) {
+ namespaceClient.dropTable(dbName, tableName);
+ }
+
+ void addColumns(String dbName, String tableName, List<AddColumnsEntry>
columns) {
+ namespaceClient.addColumns(dbName, tableName, columns);
+ }
+
+ void alterColumns(String dbName, String tableName, List<AlterColumnsEntry>
alterations) {
+ namespaceClient.alterColumns(dbName, tableName, alterations);
+ }
+
+ void dropColumns(String dbName, String tableName, List<String> columns) {
+ namespaceClient.dropColumns(dbName, tableName, columns);
+ }
+
+ DescribeTableResponse describeTable(String dbName, String tableName) {
+ try {
+ return namespaceClient.describeTable(dbName, tableName, false);
+ } catch (RuntimeException e) {
+ throw LanceErrorMessages.failure("Failed to describe Lance table "
+ dbName + "." + tableName,
+ e, null, namespaceStorageOptions, catalogSecrets);
+ }
+ }
+
+ DdlException ddlFailure(String prefix, Throwable error) {
+ RuntimeException failure = LanceErrorMessages.failure(
+ prefix, error, null, namespaceStorageOptions, catalogSecrets);
+ return new DdlException(failure.getMessage(), failure.getCause());
+ }
+
+ private byte[] emptyArrowStream(Schema schema) {
+ try (BufferAllocator allocator = namespaceAllocator.newChildAllocator(
+ "lance-create-table", 0, namespaceAllocator.getLimit());
+ VectorSchemaRoot root = VectorSchemaRoot.create(schema,
allocator);
+ ByteArrayOutputStream output = new ByteArrayOutputStream();
Review Comment:
[P1] Send a zero-row Arrow batch when creating a filesystem table.
`emptyArrowStream` calls `start()` and `end()` without `writeBatch()`, and the
new test confirms `loadNextBatch()` is false. The pinned Lance 11.0.0
DirectoryNamespace forwards CREATE TABLE to its manifest implementation, which
rejects an IPC stream with no batches as `No data provided for table creation`
([source](https://github.com/lance-format/lance/blob/v11.0.0/rust/lance-namespace-impls/src/dir/manifest.rs#L3044-L3064)).
Consequently every real filesystem CREATE TABLE fails although the
request-mock test passes. Emit a zero-row batch or use a supported schema-only
API, then exercise the real namespace path.
##########
fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceMetadataOpsTest.java:
##########
@@ -0,0 +1,568 @@
+// 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.doris.analysis.ColumnPosition;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.RefreshManager;
+import org.apache.doris.catalog.Type;
+import org.apache.doris.common.DdlException;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.ExternalDatabase;
+import org.apache.doris.datasource.ExternalTable;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo;
+
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.lance.Session;
+import org.lance.namespace.LanceNamespace;
+import org.lance.namespace.errors.NamespaceAlreadyExistsException;
+import org.lance.namespace.errors.NamespaceNotFoundException;
+import org.lance.namespace.errors.TableAlreadyExistsException;
+import org.lance.namespace.errors.TableNotFoundException;
+import org.lance.namespace.model.AlterTableAddColumnsRequest;
+import org.lance.namespace.model.CreateTableRequest;
+import org.lance.namespace.model.DropNamespaceRequest;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Optional;
+
+public class LanceMetadataOpsTest {
+ @Test
+ public void testAddColumnValidation() throws UserException {
+ Column nullable = new Column("score", Type.INT, true);
+ Assertions.assertDoesNotThrow(
+ () -> LanceMetadataOps.validateAddColumn(nullable, null));
+
+ Column required = new Column("required", Type.INT, false);
+ assertRejected(() -> LanceMetadataOps.validateAddColumn(required,
null),
+ "only supports nullable columns");
+
+ Column defaulted = new Column("defaulted", Type.INT, false, null,
true, "1", "");
+ assertRejected(() -> LanceMetadataOps.validateAddColumn(defaulted,
null),
+ "does not support default values");
+
+ Column commented = new Column("commented", Type.INT, true, "comment");
+ commented.setCommentSpecified(true);
+ assertRejected(() -> LanceMetadataOps.validateAddColumn(commented,
null),
+ "does not support column comments");
+
+ assertRejected(() -> LanceMetadataOps.validateAddColumn(
+ nullable, new ColumnPosition("id")),
+ "does not support column positions");
+ }
+
+ @Test
+ public void testModifyColumnValidation() throws UserException {
+ Column column = new Column("score", Type.BIGINT, true);
+ column.setNullableSpecified(true);
+ Assertions.assertDoesNotThrow(
+ () -> LanceMetadataOps.validateModifyColumn(column, null));
+
+ column.setCommentSpecified(true);
+ assertRejected(() -> LanceMetadataOps.validateModifyColumn(column,
null),
+ "does not support column comments");
+ }
+
+ @Test
+ public void testCreateDatabaseIfNotExistsSkipsExistingDatabase() throws
DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+
+ try {
+ Assertions.assertTrue(new LanceMetadataOps(catalog).createDb(
+ "analytics", true, Collections.singletonMap("owner",
"doris")));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace,
Mockito.never()).createNamespace(Mockito.any());
+ Mockito.verify(catalog).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testCreateRootDatabaseIsRejected() {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ for (boolean ifNotExists : Arrays.asList(false, true)) {
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> ops.createDb("default", ifNotExists,
Collections.emptyMap()));
+ Assertions.assertTrue(exception.getMessage().contains("root
database"));
+ }
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace,
Mockito.never()).namespaceExists(Mockito.any());
+ Mockito.verify(namespace,
Mockito.never()).createNamespace(Mockito.any());
+ Mockito.verify(catalog, Mockito.never()).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testCreateDatabaseRejectsMappedLocalNameConflict() throws
DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("sales_db");
+ Mockito.when(database.getRemoteName()).thenReturn("Sales");
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertTrue(ops.createDb("sales_db", true,
Collections.emptyMap()));
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> ops.createDb("sales_db", false,
Collections.emptyMap()));
+ Assertions.assertTrue(exception.getMessage().contains("exist"));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace,
Mockito.never()).namespaceExists(Mockito.any());
+ Mockito.verify(namespace,
Mockito.never()).createNamespace(Mockito.any());
+ }
+
+ @Test
+ public void testCreateDatabaseHandlesConcurrentCreate() throws
DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new NamespaceNotFoundException("missing"))
+ .when(namespace).namespaceExists(Mockito.any());
+ Mockito.doThrow(new NamespaceAlreadyExistsException("created
concurrently"))
+ .when(namespace).createNamespace(Mockito.any());
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertTrue(ops.createDb("analytics", true,
Collections.emptyMap()));
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> ops.createDb("analytics", false,
Collections.emptyMap()));
+ Assertions.assertTrue(exception.getMessage().contains("exist"));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace,
Mockito.times(2)).createNamespace(Mockito.any());
+ Mockito.verify(catalog).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testCreateTableIfNotExistsSkipsExistingRemoteTable() throws
UserException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("local_db");
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database));
+ Mockito.when(database.getRemoteName()).thenReturn("analytics");
+ Mockito.when(database.getTableNullable("events")).thenReturn(null);
+ CreateTableInfo createTableInfo = createTableInfo(true);
+
+ try {
+ Assertions.assertTrue(new
LanceMetadataOps(catalog).createTable(createTableInfo));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace, Mockito.never())
+ .createTable(Mockito.any(), Mockito.any(byte[].class));
+ Mockito.verify(database).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testCreateTableHandlesConcurrentCreate() throws UserException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new TableNotFoundException("missing"))
+ .when(namespace).tableExists(Mockito.any());
+ Mockito.doThrow(new TableAlreadyExistsException("created
concurrently"))
+ .when(namespace).createTable(Mockito.any(),
Mockito.any(byte[].class));
+ LanceCatalogClient client = newClient(namespace, new
RootAllocator(1024 * 1024));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("local_db");
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database));
+ Mockito.when(database.getRemoteName()).thenReturn("analytics");
+ Mockito.when(database.getTableNullable("events")).thenReturn(null);
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertTrue(ops.createTable(createTableInfo(true)));
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> ops.createTable(createTableInfo(false)));
+ Assertions.assertTrue(exception.getMessage().contains("already
exists"));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace, Mockito.times(2))
+ .createTable(Mockito.any(), Mockito.any(byte[].class));
+ Mockito.verify(database).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testCreateTableIgnoresStaleLocalCacheWhenRemoteTableIsGone()
throws UserException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new TableNotFoundException("missing"))
+ .when(namespace).tableExists(Mockito.any());
+ LanceCatalogClient client = newClient(namespace, new
RootAllocator(1024 * 1024));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("local_db");
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database));
+ Mockito.when(database.getRemoteName()).thenReturn("analytics");
+ ExternalTable staleTable = table("local_db", "events", "analytics",
"stale_events");
+ Mockito.when(database.getTableNullable("events"))
+ .thenReturn(staleTable, null);
+
+ try {
+ Assertions.assertFalse(new
LanceMetadataOps(catalog).createTable(createTableInfo(true)));
+ } finally {
+ client.close();
+ }
+
+ ArgumentCaptor<CreateTableRequest> request =
ArgumentCaptor.forClass(CreateTableRequest.class);
+ Mockito.verify(namespace).createTable(request.capture(),
Mockito.any(byte[].class));
+ Assertions.assertEquals(Arrays.asList("tenant", "analytics", "events"),
+ request.getValue().getId());
+ Mockito.verify(database, Mockito.atLeastOnce()).resetMetaCacheNames();
+ Mockito.verify(catalog).invalidateTableAccessCache();
+ }
+
+ @Test
+ public void testCreateTableRejectsLocalNameConflictAfterRefresh() throws
UserException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new TableNotFoundException("missing"))
+ .when(namespace).tableExists(Mockito.any());
+ LanceCatalogClient client = newClient(namespace, new
RootAllocator(1024 * 1024));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("local_db");
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database));
+ Mockito.when(database.getRemoteName()).thenReturn("analytics");
+ ExternalTable conflictingTable = table("local_db", "events",
"analytics", "Events");
+
Mockito.when(database.getTableNullable("events")).thenReturn(conflictingTable);
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertTrue(ops.createTable(createTableInfo(true)));
+ DdlException exception =
Assertions.assertThrows(DdlException.class,
+ () -> ops.createTable(createTableInfo(false)));
+ Assertions.assertTrue(exception.getMessage().contains("already
exists"));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace, Mockito.never()).createTable(Mockito.any(),
Mockito.any(byte[].class));
+ Mockito.verify(database, Mockito.times(2)).resetMetaCacheNames();
+ }
+
+ @Test
+ public void testIfExistsHandlesConcurrentDropAndRefreshesLocalNames()
throws DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new NamespaceNotFoundException("dropped concurrently"))
+ .when(namespace).dropNamespace(Mockito.any());
+ Mockito.doThrow(new TableNotFoundException("dropped concurrently"))
+ .when(namespace).dropTable(Mockito.any());
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("analytics");
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database));
+ Mockito.when(database.getRemoteName()).thenReturn("analytics");
+ ExternalTable table = table("local_db", "local_table", "analytics",
"events");
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertDoesNotThrow(() -> ops.dropDb("analytics", true,
false));
+ Assertions.assertThrows(DdlException.class,
+ () -> ops.dropDb("analytics", false, false));
+ Assertions.assertDoesNotThrow(() -> ops.dropTable(table, true));
+ Assertions.assertThrows(DdlException.class,
+ () -> ops.dropTable(table, false));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(catalog,
Mockito.never()).unregisterDatabase(Mockito.anyString());
+
Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation();
+ Mockito.verify(database).unregisterTable("local_table");
+ }
+
+ @Test
+ public void testDropDatabaseReturnsFalseForIfExistsNoOp() throws
DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertFalse(ops.dropDb("missing_db", true, false));
+ Assertions.assertThrows(DdlException.class,
+ () -> ops.dropDb("missing_db", false, false));
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(namespace,
Mockito.never()).namespaceExists(Mockito.any());
+ Mockito.verify(namespace,
Mockito.never()).dropNamespace(Mockito.any());
+ Mockito.verify(catalog,
Mockito.never()).unregisterDatabase(Mockito.anyString());
+
Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Test
+ public void testDropDatabaseUsesRemoteDatabaseName() throws DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("sales_db");
+
Mockito.doReturn(Optional.of(database)).when(catalog).getDbForReplay("sales_db");
+ Mockito.when(database.getFullName()).thenReturn("sales_db");
+ Mockito.when(database.getRemoteName()).thenReturn("Sales");
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ Assertions.assertTrue(ops.dropDb("sales_db", false, true));
+ } finally {
+ client.close();
+ }
+
+ ArgumentCaptor<DropNamespaceRequest> request =
ArgumentCaptor.forClass(DropNamespaceRequest.class);
+ Mockito.verify(namespace).dropNamespace(request.capture());
+ Assertions.assertEquals(Arrays.asList("tenant", "Sales"),
request.getValue().getId());
+ Mockito.verify(catalog).unregisterDatabase("sales_db");
+ }
+
+ @Test
+ public void testAfterDropDatabaseUsesCanonicalLocalName() {
+ LanceExternalCatalog catalog =
Mockito.mock(LanceExternalCatalog.class);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+
Mockito.doReturn(Optional.of(database)).when(catalog).getDbForReplay("sales");
+ Mockito.when(database.getFullName()).thenReturn("Sales");
+
+ new LanceMetadataOps(catalog).afterDropDb("sales");
+
+ Mockito.verify(catalog).unregisterDatabase("Sales");
+ Mockito.verify(catalog,
Mockito.never()).retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Test
+ public void
testAfterDropDatabaseRetiresAllObjectsWhenCanonicalNameIsUnavailable() {
+ LanceExternalCatalog catalog =
Mockito.mock(LanceExternalCatalog.class);
+
Mockito.doReturn(Optional.empty()).when(catalog).getDbForReplay("sales");
+
+ new LanceMetadataOps(catalog).afterDropDb("sales");
+
+ Mockito.verify(catalog).unregisterDatabase("sales");
+
Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Test
+ public void
testDropAndRenameInvalidateTableAccessCacheEvenWhenReplayCacheMisses() throws
DdlException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ LanceCatalogClient client = newClient(namespace,
Mockito.mock(BufferAllocator.class));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+
Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.empty());
+ ExternalTable table = table("local_db", "local_table", "analytics",
"events");
+ LanceMetadataOps ops = new LanceMetadataOps(catalog);
+
+ try {
+ ops.dropTable(table, false);
+ ops.afterRenameTable("local_db", "local_table", "renamed_events");
+ } finally {
+ client.close();
+ }
+
+ Mockito.verify(catalog, Mockito.times(2)).invalidateTableAccessCache();
+ }
+
+ @Test
+ public void testSuccessfulNamespaceAndTableMutationsRefreshLocalNames()
throws UserException {
+ LanceNamespace namespace = Mockito.mock(LanceNamespace.class);
+ Mockito.doThrow(new NamespaceNotFoundException("missing")).doNothing()
+ .when(namespace).namespaceExists(Mockito.any());
+ Mockito.doThrow(new TableNotFoundException("missing"))
+ .when(namespace).tableExists(Mockito.any());
+ LanceCatalogClient client = newClient(namespace, new
RootAllocator(1024 * 1024));
+ LanceExternalCatalog catalog = catalogWithClient(client);
+ ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class);
+ Mockito.doReturn(database).when(catalog).getDbNullable("local_db");
+ Mockito.doReturn(database).when(catalog).getDbNullable("analytics");
Review Comment:
[P1] Correct both contradictory expectations in this new mutation test. The
`getDbNullable("analytics")` stub returns a database when line 426 calls
`createDb("analytics", false)`, so `createDbImpl` throws `ERR_DB_CREATE_EXISTS`
before creation. After that setup is fixed, line 436 still expects two
`database.resetMetaCacheNames()` calls, but successful CREATE TABLE calls it
once and DROP TABLE only calls `unregisterTable` on this mock. As written, the
focused test cannot pass; split the create/drop fixture or stage the stub, then
correct the reset count.
--
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]