This is an automated email from the ASF dual-hosted git repository. diqiu50 pushed a commit to branch trino-irc-1.3 in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit 04eb5171467fad4529d851bc5746ee98d863696b Author: diqiu50 <[email protected]> AuthorDate: Fri Aug 21 17:44:55 2026 +0800 [#12546] improvement(trino-connector): Scope system tables to their metalake and retry failed init Build the system table factory per entry catalog with its metalake, the way the stored procedures already were, so two entry catalogs no longer report each other's state. Publish the catalog connector manager only once its initialization succeeded, so a failed init is retried instead of leaving an entry catalog whose load loop never started. --- .../trino/connector/GravitinoConnectorFactory.java | 14 +++-- .../system/table/GravitinoSystemTableCatalog.java | 12 +++- .../table/GravitinoSystemTableCatalogStatus.java | 13 +++- .../system/table/GravitinoSystemTableFactory.java | 13 ++-- .../table/GravitinoSystemTableLoadStatus.java | 12 +++- .../TestGravitinoConnectorFactoryInit.java | 70 ++++++++++++++++++++++ .../table/TestGravitinoSystemStatusTables.java | 10 ++-- 7 files changed, 122 insertions(+), 22 deletions(-) diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java index 59b58e131b..ac9302e343 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory.java @@ -62,8 +62,6 @@ public class GravitinoConnectorFactory implements ConnectorFactory { public static final String DEFAULT_CONNECTOR_NAME = "gravitino"; @SuppressWarnings("UnusedVariable") - private GravitinoSystemTableFactory gravitinoSystemTableFactory; - private CatalogConnectorManager catalogConnectorManager; private boolean catalogConnectorManagerStartTriggered = false; @@ -127,12 +125,11 @@ public class GravitinoConnectorFactory implements ConnectorFactory { CatalogRegister catalogRegister = new CatalogRegister(); CatalogConnectorFactory catalogConnectorFactory = createCatalogConnectorFactory(config); - catalogConnectorManager = + CatalogConnectorManager manager = new CatalogConnectorManager( catalogRegister, catalogConnectorFactory, this::getTrinoCatalogName); - catalogConnectorManager.config(config, client); + manager.config(config, client); - gravitinoSystemTableFactory = new GravitinoSystemTableFactory(catalogConnectorManager); if (isCoordinator(trinoConnectorContext)) { // Pin the system table splits here: the registration state the system tables report // is only recorded on the coordinator by the load loop started below. Starting the @@ -140,6 +137,7 @@ public class GravitinoConnectorFactory implements ConnectorFactory { GravitinoSystemConnector.Split.setCoordinatorAddress( getCurrentNodeAddress(trinoConnectorContext)); } + catalogConnectorManager = manager; } // The `trino.jdbc.*` settings that CatalogRegister needs to connect back to the @@ -184,9 +182,13 @@ public class GravitinoConnectorFactory implements ConnectorFactory { throw new TrinoException( GravitinoErrorCode.GRAVITINO_METALAKE_NOT_EXISTS, "No gravitino metalake selected"); } + // Built per entry catalog, like the stored procedures: both are scoped to this catalog's + // metalake even though the underlying manager is shared. GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory = new GravitinoStoredProcedureFactory(catalogConnectorManager, metalake); - return createSystemConnector(gravitinoStoredProcedureFactory, gravitinoSystemTableFactory); + GravitinoSystemTableFactory systemTableFactory = + new GravitinoSystemTableFactory(catalogConnectorManager, metalake); + return createSystemConnector(gravitinoStoredProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalog.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalog.java index aac74a4762..4105e5279a 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalog.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalog.java @@ -54,14 +54,18 @@ public class GravitinoSystemTableCatalog extends GravitinoSystemTable { ColumnMetadata.builder().setName("properties").setType(VARCHAR).build())); private final CatalogConnectorManager catalogConnectorManager; + private final String metalake; /** * Constructs a new GravitinoSystemTableCatalog. * * @param catalogConnectorManager the manager for catalog connectors + * @param metalake the metalake to report on */ - public GravitinoSystemTableCatalog(CatalogConnectorManager catalogConnectorManager) { + public GravitinoSystemTableCatalog( + CatalogConnectorManager catalogConnectorManager, String metalake) { this.catalogConnectorManager = catalogConnectorManager; + this.metalake = metalake; } @Override @@ -69,8 +73,10 @@ public class GravitinoSystemTableCatalog extends GravitinoSystemTable { List<GravitinoCatalog> gravitinoCatalogs = new ArrayList<>(); // retrieve catalogs form the Gravitino server with the configuration metalakes, // the catalogConnectorManager does not manager catalogs in worker nodes - catalogConnectorManager - .getUsedMetalakes() + // Only the metalake this connector is configured with: the manager is shared by every entry + // catalog in this Trino. + catalogConnectorManager.getUsedMetalakes().stream() + .filter(metalake::equals) .forEach( (metalakeName) -> { GravitinoMetalake metalake = catalogConnectorManager.getMetalake(metalakeName); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalogStatus.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalogStatus.java index 2d2722f506..bd75caa5ca 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalogStatus.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableCatalogStatus.java @@ -57,21 +57,30 @@ public class GravitinoSystemTableCatalogStatus extends GravitinoSystemTable { ColumnMetadata.builder().setName("failure_count").setType(BIGINT).build())); private final CatalogConnectorManager catalogConnectorManager; + private final String metalake; /** * Constructs a new GravitinoSystemTableCatalogStatus. * * @param catalogConnectorManager the manager for catalog connectors + * @param metalake the metalake to report on */ - public GravitinoSystemTableCatalogStatus(CatalogConnectorManager catalogConnectorManager) { + public GravitinoSystemTableCatalogStatus( + CatalogConnectorManager catalogConnectorManager, String metalake) { this.catalogConnectorManager = catalogConnectorManager; + this.metalake = metalake; } @Override public Page loadPageData() { // Take a snapshot first, the load loop writes these states concurrently and the column // builders must all end up with the same number of positions. - List<CatalogRegistrationState> states = catalogConnectorManager.getCatalogRegistrationStates(); + // The load loop is shared by every entry catalog in this Trino, so report only the metalake + // this connector is configured with. + List<CatalogRegistrationState> states = + catalogConnectorManager.getCatalogRegistrationStates().stream() + .filter(state -> state.getMetalake().equals(metalake)) + .toList(); int size = states.size(); BlockBuilder metalakeColumnBuilder = VARCHAR.createBlockBuilder(null, size); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableFactory.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableFactory.java index adc82731a1..4b41c8d623 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableFactory.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableFactory.java @@ -31,6 +31,7 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; public class GravitinoSystemTableFactory { private final CatalogConnectorManager catalogConnectorManager; + private final String metalake; // Per instance, not static: the tables are bound to one CatalogConnectorManager, and only the // manager on the coordinator runs the load loop that fills in the registration state. A shared @@ -41,9 +42,13 @@ public class GravitinoSystemTableFactory { * Constructs a new GravitinoSystemTableFactory. * * @param catalogConnectorManager the manager for catalog connectors + * @param metalake the metalake this connector is configured with; the tables only report on it, + * so that two entry catalogs pointed at different metalakes do not report each other's state */ - public GravitinoSystemTableFactory(CatalogConnectorManager catalogConnectorManager) { + public GravitinoSystemTableFactory( + CatalogConnectorManager catalogConnectorManager, String metalake) { this.catalogConnectorManager = catalogConnectorManager; + this.metalake = metalake; registerSystemTables(); } @@ -52,13 +57,13 @@ public class GravitinoSystemTableFactory { private void registerSystemTables() { systemTables.put( GravitinoSystemTableCatalog.TABLE_NAME, - new GravitinoSystemTableCatalog(catalogConnectorManager)); + new GravitinoSystemTableCatalog(catalogConnectorManager, metalake)); systemTables.put( GravitinoSystemTableCatalogStatus.TABLE_NAME, - new GravitinoSystemTableCatalogStatus(catalogConnectorManager)); + new GravitinoSystemTableCatalogStatus(catalogConnectorManager, metalake)); systemTables.put( GravitinoSystemTableLoadStatus.TABLE_NAME, - new GravitinoSystemTableLoadStatus(catalogConnectorManager)); + new GravitinoSystemTableLoadStatus(catalogConnectorManager, metalake)); } /** diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableLoadStatus.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableLoadStatus.java index dc4cbfe882..de8d511bd6 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableLoadStatus.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/table/GravitinoSystemTableLoadStatus.java @@ -61,14 +61,18 @@ public class GravitinoSystemTableLoadStatus extends GravitinoSystemTable { ColumnMetadata.builder().setName("metalake_errors").setType(VARCHAR).build())); private final CatalogConnectorManager catalogConnectorManager; + private final String metalake; /** * Constructs a new GravitinoSystemTableLoadStatus. * * @param catalogConnectorManager the manager for catalog connectors + * @param metalake the metalake to report errors for */ - public GravitinoSystemTableLoadStatus(CatalogConnectorManager catalogConnectorManager) { + public GravitinoSystemTableLoadStatus( + CatalogConnectorManager catalogConnectorManager, String metalake) { this.catalogConnectorManager = catalogConnectorManager; + this.metalake = metalake; } @Override @@ -87,7 +91,11 @@ public class GravitinoSystemTableLoadStatus extends GravitinoSystemTable { consecutiveFailuresColumnBuilder, catalogConnectorManager.getConsecutiveLoadFailures()); writeNullableString(lastErrorColumnBuilder, catalogConnectorManager.getLastLoadError()); - Map<String, String> metalakeErrors = catalogConnectorManager.getMetalakeErrors(); + // The load loop itself is shared by every entry catalog, so the columns above are global. + // Only the per metalake errors are narrowed to the metalake this connector reports on. + Map<String, String> allErrors = catalogConnectorManager.getMetalakeErrors(); + Map<String, String> metalakeErrors = + allErrors.containsKey(metalake) ? Map.of(metalake, allErrors.get(metalake)) : Map.of(); if (metalakeErrors.isEmpty()) { metalakeErrorsColumnBuilder.appendNull(); } else { diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorFactoryInit.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorFactoryInit.java new file mode 100644 index 0000000000..ff00e1c93f --- /dev/null +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorFactoryInit.java @@ -0,0 +1,70 @@ +/* + * 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.gravitino.trino.connector; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import com.google.common.collect.ImmutableMap; +import io.trino.spi.HostAddress; +import io.trino.spi.Node; +import io.trino.spi.NodeManager; +import io.trino.spi.TrinoException; +import io.trino.spi.connector.ConnectorContext; +import java.util.Map; +import org.apache.gravitino.client.GravitinoAdminClient; +import org.junit.jupiter.api.Test; + +public class TestGravitinoConnectorFactoryInit { + + @Test + public void testFailedInitializationIsRetriedOnTheNextCreate() { + GravitinoConnectorFactory factory = + new GravitinoConnectorFactory(mock(GravitinoAdminClient.class)); + Map<String, String> config = + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", + "gravitino.metalake", "test", + "gravitino.trino.skip-version-validation", "true", + // Makes CatalogRegister.init() fail, so the whole initialization fails. + "catalog.config-dir", "/nonexistent-gravitino-catalog-dir", + "discovery.uri", "http://127.0.0.1:8080"); + + assertThrows(TrinoException.class, () -> factory.create("gravitino", config, mockContext())); + + // The second attempt must fail the same way. If the manager was published before the + // initialization completed, this call skips the whole block and hands back an entry catalog + // whose load loop never started, which looks healthy and registers nothing. + assertThrows(TrinoException.class, () -> factory.create("gravitino", config, mockContext())); + } + + @SuppressWarnings("deprecation") + private static ConnectorContext mockContext() { + ConnectorContext context = mock(ConnectorContext.class); + when(context.getSpiVersion()).thenReturn("435"); + Node node = mock(Node.class); + when(node.isCoordinator()).thenReturn(true); + when(node.getHostAndPort()).thenReturn(HostAddress.fromParts("127.0.0.1", 8080)); + NodeManager nodeManager = mock(NodeManager.class); + when(nodeManager.getCurrentNode()).thenReturn(node); + when(context.getNodeManager()).thenReturn(nodeManager); + return context; + } +} diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/system/table/TestGravitinoSystemStatusTables.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/system/table/TestGravitinoSystemStatusTables.java index ae5ab372cd..6973d8f4d7 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/system/table/TestGravitinoSystemStatusTables.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/system/table/TestGravitinoSystemStatusTables.java @@ -92,7 +92,7 @@ public class TestGravitinoSystemStatusTables { when(manager.getLastLoadError()).thenReturn(null); when(manager.getMetalakeErrors()).thenReturn(Map.of()); - Page page = new GravitinoSystemTableLoadStatus(manager).loadPageData(); + Page page = new GravitinoSystemTableLoadStatus(manager, "test").loadPageData(); assertEquals(1, page.getPositionCount()); assertEquals(6, page.getChannelCount()); @@ -113,7 +113,7 @@ public class TestGravitinoSystemStatusTables { when(manager.getLastLoadError()).thenReturn("Connection refused"); when(manager.getMetalakeErrors()).thenReturn(Map.of("test", "Connection refused")); - Page page = new GravitinoSystemTableLoadStatus(manager).loadPageData(); + Page page = new GravitinoSystemTableLoadStatus(manager, "test").loadPageData(); // last_success_time stays null while the server is unreachable. assertTrue(page.getBlock(2).isNull(0)); @@ -136,8 +136,8 @@ public class TestGravitinoSystemStatusTables { CatalogConnectorManager second = mock(CatalogConnectorManager.class); when(second.getCatalogRegistrationStates()).thenReturn(List.of()); - GravitinoSystemTableFactory firstFactory = new GravitinoSystemTableFactory(first); - GravitinoSystemTableFactory secondFactory = new GravitinoSystemTableFactory(second); + GravitinoSystemTableFactory firstFactory = new GravitinoSystemTableFactory(first, "prod"); + GravitinoSystemTableFactory secondFactory = new GravitinoSystemTableFactory(second, "prod"); assertEquals( 1, @@ -154,7 +154,7 @@ public class TestGravitinoSystemStatusTables { private static Page loadCatalogStatusPage(List<CatalogRegistrationState> states) { CatalogConnectorManager manager = mock(CatalogConnectorManager.class); when(manager.getCatalogRegistrationStates()).thenReturn(states); - return new GravitinoSystemTableCatalogStatus(manager).loadPageData(); + return new GravitinoSystemTableCatalogStatus(manager, "test").loadPageData(); } private static String varchar(Page page, int channel) {
