github-actions[bot] commented on code in PR #68349: URL: https://github.com/apache/doris/pull/68349#discussion_r4140827531
########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java: ########## @@ -0,0 +1,581 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance; + +import org.apache.doris.analysis.ColumnPosition; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.ErrorCode; +import org.apache.doris.common.ErrorReport; +import org.apache.doris.common.UserException; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.datasource.lance.metadata.LanceTypeConverter; +import org.apache.doris.datasource.operations.ExternalMetadataOps; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo; + +import org.apache.arrow.vector.types.pojo.Schema; +import org.apache.commons.lang3.StringUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.lance.namespace.errors.NamespaceAlreadyExistsException; +import org.lance.namespace.errors.NamespaceNotFoundException; +import org.lance.namespace.errors.TableAlreadyExistsException; +import org.lance.namespace.errors.TableNotFoundException; +import org.lance.namespace.model.AddColumnsEntry; +import org.lance.namespace.model.AlterColumnsEntry; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.TreeSet; + +/** Doris external metadata operations backed by the Lance Namespace API. */ +public class LanceMetadataOps implements ExternalMetadataOps { + private static final Logger LOG = LogManager.getLogger(LanceMetadataOps.class); + private static final String TABLE_COMMENT_PROPERTY = "comment"; + + private final LanceExternalCatalog catalog; + + public LanceMetadataOps(LanceExternalCatalog catalog) { + this.catalog = catalog; + } + + @Override + public boolean createDbImpl(String dbName, boolean ifNotExists, Map<String, String> properties) + throws DdlException { + return execute("Failed to create Lance database " + dbName, client -> { + if (client.isRootDatabase(dbName)) { + throw new DdlException("Cannot create the configured Lance root database: " + dbName); + } + if (catalog.getDbNullable(dbName) != null) { Review Comment: [P1] Check whether a cached database still exists remotely before skipping CREATE. If `analytics` was cached and then removed outside Doris, `getDbNullable` still returns its local object, so `CREATE DATABASE IF NOT EXISTS analytics` returns success without creating anything; plain CREATE falsely reports that it exists. Reconcile the cached object with `databaseExists` before treating it as a conflict or no-op. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -139,6 +170,82 @@ List<String> listDatabaseNames() { return new ArrayList<>(databases); } + boolean isRootDatabase(String dbName) { + return rootDatabase.equals(dbName); + } + + boolean databaseExists(String dbName) { + if (isRootDatabase(dbName)) { + return true; + } + try { + NamespaceExistsRequest request = new NamespaceExistsRequest().id(buildNamespaceId(dbName)); + synchronized (namespaceLock) { + namespace.namespaceExists(request); + } + return true; + } catch (NamespaceNotFoundException e) { + return false; + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void createDatabase(String dbName, Map<String, String> properties) { + try { + CreateNamespaceRequest request = new CreateNamespaceRequest() + .id(buildNamespaceId(dbName)) + .mode("Create") + .properties(properties == null ? Collections.emptyMap() : properties); + synchronized (namespaceLock) { + namespace.createNamespace(request); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void dropDatabase(String dbName, boolean ifExists, boolean force) { + try { + List<String> namespaceId = buildNamespaceId(dbName); + synchronized (namespaceLock) { + if (force) { + try { + dropNamespaceCascade(namespaceId, ifExists ? "Skip" : "Fail"); + } catch (NamespaceNotFoundException e) { Review Comment: [P1] Apply IF EXISTS only to the target namespace. If another client removes a listed child after `listChildNamespaces` but before this recursion drops it, Lance raises `NamespaceNotFound` for the child. This catch treats it as a missing parent and returns before dropping the parent; `dropDbImpl` reports success and journals a DROP while the parent remains. Reconcile or retry a vanished child, and suppress the error only after confirming the requested parent is absent. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java: ########## @@ -0,0 +1,581 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance; + +import org.apache.doris.analysis.ColumnPosition; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.ErrorCode; +import org.apache.doris.common.ErrorReport; +import org.apache.doris.common.UserException; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.datasource.lance.metadata.LanceTypeConverter; +import org.apache.doris.datasource.operations.ExternalMetadataOps; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo; + +import org.apache.arrow.vector.types.pojo.Schema; +import org.apache.commons.lang3.StringUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.lance.namespace.errors.NamespaceAlreadyExistsException; +import org.lance.namespace.errors.NamespaceNotFoundException; +import org.lance.namespace.errors.TableAlreadyExistsException; +import org.lance.namespace.errors.TableNotFoundException; +import org.lance.namespace.model.AddColumnsEntry; +import org.lance.namespace.model.AlterColumnsEntry; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.TreeSet; + +/** Doris external metadata operations backed by the Lance Namespace API. */ +public class LanceMetadataOps implements ExternalMetadataOps { + private static final Logger LOG = LogManager.getLogger(LanceMetadataOps.class); + private static final String TABLE_COMMENT_PROPERTY = "comment"; + + private final LanceExternalCatalog catalog; + + public LanceMetadataOps(LanceExternalCatalog catalog) { + this.catalog = catalog; + } + + @Override + public boolean createDbImpl(String dbName, boolean ifNotExists, Map<String, String> properties) + throws DdlException { + return execute("Failed to create Lance database " + dbName, client -> { + if (client.isRootDatabase(dbName)) { + throw new DdlException("Cannot create the configured Lance root database: " + dbName); + } + if (catalog.getDbNullable(dbName) != null) { + if (ifNotExists) { + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + } + if (client.databaseExists(dbName)) { + if (ifNotExists) { + catalog.resetMetaCacheNames(); + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + } + try { + client.createDatabase(dbName, new HashMap<>( + Optional.ofNullable(properties).orElse(Collections.emptyMap()))); + return false; + } catch (NamespaceAlreadyExistsException e) { + if (ifNotExists) { + catalog.resetMetaCacheNames(); + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + throw new IllegalStateException("unreachable"); + } + }); + } + + @Override + public void afterCreateDb() { + catalog.resetMetaCacheNames(); + } + + @Override + public boolean dropDbImpl(String dbName, boolean ifExists, boolean force) throws DdlException { + ExternalDatabase<?> db = catalog.getDbNullable(dbName); + return execute("Failed to drop Lance database " + dbName, client -> { + if (client.isRootDatabase(dbName)) { + throw new DdlException("Cannot drop the configured Lance root database: " + dbName); + } + if (db == null) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + return false; + } + String remoteDbName = db.getRemoteName(); + if (client.isRootDatabase(remoteDbName)) { + throw new DdlException("Cannot drop the configured Lance root database: " + dbName); + } + if (!client.databaseExists(remoteDbName)) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + } + try { + client.dropDatabase(remoteDbName, ifExists, force); + } catch (NamespaceNotFoundException e) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + } + return true; + }); + } + + @Override + public void afterDropDb(String dbName) { + Optional<ExternalDatabase<? extends ExternalTable>> db = catalog.getDbForReplay(dbName); + if (db.isPresent()) { + catalog.unregisterDatabase(db.get().getFullName()); + return; + } + catalog.unregisterDatabase(dbName); + catalog.retireAllDatabaseObjectsWithoutEngineInvalidation(); + } + + @Override + public void afterDropDbNoOp(String dbName) { + catalog.retireAllDatabaseObjectsWithoutEngineInvalidation(); + } + + @Override + public boolean createTableImpl(CreateTableInfo createTableInfo) throws UserException { + String dbName = createTableInfo.getDbName(); + String tableName = createTableInfo.getTableName(); + ExternalDatabase<?> db = catalog.getDbNullable(dbName); + if (db == null) { + throw new DdlException("Failed to get database: '" + dbName + + "' in catalog: " + catalog.getName()); + } + List<Column> columns = createTableInfo.getColumns(); + validateCreateColumns(columns); + Schema schema = LanceTypeConverter.toArrowSchema(columns); + Map<String, String> properties = new HashMap<>( + Optional.ofNullable(createTableInfo.getProperties()).orElse(Collections.emptyMap())); + if (StringUtils.isNotBlank(createTableInfo.getComment())) { + properties.put(TABLE_COMMENT_PROPERTY, createTableInfo.getComment()); + } + + return execute("Failed to create Lance table " + dbName + "." + tableName, client -> { + if (client.tableExists(db.getRemoteName(), tableName)) { + if (createTableInfo.isIfNotExists()) { + resetTableNameCache(dbName); + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_TABLE_EXISTS_ERROR, tableName); + } + if (db.getTableNullable(tableName) != null) { + resetTableNameCache(dbName); Review Comment: [P1] Evict a stale table object before rejecting CREATE. If `events` was cached and then deleted outside Doris, `tableExists` returns false but both `getTableNullable` calls return the same cached object: `resetMetaCacheNames()` only clears the names snapshot. `CREATE TABLE IF NOT EXISTS events` then silently leaves the table absent (and plain CREATE reports it exists). Retire the cached object when the remote check proves it is gone, and exercise this with a real MetaCache-backed fixture. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -213,6 +323,96 @@ 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)) Review Comment: [P1] Reject table names containing the Lance manifest delimiter before CREATE. Doris allows `t$1`, and this request forwards it unchanged; Lance 12 stores child `analytics.t$1` as object ID `analytics$t$1`, but its child `list_tables` excludes names with another `$` after the namespace prefix. CREATE succeeds, then SHOW TABLES and Doris name resolution lose the table. Reject `$` for filesystem tables until the SDK encodes name components, and cover a child namespace round trip. ########## 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 = "REAL"; + break; + case DOUBLE: + sqlType = "DOUBLE"; + break; + case CHAR: + case VARCHAR: + case STRING: + sqlType = "VARCHAR"; + 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() + ", " Review Comment: [P1] Handle DECIMAL256 before constructing this ADD expression. Doris accepts `DECIMAL(40,4)` and CREATE maps it to Arrow Decimal256, but Lance 12 parses `CAST(NULL AS DECIMAL(40,4))` into Decimal128, whose maximum precision is 38. ADD therefore cannot create the requested schema (or fails when Arrow validates it). Use a Decimal256-capable namespace path or reject this ADD type before remote mutation, and cover precision above 38 with the pinned SDK. ########## regression-test/suites/external_table_p0/lance/test_lance_ddl.groovy: ########## @@ -0,0 +1,133 @@ +// 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. + +suite("test_lance_ddl", "p0,external") { + String enabled = context.config.otherConfigs.get("enableIcebergTest") + if (enabled == null || !enabled.equalsIgnoreCase("true")) { + logger.info("disable Lance DDL test because the Iceberg MinIO environment is disabled.") + return + } + + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_lance_ddl" + String databaseName = "lance_ddl_db" + String tableName = "events" + + sql """DROP CATALOG IF EXISTS `${catalogName}`""" + try { + sql """ + CREATE CATALOG `${catalogName}` PROPERTIES ( + "type" = "lance", + "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance_ddl", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true" + ) + """ + + sql """DROP DATABASE IF EXISTS `${catalogName}`.`${databaseName}` FORCE""" + sql """CREATE DATABASE `${catalogName}`.`${databaseName}`""" + sql """ + CREATE TABLE `${catalogName}`.`${databaseName}`.`${tableName}` ( + id INT NOT NULL COMMENT 'identifier', + name STRING NULL, + tags ARRAY<STRING> NULL, + amount DECIMAL(10, 2) NULL + ) ENGINE=LANCE + COMMENT 'Lance DDL regression table' + PROPERTIES ("owner" = "doris") + """ + + def columns = sql """DESC `${catalogName}`.`${databaseName}`.`${tableName}`""" + assertEquals(["id", "name", "tags", "amount"], columns.collect { it[0] }) + assertEquals("int", columns[0][1].toString().toLowerCase()) + assertEquals("array<text>", columns[2][1].toString().toLowerCase()) + + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + ADD COLUMN score INT NULL + """ + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + ADD COLUMN ( + active BOOLEAN NULL, + note STRING NULL + ) + """ + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + MODIFY COLUMN score BIGINT NOT NULL + """ + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + RENAME COLUMN score TO ranking + """ + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + DROP COLUMN note + """ + + columns = sql """DESC `${catalogName}`.`${databaseName}`.`${tableName}`""" + assertEquals( + ["id", "name", "tags", "amount", "ranking", "active"], + columns.collect { it[0] }) + assertEquals("bigint", columns[4][1].toString().toLowerCase()) + assertEquals("NO", columns[4][2]) + + test { + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + ADD COLUMN positioned INT NULL FIRST + """ + exception "does not support column positions" + } + test { + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + ADD COLUMN defaulted INT NULL DEFAULT '1' + """ + exception "does not support default values" + } + test { + sql """ + ALTER TABLE `${catalogName}`.`${databaseName}`.`${tableName}` + MODIFY COLUMN ranking BIGINT NOT NULL COMMENT 'rank' + """ + exception "does not support column comments" + } + + String showCreate = sql("""SHOW CREATE TABLE `${catalogName}`.`${databaseName}`.`${tableName}`""")[0][1] + assertTrue(showCreate.contains("ENGINE=LANCE")) + assertTrue(showCreate.contains("COMMENT 'Lance DDL regression table'")) Review Comment: [P1] Match the new SHOW CREATE comment quoting in this regression. The child table's stored `comment` property is rendered through `SqlLiteralUtils.quoteStringLiteral`, which emits `COMMENT "Lance DDL regression table"`; this assertion requires single quotes and therefore fails when the test reaches SHOW CREATE. Assert the actual escaped literal or parse the emitted statement and check that its comment round trips. ########## fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceMetadataOpsTest.java: ########## @@ -0,0 +1,652 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance; + +import org.apache.doris.analysis.ColumnPosition; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.RefreshManager; +import org.apache.doris.catalog.Type; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.UserException; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo; + +import org.apache.arrow.memory.BufferAllocator; +import org.apache.arrow.memory.RootAllocator; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.lance.Session; +import org.lance.namespace.LanceNamespace; +import org.lance.namespace.errors.NamespaceAlreadyExistsException; +import org.lance.namespace.errors.NamespaceNotFoundException; +import org.lance.namespace.errors.TableAlreadyExistsException; +import org.lance.namespace.errors.TableNotFoundException; +import org.lance.namespace.model.AlterTableAddColumnsRequest; +import org.lance.namespace.model.AlterTableAlterColumnsRequest; +import org.lance.namespace.model.CreateTableRequest; +import org.lance.namespace.model.DropNamespaceRequest; +import org.mockito.ArgumentCaptor; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +import java.util.Arrays; +import java.util.Collections; +import java.util.Optional; + +public class LanceMetadataOpsTest { + @Test + public void testAddColumnValidation() throws UserException { + Column nullable = new Column("score", Type.INT, true); + Assertions.assertDoesNotThrow( + () -> LanceMetadataOps.validateAddColumn(nullable, null)); + + Column required = new Column("required", Type.INT, false); + assertRejected(() -> LanceMetadataOps.validateAddColumn(required, null), + "only supports nullable columns"); + + Column defaulted = new Column("defaulted", Type.INT, false, null, true, "1", ""); + assertRejected(() -> LanceMetadataOps.validateAddColumn(defaulted, null), + "does not support default values"); + + Column commented = new Column("commented", Type.INT, true, "comment"); + commented.setCommentSpecified(true); + assertRejected(() -> LanceMetadataOps.validateAddColumn(commented, null), + "does not support column comments"); + + assertRejected(() -> LanceMetadataOps.validateAddColumn( + nullable, new ColumnPosition("id")), + "does not support column positions"); + } + + @Test + public void testModifyColumnValidation() throws UserException { + Column column = new Column("score", Type.BIGINT, true); + column.setNullableSpecified(true); + Assertions.assertDoesNotThrow( + () -> LanceMetadataOps.validateModifyColumn(column, null)); + + column.setCommentSpecified(true); + assertRejected(() -> LanceMetadataOps.validateModifyColumn(column, null), + "does not support column comments"); + } + + @Test + public void testCreateDatabaseIfNotExistsSkipsExistingDatabase() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + + try { + Assertions.assertTrue(new LanceMetadataOps(catalog).createDb( + "analytics", true, Collections.singletonMap("owner", "doris"))); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + Mockito.verify(catalog).resetMetaCacheNames(); + } + + @Test + public void testCreateRootDatabaseIsRejected() { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + for (boolean ifNotExists : Arrays.asList(false, true)) { + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("default", ifNotExists, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("root database")); + } + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).namespaceExists(Mockito.any()); + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + Mockito.verify(catalog, Mockito.never()).resetMetaCacheNames(); + } + + @Test + public void testCreateDatabaseRejectsMappedLocalNameConflict() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("sales_db"); + Mockito.when(database.getRemoteName()).thenReturn("Sales"); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createDb("sales_db", true, Collections.emptyMap())); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("sales_db", false, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("exist")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).namespaceExists(Mockito.any()); + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + } + + @Test + public void testCreateDatabaseHandlesConcurrentCreate() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new NamespaceNotFoundException("missing")) + .when(namespace).namespaceExists(Mockito.any()); + Mockito.doThrow(new NamespaceAlreadyExistsException("created concurrently")) + .when(namespace).createNamespace(Mockito.any()); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createDb("analytics", true, Collections.emptyMap())); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("analytics", false, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("exist")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.times(2)).createNamespace(Mockito.any()); + Mockito.verify(catalog).resetMetaCacheNames(); + } + + @Test + public void testCreateTableIfNotExistsSkipsExistingRemoteTable() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + Mockito.when(database.getTableNullable("events")).thenReturn(null); + CreateTableInfo createTableInfo = createTableInfo(true); + + try { + Assertions.assertTrue(new LanceMetadataOps(catalog).createTable(createTableInfo)); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()) + .createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database).resetMetaCacheNames(); + } + + @Test + public void testCreateTableHandlesConcurrentCreate() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + Mockito.doThrow(new TableAlreadyExistsException("created concurrently")) + .when(namespace).createTable(Mockito.any(), Mockito.any(byte[].class)); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + Mockito.when(database.getTableNullable("events")).thenReturn(null); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createTable(createTableInfo(true))); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createTable(createTableInfo(false))); + Assertions.assertTrue(exception.getMessage().contains("already exists")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.times(2)) + .createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database).resetMetaCacheNames(); + } + + @Test + public void testCreateTableIgnoresStaleLocalCacheWhenRemoteTableIsGone() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable staleTable = table("local_db", "events", "analytics", "stale_events"); + Mockito.when(database.getTableNullable("events")) + .thenReturn(staleTable, null); + + try { + Assertions.assertFalse(new LanceMetadataOps(catalog).createTable(createTableInfo(true))); + } finally { + client.close(); + } + + ArgumentCaptor<CreateTableRequest> request = ArgumentCaptor.forClass(CreateTableRequest.class); + Mockito.verify(namespace).createTable(request.capture(), Mockito.any(byte[].class)); + Assertions.assertEquals(Arrays.asList("tenant", "analytics", "events"), + request.getValue().getId()); + Mockito.verify(database, Mockito.atLeastOnce()).resetMetaCacheNames(); + Mockito.verify(catalog).invalidateTableAccessCache(); + } + + @Test + public void testCreateTableRejectsLocalNameConflictAfterRefresh() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable conflictingTable = table("local_db", "events", "analytics", "Events"); + Mockito.when(database.getTableNullable("events")).thenReturn(conflictingTable); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createTable(createTableInfo(true))); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createTable(createTableInfo(false))); + Assertions.assertTrue(exception.getMessage().contains("already exists")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database, Mockito.times(2)).resetMetaCacheNames(); + } + + @Test + public void testIfExistsHandlesConcurrentDropAndRefreshesLocalNames() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new NamespaceNotFoundException("dropped concurrently")) + .when(namespace).dropNamespace(Mockito.any()); + Mockito.doThrow(new TableNotFoundException("dropped concurrently")) + .when(namespace).dropTable(Mockito.any()); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("analytics"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable table = table("local_db", "local_table", "analytics", "events"); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertDoesNotThrow(() -> ops.dropDb("analytics", true, false)); + Assertions.assertThrows(DdlException.class, + () -> ops.dropDb("analytics", false, false)); + Assertions.assertDoesNotThrow(() -> ops.dropTable(table, true)); + Assertions.assertThrows(DdlException.class, + () -> ops.dropTable(table, false)); + } finally { + client.close(); + } + + Mockito.verify(catalog, Mockito.never()).unregisterDatabase(Mockito.anyString()); + Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation(); + Mockito.verify(database).unregisterTable("local_table"); Review Comment: [P1] Verify the replay-safe table eviction method used here. `ops.dropTable(table, true)` runs `afterDropTable`, which calls `database.unregisterTableForReplay("local_table")`; this test instead verifies `unregisterTable` on a plain Mockito mock. Nothing forwards that call, so this new test fails even when the implementation behaves as written. Verify the actual method and cover the fallback separately. ########## fe/fe-core/src/test/java/org/apache/doris/datasource/lance/LanceMetadataOpsTest.java: ########## @@ -0,0 +1,652 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance; + +import org.apache.doris.analysis.ColumnPosition; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.RefreshManager; +import org.apache.doris.catalog.Type; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.UserException; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo; + +import org.apache.arrow.memory.BufferAllocator; +import org.apache.arrow.memory.RootAllocator; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.lance.Session; +import org.lance.namespace.LanceNamespace; +import org.lance.namespace.errors.NamespaceAlreadyExistsException; +import org.lance.namespace.errors.NamespaceNotFoundException; +import org.lance.namespace.errors.TableAlreadyExistsException; +import org.lance.namespace.errors.TableNotFoundException; +import org.lance.namespace.model.AlterTableAddColumnsRequest; +import org.lance.namespace.model.AlterTableAlterColumnsRequest; +import org.lance.namespace.model.CreateTableRequest; +import org.lance.namespace.model.DropNamespaceRequest; +import org.mockito.ArgumentCaptor; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +import java.util.Arrays; +import java.util.Collections; +import java.util.Optional; + +public class LanceMetadataOpsTest { + @Test + public void testAddColumnValidation() throws UserException { + Column nullable = new Column("score", Type.INT, true); + Assertions.assertDoesNotThrow( + () -> LanceMetadataOps.validateAddColumn(nullable, null)); + + Column required = new Column("required", Type.INT, false); + assertRejected(() -> LanceMetadataOps.validateAddColumn(required, null), + "only supports nullable columns"); + + Column defaulted = new Column("defaulted", Type.INT, false, null, true, "1", ""); + assertRejected(() -> LanceMetadataOps.validateAddColumn(defaulted, null), + "does not support default values"); + + Column commented = new Column("commented", Type.INT, true, "comment"); + commented.setCommentSpecified(true); + assertRejected(() -> LanceMetadataOps.validateAddColumn(commented, null), + "does not support column comments"); + + assertRejected(() -> LanceMetadataOps.validateAddColumn( + nullable, new ColumnPosition("id")), + "does not support column positions"); + } + + @Test + public void testModifyColumnValidation() throws UserException { + Column column = new Column("score", Type.BIGINT, true); + column.setNullableSpecified(true); + Assertions.assertDoesNotThrow( + () -> LanceMetadataOps.validateModifyColumn(column, null)); + + column.setCommentSpecified(true); + assertRejected(() -> LanceMetadataOps.validateModifyColumn(column, null), + "does not support column comments"); + } + + @Test + public void testCreateDatabaseIfNotExistsSkipsExistingDatabase() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + + try { + Assertions.assertTrue(new LanceMetadataOps(catalog).createDb( + "analytics", true, Collections.singletonMap("owner", "doris"))); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + Mockito.verify(catalog).resetMetaCacheNames(); + } + + @Test + public void testCreateRootDatabaseIsRejected() { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + for (boolean ifNotExists : Arrays.asList(false, true)) { + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("default", ifNotExists, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("root database")); + } + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).namespaceExists(Mockito.any()); + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + Mockito.verify(catalog, Mockito.never()).resetMetaCacheNames(); + } + + @Test + public void testCreateDatabaseRejectsMappedLocalNameConflict() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("sales_db"); + Mockito.when(database.getRemoteName()).thenReturn("Sales"); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createDb("sales_db", true, Collections.emptyMap())); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("sales_db", false, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("exist")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).namespaceExists(Mockito.any()); + Mockito.verify(namespace, Mockito.never()).createNamespace(Mockito.any()); + } + + @Test + public void testCreateDatabaseHandlesConcurrentCreate() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new NamespaceNotFoundException("missing")) + .when(namespace).namespaceExists(Mockito.any()); + Mockito.doThrow(new NamespaceAlreadyExistsException("created concurrently")) + .when(namespace).createNamespace(Mockito.any()); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createDb("analytics", true, Collections.emptyMap())); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createDb("analytics", false, Collections.emptyMap())); + Assertions.assertTrue(exception.getMessage().contains("exist")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.times(2)).createNamespace(Mockito.any()); + Mockito.verify(catalog).resetMetaCacheNames(); + } + + @Test + public void testCreateTableIfNotExistsSkipsExistingRemoteTable() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + Mockito.when(database.getTableNullable("events")).thenReturn(null); + CreateTableInfo createTableInfo = createTableInfo(true); + + try { + Assertions.assertTrue(new LanceMetadataOps(catalog).createTable(createTableInfo)); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()) + .createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database).resetMetaCacheNames(); + } + + @Test + public void testCreateTableHandlesConcurrentCreate() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + Mockito.doThrow(new TableAlreadyExistsException("created concurrently")) + .when(namespace).createTable(Mockito.any(), Mockito.any(byte[].class)); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + Mockito.when(database.getTableNullable("events")).thenReturn(null); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createTable(createTableInfo(true))); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createTable(createTableInfo(false))); + Assertions.assertTrue(exception.getMessage().contains("already exists")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.times(2)) + .createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database).resetMetaCacheNames(); + } + + @Test + public void testCreateTableIgnoresStaleLocalCacheWhenRemoteTableIsGone() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable staleTable = table("local_db", "events", "analytics", "stale_events"); + Mockito.when(database.getTableNullable("events")) + .thenReturn(staleTable, null); + + try { + Assertions.assertFalse(new LanceMetadataOps(catalog).createTable(createTableInfo(true))); + } finally { + client.close(); + } + + ArgumentCaptor<CreateTableRequest> request = ArgumentCaptor.forClass(CreateTableRequest.class); + Mockito.verify(namespace).createTable(request.capture(), Mockito.any(byte[].class)); + Assertions.assertEquals(Arrays.asList("tenant", "analytics", "events"), + request.getValue().getId()); + Mockito.verify(database, Mockito.atLeastOnce()).resetMetaCacheNames(); + Mockito.verify(catalog).invalidateTableAccessCache(); + } + + @Test + public void testCreateTableRejectsLocalNameConflictAfterRefresh() throws UserException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new TableNotFoundException("missing")) + .when(namespace).tableExists(Mockito.any()); + LanceCatalogClient client = newClient(namespace, new RootAllocator(1024 * 1024)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("local_db"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable conflictingTable = table("local_db", "events", "analytics", "Events"); + Mockito.when(database.getTableNullable("events")).thenReturn(conflictingTable); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.createTable(createTableInfo(true))); + DdlException exception = Assertions.assertThrows(DdlException.class, + () -> ops.createTable(createTableInfo(false))); + Assertions.assertTrue(exception.getMessage().contains("already exists")); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).createTable(Mockito.any(), Mockito.any(byte[].class)); + Mockito.verify(database, Mockito.times(2)).resetMetaCacheNames(); + } + + @Test + public void testIfExistsHandlesConcurrentDropAndRefreshesLocalNames() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + Mockito.doThrow(new NamespaceNotFoundException("dropped concurrently")) + .when(namespace).dropNamespace(Mockito.any()); + Mockito.doThrow(new TableNotFoundException("dropped concurrently")) + .when(namespace).dropTable(Mockito.any()); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("analytics"); + Mockito.when(catalog.getDbForReplay("local_db")).thenReturn(Optional.of(database)); + Mockito.when(database.getRemoteName()).thenReturn("analytics"); + ExternalTable table = table("local_db", "local_table", "analytics", "events"); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertDoesNotThrow(() -> ops.dropDb("analytics", true, false)); + Assertions.assertThrows(DdlException.class, + () -> ops.dropDb("analytics", false, false)); + Assertions.assertDoesNotThrow(() -> ops.dropTable(table, true)); + Assertions.assertThrows(DdlException.class, + () -> ops.dropTable(table, false)); + } finally { + client.close(); + } + + Mockito.verify(catalog, Mockito.never()).unregisterDatabase(Mockito.anyString()); + Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation(); + Mockito.verify(database).unregisterTable("local_table"); + } + + @Test + public void testDropDatabaseReturnsFalseForIfExistsNoOp() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertFalse(ops.dropDb("missing_db", true, false)); + Assertions.assertThrows(DdlException.class, + () -> ops.dropDb("missing_db", false, false)); + } finally { + client.close(); + } + + Mockito.verify(namespace, Mockito.never()).namespaceExists(Mockito.any()); + Mockito.verify(namespace, Mockito.never()).dropNamespace(Mockito.any()); + Mockito.verify(catalog, Mockito.never()).unregisterDatabase(Mockito.anyString()); + Mockito.verify(catalog).retireAllDatabaseObjectsWithoutEngineInvalidation(); + } + + @Test + public void testDropDatabaseUsesRemoteDatabaseName() throws DdlException { + LanceNamespace namespace = Mockito.mock(LanceNamespace.class); + LanceCatalogClient client = newClient(namespace, Mockito.mock(BufferAllocator.class)); + LanceExternalCatalog catalog = catalogWithClient(client); + ExternalDatabase<?> database = Mockito.mock(ExternalDatabase.class); + Mockito.doReturn(database).when(catalog).getDbNullable("sales_db"); + Mockito.doReturn(Optional.of(database)).when(catalog).getDbForReplay("sales_db"); + Mockito.when(database.getFullName()).thenReturn("sales_db"); + Mockito.when(database.getRemoteName()).thenReturn("Sales"); + LanceMetadataOps ops = new LanceMetadataOps(catalog); + + try { + Assertions.assertTrue(ops.dropDb("sales_db", false, true)); Review Comment: [P1] Stub the listings required by this FORCE test. `dropDb(..., true)` now traverses children before calling `dropNamespace`, but this bare `LanceNamespace` mock has no `listNamespaces` response. Mockito returns null and `listChildNamespaces` throws before this `assertTrue` or the remote-name assertion can pass. Return empty namespace and table listing responses, then verify the captured ID. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -139,6 +170,82 @@ List<String> listDatabaseNames() { return new ArrayList<>(databases); } + boolean isRootDatabase(String dbName) { + return rootDatabase.equals(dbName); + } + + boolean databaseExists(String dbName) { + if (isRootDatabase(dbName)) { + return true; + } + try { + NamespaceExistsRequest request = new NamespaceExistsRequest().id(buildNamespaceId(dbName)); + synchronized (namespaceLock) { + namespace.namespaceExists(request); + } + return true; + } catch (NamespaceNotFoundException e) { + return false; + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void createDatabase(String dbName, Map<String, String> properties) { + try { + CreateNamespaceRequest request = new CreateNamespaceRequest() + .id(buildNamespaceId(dbName)) + .mode("Create") + .properties(properties == null ? Collections.emptyMap() : properties); + synchronized (namespaceLock) { + namespace.createNamespace(request); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void dropDatabase(String dbName, boolean ifExists, boolean force) { + try { + List<String> namespaceId = buildNamespaceId(dbName); + synchronized (namespaceLock) { + if (force) { + try { + dropNamespaceCascade(namespaceId, ifExists ? "Skip" : "Fail"); + } catch (NamespaceNotFoundException e) { + if (!ifExists) { + throw e; + } + } + return; + } + namespace.dropNamespace(new DropNamespaceRequest() + .id(namespaceId) + .mode(ifExists ? "Skip" : "Fail") + .behavior("Restrict")); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + private void dropNamespaceCascade(List<String> namespaceId, String mode) { + for (String child : listChildNamespaces(namespaceId)) { + List<String> childId = new ArrayList<>(namespaceId); + childId.add(child); + dropNamespaceCascade(childId, "Fail"); + } + for (String table : listTableNames(namespaceId)) { Review Comment: [P2] Bound or batch this FORCE deletion path for large namespaces. Every Lance 12 `drop_table` rewrites the remaining manifest, so deleting N tables here performs roughly N(N-1)/2 manifest-row visits and N object-store commits; the outer `namespaceLock` blocks this catalog's listings and existence checks throughout. A 1,000-table database means about 500,000 row visits before counting other manifest objects. Use a batch SDK path or reject large FORCE operations until one is available. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -139,6 +170,82 @@ List<String> listDatabaseNames() { return new ArrayList<>(databases); } + boolean isRootDatabase(String dbName) { + return rootDatabase.equals(dbName); + } + + boolean databaseExists(String dbName) { + if (isRootDatabase(dbName)) { + return true; + } + try { + NamespaceExistsRequest request = new NamespaceExistsRequest().id(buildNamespaceId(dbName)); + synchronized (namespaceLock) { + namespace.namespaceExists(request); + } + return true; + } catch (NamespaceNotFoundException e) { + return false; + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void createDatabase(String dbName, Map<String, String> properties) { + try { + CreateNamespaceRequest request = new CreateNamespaceRequest() + .id(buildNamespaceId(dbName)) + .mode("Create") + .properties(properties == null ? Collections.emptyMap() : properties); + synchronized (namespaceLock) { + namespace.createNamespace(request); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void dropDatabase(String dbName, boolean ifExists, boolean force) { + try { + List<String> namespaceId = buildNamespaceId(dbName); + synchronized (namespaceLock) { + if (force) { + try { + dropNamespaceCascade(namespaceId, ifExists ? "Skip" : "Fail"); + } catch (NamespaceNotFoundException e) { + if (!ifExists) { + throw e; + } + } + return; + } + namespace.dropNamespace(new DropNamespaceRequest() + .id(namespaceId) + .mode(ifExists ? "Skip" : "Fail") + .behavior("Restrict")); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + private void dropNamespaceCascade(List<String> namespaceId, String mode) { + for (String child : listChildNamespaces(namespaceId)) { + List<String> childId = new ArrayList<>(namespaceId); + childId.add(child); + dropNamespaceCascade(childId, "Fail"); + } + for (String table : listTableNames(namespaceId)) { + List<String> tableId = new ArrayList<>(namespaceId); + tableId.add(table); + namespace.dropTable(new DropTableRequest().id(tableId)); Review Comment: [P1] Protect table locations referenced outside the FORCE deletion set. Lance 12 allows `archive.events_alias` to register the same relative `events.lance` location as root table `events`. This loop drops the alias, and the SDK deletes `events.lance` physically while the root table's manifest entry remains. `DROP DATABASE archive FORCE` therefore destroys an unrelated visible table. Preflight shared locations or use an SDK operation that deregisters shared aliases without deleting their data. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -139,6 +170,82 @@ List<String> listDatabaseNames() { return new ArrayList<>(databases); } + boolean isRootDatabase(String dbName) { + return rootDatabase.equals(dbName); + } + + boolean databaseExists(String dbName) { + if (isRootDatabase(dbName)) { + return true; + } + try { + NamespaceExistsRequest request = new NamespaceExistsRequest().id(buildNamespaceId(dbName)); + synchronized (namespaceLock) { + namespace.namespaceExists(request); + } + return true; + } catch (NamespaceNotFoundException e) { + return false; + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void createDatabase(String dbName, Map<String, String> properties) { + try { + CreateNamespaceRequest request = new CreateNamespaceRequest() + .id(buildNamespaceId(dbName)) + .mode("Create") + .properties(properties == null ? Collections.emptyMap() : properties); + synchronized (namespaceLock) { + namespace.createNamespace(request); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + void dropDatabase(String dbName, boolean ifExists, boolean force) { + try { + List<String> namespaceId = buildNamespaceId(dbName); + synchronized (namespaceLock) { + if (force) { + try { + dropNamespaceCascade(namespaceId, ifExists ? "Skip" : "Fail"); + } catch (NamespaceNotFoundException e) { + if (!ifExists) { + throw e; + } + } + return; + } + namespace.dropNamespace(new DropNamespaceRequest() + .id(namespaceId) + .mode(ifExists ? "Skip" : "Fail") + .behavior("Restrict")); + } + } catch (DdlException e) { + throw new RuntimeException(e); + } + } + + private void dropNamespaceCascade(List<String> namespaceId, String mode) { + for (String child : listChildNamespaces(namespaceId)) { + List<String> childId = new ArrayList<>(namespaceId); + childId.add(child); + dropNamespaceCascade(childId, "Fail"); + } + for (String table : listTableNames(namespaceId)) { + List<String> tableId = new ArrayList<>(namespaceId); + tableId.add(table); + namespace.dropTable(new DropTableRequest().id(tableId)); Review Comment: [P1] Keep FORCE retries able to clean up a table whose physical delete failed. Lance 12 `drop_table` removes its manifest row before deleting the directory. If that directory deletion fails, this loop throws; on retry, child `list_tables` no longer reports the table, so FORCE can finish while its data remains orphaned. Preserve the location for retry or require an SDK drop operation with recoverable cleanup, and test a failure after manifest removal. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -213,6 +323,96 @@ 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; + } + try { + // DirectoryNamespace only accepts locations relative to its warehouse root. + namespace.registerTable(new RegisterTableRequest() + .id(tableId) + .location(tableName + ".lance") + .mode("Create")); + } catch (TableAlreadyExistsException e) { Review Comment: [P1] Verify the winning registration before retrying this mutation. After Doris observes an unregistered root `events.lance`, another namespace client can register the name `events` at another valid relative location such as `other.lance`. This catch ignores the conflict and retries by name, so `DROP TABLE default.events` can delete `other.lance` (or ALTER can change it) instead of the dataset Doris observed. Describe and compare the registered location to `events.lance`, or fail without mutating it. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceMetadataOps.java: ########## @@ -0,0 +1,581 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.datasource.lance; + +import org.apache.doris.analysis.ColumnPosition; +import org.apache.doris.catalog.Column; +import org.apache.doris.catalog.Env; +import org.apache.doris.common.DdlException; +import org.apache.doris.common.ErrorCode; +import org.apache.doris.common.ErrorReport; +import org.apache.doris.common.UserException; +import org.apache.doris.datasource.ExternalDatabase; +import org.apache.doris.datasource.ExternalTable; +import org.apache.doris.datasource.lance.metadata.LanceTypeConverter; +import org.apache.doris.datasource.operations.ExternalMetadataOps; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateOrReplaceTagInfo; +import org.apache.doris.nereids.trees.plans.commands.info.CreateTableInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropBranchInfo; +import org.apache.doris.nereids.trees.plans.commands.info.DropTagInfo; + +import org.apache.arrow.vector.types.pojo.Schema; +import org.apache.commons.lang3.StringUtils; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.lance.namespace.errors.NamespaceAlreadyExistsException; +import org.lance.namespace.errors.NamespaceNotFoundException; +import org.lance.namespace.errors.TableAlreadyExistsException; +import org.lance.namespace.errors.TableNotFoundException; +import org.lance.namespace.model.AddColumnsEntry; +import org.lance.namespace.model.AlterColumnsEntry; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.TreeSet; + +/** Doris external metadata operations backed by the Lance Namespace API. */ +public class LanceMetadataOps implements ExternalMetadataOps { + private static final Logger LOG = LogManager.getLogger(LanceMetadataOps.class); + private static final String TABLE_COMMENT_PROPERTY = "comment"; + + private final LanceExternalCatalog catalog; + + public LanceMetadataOps(LanceExternalCatalog catalog) { + this.catalog = catalog; + } + + @Override + public boolean createDbImpl(String dbName, boolean ifNotExists, Map<String, String> properties) + throws DdlException { + return execute("Failed to create Lance database " + dbName, client -> { + if (client.isRootDatabase(dbName)) { + throw new DdlException("Cannot create the configured Lance root database: " + dbName); + } + if (catalog.getDbNullable(dbName) != null) { + if (ifNotExists) { + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + } + if (client.databaseExists(dbName)) { + if (ifNotExists) { + catalog.resetMetaCacheNames(); + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + } + try { + client.createDatabase(dbName, new HashMap<>( + Optional.ofNullable(properties).orElse(Collections.emptyMap()))); + return false; + } catch (NamespaceAlreadyExistsException e) { + if (ifNotExists) { + catalog.resetMetaCacheNames(); + return true; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_CREATE_EXISTS, dbName); + throw new IllegalStateException("unreachable"); + } + }); + } + + @Override + public void afterCreateDb() { + catalog.resetMetaCacheNames(); + } + + @Override + public boolean dropDbImpl(String dbName, boolean ifExists, boolean force) throws DdlException { + ExternalDatabase<?> db = catalog.getDbNullable(dbName); + return execute("Failed to drop Lance database " + dbName, client -> { + if (client.isRootDatabase(dbName)) { + throw new DdlException("Cannot drop the configured Lance root database: " + dbName); + } + if (db == null) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + return false; + } + String remoteDbName = db.getRemoteName(); + if (client.isRootDatabase(remoteDbName)) { + throw new DdlException("Cannot drop the configured Lance root database: " + dbName); + } + if (!client.databaseExists(remoteDbName)) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + } + try { + client.dropDatabase(remoteDbName, ifExists, force); + } catch (NamespaceNotFoundException e) { + if (ifExists) { + return false; + } + ErrorReport.reportDdlException(ErrorCode.ERR_DB_DROP_EXISTS, dbName); + } + return true; + }); + } + + @Override + public void afterDropDb(String dbName) { + Optional<ExternalDatabase<? extends ExternalTable>> db = catalog.getDbForReplay(dbName); + if (db.isPresent()) { + catalog.unregisterDatabase(db.get().getFullName()); + return; + } + catalog.unregisterDatabase(dbName); + catalog.retireAllDatabaseObjectsWithoutEngineInvalidation(); + } + + @Override + public void afterDropDbNoOp(String dbName) { Review Comment: [P1] Complete cache reconciliation for this IF EXISTS no-op. If a cached namespace was removed outside Doris, `DROP DATABASE IF EXISTS` retires database objects on the leader here, but `ExternalCatalog.dropDb` skips the edit log for a false result, leaving follower FEs with stale objects. This hook also bypasses Lance's table-access URI cache invalidation, so a later external recreation under the same name can still read the old URI until cache expiry. Propagate a targeted refresh to followers and invalidate the engine access cache when retiring these objects. ########## 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 = "REAL"; + break; + case DOUBLE: + sqlType = "DOUBLE"; + break; + case CHAR: + case VARCHAR: + case STRING: + sqlType = "VARCHAR"; + break; + case VARBINARY: Review Comment: [P1] Preserve the declared VARBINARY bound on ADD. `ALTER TABLE t ADD COLUMN b VARBINARY(4) NULL` passes the ALTER and Lance checks, but this branch emits `CAST(NULL AS BINARY)`. Lance 12 adds an unbounded Arrow Binary field, and the next Doris metadata load reports maximum-width VARBINARY instead of length 4. Reject bounded VARBINARY ADD until a schema path can retain its bound, or persist and round-trip that bound; add a namespace-backed ALTER test. ########## fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceNamespaceClient.java: ########## @@ -213,6 +323,96 @@ 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); Review Comment: [P2] Clean up a child dataset when its manifest registration fails. In Lance 12 this `createTable` call writes the dataset to a fresh hashed directory before inserting its manifest row. If the manifest commit fails after the write, CREATE reports an error but leaves an invisible physical table; a retry chooses another directory, and later DROP/FORCE cannot find the first. Require an SDK path that removes the staged directory on failure or retain a recoverable cleanup record, and cover a manifest-insert failure after dataset creation. ########## 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 = "REAL"; + break; + case DOUBLE: + sqlType = "DOUBLE"; + break; + case CHAR: + case VARCHAR: + case STRING: Review Comment: [P1] Emit type names the pinned Lance ADD parser accepts. This branch produces `CAST(NULL AS VARCHAR)` for `STRING`/`VARCHAR`/`CHAR`, and the FLOAT branch produces `CAST(NULL AS REAL)`. Lance 12's SQL expression planner handles `STRING` and `FLOAT` but rejects `VARCHAR` and `REAL`, so `ADD COLUMN note STRING` (including the new regression's multi-column ALTER) fails remotely. Emit the supported tokens and verify both through a real namespace ALTER. ########## 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 = "REAL"; + break; + case DOUBLE: + sqlType = "DOUBLE"; + break; + case CHAR: + case VARCHAR: + case STRING: + sqlType = "VARCHAR"; + 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: Review Comment: [P1] Preserve or reject bounded character types across Lance DDL. `CREATE TABLE t (c VARCHAR(10))` passes validation but stores Arrow Utf8 without length metadata; the next load reports `c STRING`, so DESC/SHOW CREATE lose the bound. `MODIFY COLUMN c VARCHAR(10)` likewise sends only `utf8` and cannot establish a bound on an existing STRING field. Store and round-trip supported CHAR/VARCHAR lengths, or reject bounded declarations in both CREATE and MODIFY before reporting success. -- 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]
