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;
     }

Reply via email to