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 5ccc0434527ad95bcf39149442245db7d78f024c Author: diqiu50 <[email protected]> AuthorDate: Fri Aug 21 21:12:04 2026 +0800 [#12546] improvement(trino-connector): Prune state of deleted metalakes and harden the tables Drop catalog rows and errors belonging to a metalake that no longer exists, so the tables stop reporting catalogs that are gone. Degrade instead of failing the whole load_status query when the metalake errors cannot be serialized, and hoist the shared ObjectMapper. --- .../trino/connector/GravitinoConnectorFactory.java | 8 ++-- .../connector/catalog/CatalogConnectorManager.java | 9 +++++ .../trino/connector/catalog/CatalogRegister.java | 2 +- .../system/table/GravitinoSystemTableCatalog.java | 4 +- .../table/GravitinoSystemTableLoadStatus.java | 15 +++++--- .../trino/connector/TestGravitinoConnector.java | 13 +++++++ .../catalog/TestCatalogConnectorManager.java | 45 ++++++++++++++++------ 7 files changed, 74 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 ac9302e343..9ce0efc330 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 @@ -63,6 +63,7 @@ public class GravitinoConnectorFactory implements ConnectorFactory { @SuppressWarnings("UnusedVariable") private CatalogConnectorManager catalogConnectorManager; + private boolean catalogConnectorManagerStartTriggered = false; private GravitinoAdminClient client; @@ -149,9 +150,9 @@ public class GravitinoConnectorFactory implements ConnectorFactory { if (!catalogConnectorManagerStartTriggered && !config.isDynamicConnector() && isCoordinator(trinoConnectorContext)) { - // Triggered before start() on purpose: everything that makes it fail is a - // configuration error, and retrying on the next create() would only open another - // connection. + // Mark the attempt before start() so concurrent connector creation cannot start the + // manager twice. The flag is reset below if initialization fails, allowing a corrected + // configuration to retry on the next create(). catalogConnectorManagerStartTriggered = true; // Only the configuration is re-applied here: rebuilding the Gravitino client would leak // the one a dynamic connector may have already built. @@ -159,6 +160,7 @@ public class GravitinoConnectorFactory implements ConnectorFactory { catalogConnectorManager.start(); } } catch (Exception e) { + catalogConnectorManagerStartTriggered = false; String message = "Initialization of the GravitinoConnector failed " + e.getMessage(); LOG.error(message); throw new TrinoException(GRAVITINO_RUNTIME_ERROR, message, e); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java index d07ede3d01..a3e2e703e8 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java @@ -241,6 +241,8 @@ public class CatalogConnectorManager { } } + pruneMissingMetalakes(usedMetalakes); + if (metalakeErrors.isEmpty()) { lastSuccessfulLoadTimeMs = System.currentTimeMillis(); recordLoadSuccess(); @@ -262,6 +264,13 @@ public class CatalogConnectorManager { } } + private void pruneMissingMetalakes(Set<String> usedMetalakes) { + // A metalake that was deleted, or that dropped out of the configuration, leaves its catalog + // rows behind. Without this they keep reporting REGISTERED for catalogs that no longer exist. + catalogStates.values().removeIf(state -> !usedMetalakes.contains(state.getMetalake())); + metalakeErrors.keySet().removeIf(metalakeName -> !usedMetalakes.contains(metalakeName)); + } + private void recordLoadSuccess() { if (lastLoadError != null) { LOG.info("The Gravitino catalog load loop recovered."); diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java index 4c25ef36df..5f5ce2ec02 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java @@ -37,8 +37,8 @@ import java.util.Properties; import java.util.Set; import java.util.regex.Matcher; import java.util.regex.Pattern; -import org.apache.commons.lang3.StringUtils; import javax.annotation.Nullable; +import org.apache.commons.lang3.StringUtils; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.catalog.iceberg.IcebergConnectorAdapter; 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 4105e5279a..895ceac2ca 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 @@ -45,6 +45,8 @@ public class GravitinoSystemTableCatalog extends GravitinoSystemTable { public static final SchemaTableName TABLE_NAME = new SchemaTableName(SYSTEM_TABLE_SCHEMA_NAME, "catalog"); + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + private static final ConnectorTableMetadata TABLE_METADATA = new ConnectorTableMetadata( TABLE_NAME, @@ -104,7 +106,7 @@ public class GravitinoSystemTableCatalog extends GravitinoSystemTable { try { VARCHAR.writeString( propertyColumnBuilder, - new ObjectMapper().writeValueAsString(new TreeMap<>(catalog.getProperties()))); + OBJECT_MAPPER.writeValueAsString(new TreeMap<>(catalog.getProperties()))); } catch (JsonProcessingException e) { throw new TrinoException( GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, "Invalid property format", e); // 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 de8d511bd6..058b2f965c 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 @@ -25,7 +25,6 @@ import static io.trino.spi.type.VarcharType.VARCHAR; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import io.trino.spi.Page; -import io.trino.spi.TrinoException; import io.trino.spi.block.BlockBuilder; import io.trino.spi.connector.ColumnMetadata; import io.trino.spi.connector.ConnectorTableMetadata; @@ -33,8 +32,9 @@ import io.trino.spi.connector.SchemaTableName; import java.util.List; import java.util.Map; import java.util.TreeMap; -import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * An implementation of the load status system table. @@ -45,6 +45,9 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; */ public class GravitinoSystemTableLoadStatus extends GravitinoSystemTable { + private static final Logger LOG = LoggerFactory.getLogger(GravitinoSystemTableLoadStatus.class); + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + /** The name of the load status system table. */ public static final SchemaTableName TABLE_NAME = new SchemaTableName(SYSTEM_TABLE_SCHEMA_NAME, "load_status"); @@ -102,10 +105,12 @@ public class GravitinoSystemTableLoadStatus extends GravitinoSystemTable { try { VARCHAR.writeString( metalakeErrorsColumnBuilder, - new ObjectMapper().writeValueAsString(new TreeMap<>(metalakeErrors))); + OBJECT_MAPPER.writeValueAsString(new TreeMap<>(metalakeErrors))); } catch (JsonProcessingException e) { - throw new TrinoException( - GravitinoErrorCode.GRAVITINO_ILLEGAL_ARGUMENT, "Invalid metalake error format", e); + // Degrade rather than fail: this table is what a user reads while diagnosing a broken + // load loop, and one unserializable column must not take last_error down with it. + LOG.warn("Failed to serialize the metalake errors", e); + VARCHAR.writeString(metalakeErrorsColumnBuilder, new TreeMap<>(metalakeErrors).toString()); } } diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnector.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnector.java index 2fbd12c4c1..94a6b412b6 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnector.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnector.java @@ -326,6 +326,19 @@ public abstract class TestGravitinoConnector extends AbstractGravitinoConnectorT assertEquals(row.getField(6), 0L); } + @Test + public void testCatalogStatusSystemTableWithReorderedColumns() throws Exception { + // page.getColumns() honors the requested order; a projection that sorted or de-duplicated + // channels would return correct looking data in the wrong columns. + MaterializedResult result = + computeActual("select status, catalog_name, status from gravitino.system.catalog_status"); + assertEquals(result.getRowCount(), 1); + MaterializedRow row = result.getMaterializedRows().get(0); + assertEquals(row.getField(0), "REGISTERED"); + assertEquals(row.getField(1), "memory"); + assertEquals(row.getField(2), "REGISTERED"); + } + @Test public void testLoadStatusSystemTable() throws Exception { MaterializedResult result = diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java index 0b2a97d971..9360e6fdec 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java @@ -37,15 +37,16 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import io.trino.spi.TrinoException; import io.trino.spi.connector.ConnectorContext; -import java.util.Optional; -import org.apache.gravitino.client.GravitinoAdminClient; -import org.apache.gravitino.client.GravitinoMetalake; -import org.apache.gravitino.exceptions.RESTException; import java.time.Instant; +import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import org.apache.gravitino.Audit; import org.apache.gravitino.Catalog; +import org.apache.gravitino.client.GravitinoAdminClient; +import org.apache.gravitino.client.GravitinoMetalake; +import org.apache.gravitino.exceptions.RESTException; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; @@ -624,6 +625,25 @@ public class TestCatalogConnectorManager { assertTrue(state.getLastError().contains("IllegalStateException")); } + @Test + public void testStateOfAVanishedMetalakeIsPruned() throws Exception { + LoadFixture fixture = new LoadFixture(); + fixture.withCatalogs(mockCatalog("memory", "memory", Catalog.Type.RELATIONAL)); + + // Multi metalake mode, so the used metalakes come from listMetalakes(). + CatalogConnectorManager manager = + fixture.createManager(ImmutableMap.of("gravitino.use-single-metalake", "false")); + manager.loadMetalakeSync(); + assertEquals(1, manager.getCatalogRegistrationStates().size()); + + // The metalake is deleted in Gravitino: its catalog rows must not linger as REGISTERED. + Mockito.doReturn(new GravitinoMetalake[0]).when(fixture.client).listMetalakes(); + manager.loadMetalakeSync(); + + assertTrue(manager.getCatalogRegistrationStates().isEmpty()); + assertTrue(manager.getMetalakeErrors().isEmpty()); + } + private CatalogConnectorManager createManager(ImmutableMap<String, String> configMap) throws Exception { return createManager(createCatalogConnectorFactory(), configMap); @@ -709,6 +729,7 @@ public class TestCatalogConnectorManager { when(catalogRegister.isTrinoStarted()).thenReturn(true); when(metalake.name()).thenReturn("test"); when(client.loadMetalake(any())).thenReturn(metalake); + Mockito.doReturn(new GravitinoMetalake[] {metalake}).when(client).listMetalakes(); when(catalogFactory.getSupportedCatalogProviders()).thenReturn(ImmutableSet.of("memory")); CatalogConnectorContext.Builder builder = mock(CatalogConnectorContext.Builder.class); when(catalogFactory.createCatalogConnectorContextBuilder(any())).thenReturn(builder); @@ -727,15 +748,15 @@ public class TestCatalogConnectorManager { } CatalogConnectorManager createManager(Map<String, String> extraConfig) { - ImmutableMap<String, String> configMap = - ImmutableMap.<String, String>builder() - .put("gravitino.uri", "http://127.0.0.1:8090") - .put("gravitino.metalake", "test") - .put("gravitino.use-single-metalake", "true") - .putAll(extraConfig) - .build(); + Map<String, String> defaults = new HashMap<>(); + defaults.put("gravitino.uri", "http://127.0.0.1:8090"); + defaults.put("gravitino.metalake", "test"); + defaults.put("gravitino.use-single-metalake", "true"); + defaults.putAll(extraConfig); + ImmutableMap<String, String> configMap = ImmutableMap.copyOf(defaults); CatalogConnectorManager manager = - new CatalogConnectorManager(catalogRegister, catalogFactory, null); + new CatalogConnectorManager( + catalogRegister, catalogFactory, (metalakeName, catalog) -> catalog); manager.config(new GravitinoConfig(configMap), client); return manager; }
