github-actions[bot] commented on code in PR #68349:
URL: https://github.com/apache/doris/pull/68349#discussion_r4178113812
##########
fe/fe-core/src/test/java/org/apache/doris/datasource/lance/metadata/LanceTypeConverterTest.java:
##########
@@ -272,6 +276,125 @@ public void
testUnknownExtensionAndDictionaryMarkersAreUnsupported() {
Assertions.assertEquals(Type.UNSUPPORTED,
LanceTypeConverter.toDorisType(blobField));
}
+ @Test
+ public void testDorisSchemaRoundTripForSupportedTypes() {
+ StructType struct = new StructType(new StructField(
+ "required", Type.BOOLEAN, "nested comment", false));
+ List<Column> columns = Arrays.asList(
+ new Column("id", Type.INT, false, "identifier"),
+ new Column("tags", new ArrayType(Type.STRING), true),
+ new Column("attrs", new MapType(Type.STRING, Type.BIGINT,
false, true), true),
+ new Column("payload", struct, true),
+ new Column("document", Type.JSONB, true));
+
+ Schema schema = LanceTypeConverter.toArrowSchema(columns);
+
+ Assertions.assertEquals(5, schema.getFields().size());
+ Assertions.assertEquals("identifier",
schema.findField("id").getMetadata().get("comment"));
+ Assertions.assertFalse(schema.findField("id").isNullable());
+ Assertions.assertEquals("arrow.json",
+
schema.findField("document").getMetadata().get("ARROW:extension:name"));
+
Assertions.assertFalse(schema.findField("payload").getChildren().get(0).isNullable());
+ Assertions.assertEquals("nested comment",
+
schema.findField("payload").getChildren().get(0).getMetadata().get("comment"));
+
+ for (int i = 0; i < columns.size(); i++) {
+ Assertions.assertEquals(columns.get(i).getType(),
+ LanceTypeConverter.toDorisType(schema.getFields().get(i)));
+ }
+ }
+
+ @Test
+ public void testDorisDecimalSelectsArrowBitWidth() {
+ Schema schema = LanceTypeConverter.toArrowSchema(Arrays.asList(
+ new Column("decimal128", ScalarType.createDecimalV3Type(38,
4), true),
+ new Column("decimal256", ScalarType.createDecimalV3Type(40,
4), true)));
+
+ Assertions.assertEquals(128,
+ ((ArrowType.Decimal)
schema.findField("decimal128").getType()).getBitWidth());
+ Assertions.assertEquals(256,
+ ((ArrowType.Decimal)
schema.findField("decimal256").getType()).getBitWidth());
+ }
+
+ @Test
+ public void testAddColumnExpressionsPreserveScalarTypes() {
+ Assertions.assertEquals("CAST(NULL AS INT)",
+ LanceTypeConverter.toAddColumnExpression(Type.INT));
+ Assertions.assertEquals("CAST(NULL AS VARCHAR)",
Review Comment:
[P1] Update this expected expression to `CAST(NULL AS STRING)`.
`toAddColumnExpression(Type.STRING)` now emits the Lance-supported `STRING`
token, so this new assertion always fails before it can verify the remaining
scalar expressions. Keep the test aligned with the production parser-token fix.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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();
Review Comment:
[P2] Keep successful cold DROP replay scoped to the dropped database when
its name is known. On a follower with other warm databases but no cached object
for this target, `getDbForReplay` is empty; `unregisterDatabase(dbName)`
already invalidates the target, then this call clears every other database
object and its table caches. A routine DROP can force catalog-wide metadata
reloads. Reserve full retirement for an unresolved case mapping.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
+ throw new DdlException(
+ "CREATE TABLE is only supported in the warehouse root
namespace "
+ + "for Lance filesystem catalogs");
+ }
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ ExternalTable cachedTable = db.getTableNullable(tableName);
+ if (cachedTable != null) {
+ if (client.tableExists(db.getRemoteName(),
cachedTable.getRemoteName())) {
+ db.resetMetaCacheNames();
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ catalog.invalidateTableAccessCache();
+ db.unregisterTable(cachedTable.getName());
+ db.resetMetaCacheNames();
+ }
+ 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) {
+ validateLosslessType(column.getType(), "CREATE TABLE");
+ 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);
+ if (db.isPresent()) {
+ boolean invalidated = db.get().unregisterTableForReplay(tblName);
+ if (!invalidated && !db.get().hasLocalTableName(tblName)) {
+ db.get().retireAllTableObjectsWithoutEngineInvalidation();
+ }
+ }
+ }
+
+ @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());
+ List<AlterColumnsEntry> alterations = new ArrayList<>();
+ boolean typeChanged =
!currentColumn.getType().equals(column.getType());
+ if (typeChanged && StringUtils.isNotEmpty(currentColumn.getComment()))
{
+ throw new UserException("Lance MODIFY COLUMN does not support
changing the type "
+ + "of a column with a comment");
+ }
+ if (typeChanged) {
+ alterations.add(new AlterColumnsEntry()
+ .path(currentColumn.getName())
+
.dataType(LanceTypeConverter.toAlterColumnType(column.getType())));
+ }
+ if (column.isNullableSpecified()
+ && currentColumn.isAllowNull() != column.isAllowNull()) {
+ alterations.add(new AlterColumnsEntry()
+ .path(currentColumn.getName())
+ .nullable(column.isAllowNull()));
+ }
+ if (alterations.isEmpty()) {
+ return;
+ }
+
+ boolean modified = false;
+ DdlException mutationFailure = null;
+ try {
+ for (AlterColumnsEntry alteration : alterations) {
+ execute("Failed to modify column " + currentColumn.getName() +
" in Lance table "
+ + tableName(dorisTable),
+ client -> {
+ client.alterColumns(dorisTable.getRemoteDbName(),
dorisTable.getRemoteName(),
+ Collections.singletonList(alteration));
+ return null;
+ });
+ modified = true;
+ }
+ } catch (DdlException e) {
Review Comment:
[P1] Journal a refresh when the first MODIFY step commits but the second
fails. For `MODIFY COLUMN score BIGINT NOT NULL`, the type alteration can
succeed and the nullability alteration can fail (for example, existing NULL
rows). This catch rethrows after refreshing only the leader;
`ExternalCatalog.modifyColumn` logs the follower refresh only after this method
returns, so followers keep the old INT schema against a BIGINT Lance field.
Emit the refresh log for any committed alteration, including this failure path,
and test follower replay after the second request fails.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java:
##########
@@ -213,6 +300,102 @@ 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) {
+ try {
+ DropTableRequest request = new
DropTableRequest().id(buildTableId(dbName, tableName));
+ executeTableMutation(dbName, tableName, () ->
namespace.dropTable(request));
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void addColumns(String dbName, String tableName, List<AddColumnsEntry>
columns) {
+ try {
+ AlterTableAddColumnsRequest request = new
AlterTableAddColumnsRequest()
+ .id(buildTableId(dbName, tableName))
+ .newColumns(columns);
+ executeTableMutation(dbName, tableName, () ->
namespace.alterTableAddColumns(request));
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void alterColumns(String dbName, String tableName, List<AlterColumnsEntry>
alterations) {
+ try {
+ AlterTableAlterColumnsRequest request = new
AlterTableAlterColumnsRequest()
+ .id(buildTableId(dbName, tableName))
+ .alterations(alterations);
+ executeTableMutation(dbName, tableName, () ->
namespace.alterTableAlterColumns(request));
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ void dropColumns(String dbName, String tableName, List<String> columns) {
+ try {
+ AlterTableDropColumnsRequest request = new
AlterTableDropColumnsRequest()
+ .id(buildTableId(dbName, tableName))
+ .columns(columns);
+ executeTableMutation(dbName, tableName, () ->
namespace.alterTableDropColumns(request));
+ } catch (DdlException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void executeTableMutation(String dbName, String tableName,
Runnable mutation)
+ throws DdlException {
+ List<String> namespaceId = buildNamespaceId(dbName);
+ List<String> tableId = new ArrayList<>(namespaceId);
+ tableId.add(tableName);
+ synchronized (namespaceLock) {
+ try {
+ mutation.run();
+ return;
+ } catch (TableNotFoundException originalException) {
+ if (!LANCE_FILESYSTEM.equals(catalogType) ||
!namespaceId.isEmpty()) {
+ throw originalException;
+ }
+ try {
+ namespace.tableExists(new
TableExistsRequest().id(tableId));
+ } catch (TableNotFoundException | NamespaceNotFoundException
e) {
+ throw originalException;
+ }
+ String expectedLocation = tableName + ".lance";
+ try {
+ // DirectoryNamespace only accepts locations relative to
its warehouse root.
+ namespace.registerTable(new RegisterTableRequest()
+ .id(tableId)
+ .location(expectedLocation)
+ .mode("Create"));
+ } catch (TableAlreadyExistsException e) {
+ String registeredLocation =
describeTable(tableId).getLocation();
+ if (!expectedLocation.equals(registeredLocation)) {
Review Comment:
[P2] Compare canonical locations before rejecting a concurrent root
registration. This request registers relative `events.lance`, but Lance 12
`describe_table` returns its full warehouse URI in `location`. If another
client registers the same `events.lance` between the first mutation and this
retry, this equality check still fails and the valid DROP/ALTER is rejected.
Normalize both paths to the same URI form and make the mock return an
SDK-shaped full URI.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
+ throw new DdlException(
+ "CREATE TABLE is only supported in the warehouse root
namespace "
+ + "for Lance filesystem catalogs");
+ }
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ ExternalTable cachedTable = db.getTableNullable(tableName);
+ if (cachedTable != null) {
+ if (client.tableExists(db.getRemoteName(),
cachedTable.getRemoteName())) {
+ db.resetMetaCacheNames();
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ catalog.invalidateTableAccessCache();
+ db.unregisterTable(cachedTable.getName());
+ db.resetMetaCacheNames();
+ }
+ 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) {
+ validateLosslessType(column.getType(), "CREATE TABLE");
+ 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) {
Review Comment:
[P2] Reconcile the URI cache for a cold `DROP TABLE IF EXISTS` no-op. If the
table object was evicted and a name refresh observes an external deletion,
`ExternalCatalog.dropTable` returns before this hook or any journal entry.
Lance can still hold the previous URI in its independent access cache; a
same-name recreation at another location can then read the old dataset or fail
until TTL expiry, including on followers. Add a missing-table reconciliation
hook that invalidates access and is replayed.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
Review Comment:
[P1] Treat the configured parent as this catalog's root for filesystem
CREATE TABLE. With `lance.namespace.parent=tenant`, the exposed root database
`default` resolves to full namespace `[tenant]`; `isRootNamespace` checks that
full ID for emptiness and this guard rejects every CREATE TABLE in the catalog.
Lance 12's manifest-backed DirectoryNamespace supports child table IDs, while
the new mock tests bypass the guard by returning a null catalog type. Allow
this configured root path and cover it with a typed filesystem fixture.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
+ throw new DdlException(
+ "CREATE TABLE is only supported in the warehouse root
namespace "
+ + "for Lance filesystem catalogs");
+ }
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ ExternalTable cachedTable = db.getTableNullable(tableName);
+ if (cachedTable != null) {
+ if (client.tableExists(db.getRemoteName(),
cachedTable.getRemoteName())) {
+ db.resetMetaCacheNames();
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ catalog.invalidateTableAccessCache();
+ db.unregisterTable(cachedTable.getName());
+ db.resetMetaCacheNames();
+ }
+ 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) {
+ validateLosslessType(column.getType(), "CREATE TABLE");
+ 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);
+ if (db.isPresent()) {
+ boolean invalidated = db.get().unregisterTableForReplay(tblName);
+ if (!invalidated && !db.get().hasLocalTableName(tblName)) {
+ db.get().retireAllTableObjectsWithoutEngineInvalidation();
+ }
+ }
+ }
+
+ @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");
Review Comment:
[P1] Route REST table rename to the namespace instead of rejecting it here.
The PR advertises table rename, and Lance 12 `RestNamespace.renameTable`
implements the request through `/v1/table/{id}/rename`; this unconditional
exception makes every REST `ALTER TABLE ... RENAME` fail before the SDK. Keep
the filesystem guard where its SDK lacks rename, and exercise the REST request
plus the existing cache/replay hook.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
Review Comment:
[P1] Allow the supported filesystem namespace operations through this guard.
The PR promises database DDL for filesystem catalogs, but both CREATE DATABASE
here and DROP DATABASE at line 125 throw before contacting Lance. Pinned Lance
12 defaults to manifest mode and delegates child namespace CREATE and empty
DROP to ManifestNamespace, so even those supported operations are unavailable;
the new tests bypass this with a mock catalog type of null. Enable the
supported paths with a real filesystem fixture, and handle populated FORCE
according to the stated contract.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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();
Review Comment:
[P1] Retire a stale database object when CREATE is replayed. If an external
client deletes `analytics` and the leader recreates it, `createDbImpl` evicts
the leader's old object, but a follower receiving the CREATE log reaches this
hook and only resets the database-name list. `MetaCache.resetNames` leaves a
warm database object and its table metadata resident, so the follower can keep
resolving the previous namespace contents. Evict the matching database object
and routed caches during replay, and cover a warm follower after same-name
recreation.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/metadata/LanceTypeConverter.java:
##########
@@ -44,6 +55,254 @@ public final class LanceTypeConverter {
private LanceTypeConverter() {
}
+ /** Converts Doris columns into the Arrow schema consumed by Lance table
creation. */
+ public static Schema toArrowSchema(List<Column> columns) {
+ List<Field> fields = new ArrayList<>(columns.size());
+ for (Column column : columns) {
+ fields.add(toArrowField(column.getName(), column.getType(),
+ column.isAllowNull(), column.getComment()));
+ }
+ return new Schema(fields);
+ }
+
+ /** Builds the typed NULL expression required by the Lance Namespace
add-columns API. */
+ public static String toAddColumnExpression(Type type) {
+ String sqlType;
+ switch (type.getPrimitiveType()) {
+ case BOOLEAN:
+ sqlType = "BOOLEAN";
+ break;
+ case TINYINT:
+ sqlType = "TINYINT";
+ break;
+ case SMALLINT:
+ sqlType = "SMALLINT";
+ break;
+ case INT:
+ sqlType = "INT";
+ break;
+ case BIGINT:
+ sqlType = "BIGINT";
+ break;
+ case FLOAT:
+ sqlType = "FLOAT";
+ break;
+ case DOUBLE:
+ sqlType = "DOUBLE";
+ break;
+ case CHAR:
+ case VARCHAR:
+ case STRING:
+ sqlType = "STRING";
+ break;
+ case VARBINARY:
+ sqlType = "BINARY";
+ break;
+ case DATE:
+ case DATEV2:
+ sqlType = "DATE";
+ break;
+ case DATETIME:
+ case DATETIMEV2:
+ sqlType = "TIMESTAMP(" + supportedTemporalScale((ScalarType)
type) + ")";
+ break;
+ case DECIMALV2:
+ case DECIMAL32:
+ case DECIMAL64:
+ case DECIMAL128:
+ case DECIMAL256:
+ ScalarType decimal = (ScalarType) type;
+ sqlType = "DECIMAL(" + decimal.getScalarPrecision() + ", "
+ + decimal.getScalarScale() + ")";
+ break;
+ default:
+ throw new IllegalArgumentException(
+ "Doris type is not supported for Lance ADD COLUMN: " +
type.toSql());
+ }
+ return "CAST(NULL AS " + sqlType + ")";
+ }
+
+ /**
+ * Converts a Doris type to the scalar type name accepted by Lance
Namespace 0.7.7
+ * AlterTableAlterColumns.
+ */
+ public static String toAlterColumnType(Type type) {
+ if (type.getPrimitiveType() == PrimitiveType.JSONB) {
+ throw new IllegalArgumentException(
+ "Doris type is not supported for Lance MODIFY COLUMN: " +
type.toSql());
+ }
+ ArrowType arrowType = toArrowType(type);
+ switch (arrowType.getTypeID()) {
+ case Bool:
+ return "bool";
+ case Int:
+ ArrowType.Int integer = (ArrowType.Int) arrowType;
+ if (!integer.getIsSigned()) {
+ break;
+ }
+ return "int" + integer.getBitWidth();
+ case FloatingPoint:
+ FloatingPointPrecision precision =
+ ((ArrowType.FloatingPoint) arrowType).getPrecision();
+ if (precision == FloatingPointPrecision.SINGLE) {
+ return "float32";
+ }
+ if (precision == FloatingPointPrecision.DOUBLE) {
+ return "float64";
+ }
+ break;
+ case Utf8:
+ return "utf8";
+ case Binary:
+ return "binary";
+ case Date:
+ if (((ArrowType.Date) arrowType).getUnit() == DateUnit.DAY) {
+ return "date32";
+ }
+ break;
+ case Timestamp:
+ ArrowType.Timestamp timestamp = (ArrowType.Timestamp)
arrowType;
+ if (timestamp.getUnit() == TimeUnit.MICROSECOND
+ && StringUtils.isEmpty(timestamp.getTimezone())) {
+ return "timestamp";
+ }
+ break;
+ default:
+ break;
+ }
+ throw new IllegalArgumentException(
+ "Doris type is not supported for Lance MODIFY COLUMN: " +
type.toSql());
+ }
+
+ private static Field toArrowField(String name, Type type, boolean
nullable, String comment) {
+ if (type.getPrimitiveType() == PrimitiveType.NULL_TYPE && !nullable) {
+ throw new IllegalArgumentException("A NULL_TYPE Lance field must
be nullable: " + name);
+ }
+ Map<String, String> metadata = new HashMap<>();
+ if (comment != null && !comment.isEmpty()) {
+ metadata.put("comment", comment);
+ }
+ if (type.getPrimitiveType() == PrimitiveType.JSONB) {
+ metadata.put(ARROW_EXTENSION_NAME, ARROW_JSON_EXTENSION);
+ }
+ return new Field(name, new FieldType(nullable, toArrowType(type),
null, metadata),
+ toArrowChildren(type));
+ }
+
+ private static ArrowType toArrowType(Type type) {
+ PrimitiveType primitiveType = type.getPrimitiveType();
+ switch (primitiveType) {
+ case NULL_TYPE:
+ return ArrowType.Null.INSTANCE;
+ case BOOLEAN:
+ return ArrowType.Bool.INSTANCE;
+ case TINYINT:
+ return new ArrowType.Int(8, true);
+ case SMALLINT:
+ return new ArrowType.Int(16, true);
+ case INT:
+ return new ArrowType.Int(32, true);
+ case BIGINT:
+ return new ArrowType.Int(64, true);
+ case FLOAT:
+ return new
ArrowType.FloatingPoint(FloatingPointPrecision.SINGLE);
+ case DOUBLE:
+ return new
ArrowType.FloatingPoint(FloatingPointPrecision.DOUBLE);
+ case CHAR:
+ case VARCHAR:
+ case STRING:
+ return ArrowType.Utf8.INSTANCE;
+ case VARBINARY:
+ return ArrowType.Binary.INSTANCE;
+ case JSONB:
+ return ArrowType.Utf8.INSTANCE;
+ case DATE:
Review Comment:
[P2] Preserve or reject legacy types before CREATE commits them. When
`disable_datev1=false` or `disable_decimalv2=false` and the corresponding
conversion flag is off, `CREATE TABLE` accepts DATE, DATETIME, and DECIMALV2
here, but Arrow stores no legacy marker; the next metadata load reports DATEV2,
DATETIMEV2(0), and DecimalV3. DESC/SHOW CREATE and planning then see a
different type from the declaration. Reject those types under these settings or
persist enough metadata to round-trip them, with a reload test.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
+ throw new DdlException(
+ "CREATE TABLE is only supported in the warehouse root
namespace "
+ + "for Lance filesystem catalogs");
+ }
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ ExternalTable cachedTable = db.getTableNullable(tableName);
+ if (cachedTable != null) {
+ if (client.tableExists(db.getRemoteName(),
cachedTable.getRemoteName())) {
+ db.resetMetaCacheNames();
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ catalog.invalidateTableAccessCache();
+ db.unregisterTable(cachedTable.getName());
+ db.resetMetaCacheNames();
+ }
+ 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) {
+ validateLosslessType(column.getType(), "CREATE TABLE");
+ 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 {
Review Comment:
[P1] Do not treat a failed cold table listing as a successful Lance DROP.
Before this new implementation is called, `ExternalDatabase.buildTableForInit`
converts a non-conflict `listTables` exception into a null table, and
`ExternalCatalog.dropTable(..., true)` returns immediately. A transient
namespace failure therefore makes `DROP TABLE IF EXISTS t` report success while
`t` still exists, with no DROP request or journal record. Propagate lookup
errors or confirm remote absence before honoring IF EXISTS.
##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java:
##########
@@ -0,0 +1,655 @@
+// 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.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "CREATE DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ return execute("Failed to create Lance database " + dbName, client -> {
+ if (client.isRootDatabase(dbName)) {
+ throw new DdlException("Cannot create the configured Lance
root database: " + dbName);
+ }
+ ExternalDatabase<?> cachedDb = catalog.getDbNullable(dbName);
+ if (cachedDb != null) {
+ if (client.databaseExists(cachedDb.getRemoteName())) {
+ if (ifNotExists) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName);
+ }
+ catalog.unregisterDatabase(cachedDb.getFullName());
+ catalog.resetMetaCacheNames();
+ }
+ 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 {
+ if
(LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType())) {
+ throw new DdlException(
+ "DROP DATABASE is not supported for Lance filesystem
catalogs");
+ }
+ 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.invalidateTableAccessCache();
+ catalog.retireAllDatabaseObjectsWithoutEngineInvalidation();
+ }
+
+ @Override
+ public boolean shouldJournalDropDbNoOp() {
+ return true;
+ }
+
+ @Override
+ public boolean createTableImpl(CreateTableInfo createTableInfo) throws
UserException {
+ String dbName = createTableInfo.getDbName();
+ String tableName = createTableInfo.getTableName();
+ ExternalDatabase<?> db = catalog.getDbNullable(dbName);
+ boolean filesystemCatalog =
+
LanceExternalCatalog.LANCE_FILESYSTEM.equals(catalog.getLanceCatalogType());
+ if (filesystemCatalog && tableName.contains("$")) {
+ throw new DdlException(
+ "Lance filesystem table names must not contain the
manifest delimiter '$'");
+ }
+ 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 (filesystemCatalog &&
!client.isRootNamespace(db.getRemoteName())) {
+ throw new DdlException(
+ "CREATE TABLE is only supported in the warehouse root
namespace "
+ + "for Lance filesystem catalogs");
+ }
+ if (client.tableExists(db.getRemoteName(), tableName)) {
+ if (createTableInfo.isIfNotExists()) {
+ resetTableNameCache(dbName);
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ ExternalTable cachedTable = db.getTableNullable(tableName);
+ if (cachedTable != null) {
+ if (client.tableExists(db.getRemoteName(),
cachedTable.getRemoteName())) {
+ db.resetMetaCacheNames();
+ if (createTableInfo.isIfNotExists()) {
+ return true;
+ }
+
ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName);
+ }
+ catalog.invalidateTableAccessCache();
+ db.unregisterTable(cachedTable.getName());
+ db.resetMetaCacheNames();
+ }
+ 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) {
+ validateLosslessType(column.getType(), "CREATE TABLE");
+ 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);
Review Comment:
[P1] Retire the previous table object on CREATE replay. If an external
client deletes `events` and the leader recreates it with a new schema, the
leader evicts its stale object before CREATE, but a warm follower receiving the
CREATE log reaches this hook and resets only table names and the Lance URI
cache. `MetaCache.resetNames` leaves the old table object and schema resident,
so the follower can still plan against old columns. Invalidate the target table
object and routed caches during replay, with a two-FE recreation test.
--
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]