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 d3287502986fed165a444427a7196589fc371609 Author: diqiu50 <[email protected]> AuthorDate: Fri Aug 21 17:05:30 2026 +0800 [#12546] improvement(trino-connector): Scope the system table registry to its connector The registry was static while each table is bound to one CatalogConnectorManager, so any other connector in the JVM could take it over. Make it an instance field threaded through the system connector, which also removes the need to register the tables only on the coordinator. --- .../connector/GravitinoConnectorFactory440.java | 6 ++- .../connector/GravitinoSystemConnector440.java | 12 ++++-- .../connector/GravitinoConnectorFactory446.java | 6 ++- .../connector/GravitinoSystemConnector446.java | 12 ++++-- .../connector/GravitinoConnectorFactory452.java | 6 ++- .../connector/GravitinoSystemConnector452.java | 12 ++++-- .../connector/GravitinoConnectorFactory469.java | 6 ++- .../connector/GravitinoSystemConnector469.java | 12 ++++-- .../connector/GravitinoConnectorFactory478.java | 6 ++- .../connector/GravitinoSystemConnector478.java | 12 ++++-- .../trino/connector/GravitinoConnectorFactory.java | 7 ++-- .../connector/system/GravitinoSystemConnector.java | 33 ++++++++++++++-- .../system/GravitinoSystemConnectorMetadata.java | 25 ++++++++---- .../system/table/GravitinoSystemTableFactory.java | 45 ++++++++++++++++------ .../table/TestGravitinoSystemStatusTables.java | 29 ++++++++++++++ 15 files changed, 178 insertions(+), 51 deletions(-) diff --git a/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory440.java b/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory440.java index f167dfad92..06e52fc6eb 100644 --- a/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory440.java +++ b/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory440.java @@ -22,6 +22,7 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoConnectorFactory440 extends GravitinoConnectorFactory { @@ -46,7 +47,8 @@ public class GravitinoConnectorFactory440 extends GravitinoConnectorFactory { @Override protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector440(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector440(storedProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector440.java b/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector440.java index de902bd958..4b03a47ac3 100644 --- a/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector440.java +++ b/trino-connector/trino-connector-440-445/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector440.java @@ -28,12 +28,14 @@ import io.trino.spi.connector.ConnectorSplitManager; import io.trino.spi.connector.SchemaTableName; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoSystemConnector440 extends GravitinoSystemConnector { public GravitinoSystemConnector440( - GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { - super(gravitinoStoredProcedureFactory); + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + super(gravitinoStoredProcedureFactory, systemTableFactory); } @Override @@ -43,11 +45,15 @@ public class GravitinoSystemConnector440 extends GravitinoSystemConnector { @Override protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider440(); + return new DatasourceProvider440(getSystemTableFactory()); } static class DatasourceProvider440 extends DatasourceProvider { + DatasourceProvider440(GravitinoSystemTableFactory systemTableFactory) { + super(systemTableFactory); + } + @Override protected ConnectorPageSource createPageSource(Page page) { return new SystemTablePageSource440(page); diff --git a/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory446.java b/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory446.java index 39e273e0a3..7e02341b76 100644 --- a/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory446.java +++ b/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory446.java @@ -22,6 +22,7 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoConnectorFactory446 extends GravitinoConnectorFactory { @@ -51,7 +52,8 @@ public class GravitinoConnectorFactory446 extends GravitinoConnectorFactory { @Override protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector446(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector446(storedProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector446.java b/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector446.java index 5eee4446e7..952a6fff45 100644 --- a/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector446.java +++ b/trino-connector/trino-connector-446-451/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector446.java @@ -28,12 +28,14 @@ import io.trino.spi.connector.ConnectorSplitManager; import io.trino.spi.connector.SchemaTableName; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoSystemConnector446 extends GravitinoSystemConnector { public GravitinoSystemConnector446( - GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { - super(gravitinoStoredProcedureFactory); + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + super(gravitinoStoredProcedureFactory, systemTableFactory); } @Override @@ -43,11 +45,15 @@ public class GravitinoSystemConnector446 extends GravitinoSystemConnector { @Override protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider446(); + return new DatasourceProvider446(getSystemTableFactory()); } static class DatasourceProvider446 extends DatasourceProvider { + DatasourceProvider446(GravitinoSystemTableFactory systemTableFactory) { + super(systemTableFactory); + } + @Override protected ConnectorPageSource createPageSource(Page page) { return new SystemTablePageSource446(page); diff --git a/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory452.java b/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory452.java index 01f73440e1..7cf665bcba 100644 --- a/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory452.java +++ b/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory452.java @@ -22,6 +22,7 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoConnectorFactory452 extends GravitinoConnectorFactory { @@ -51,7 +52,8 @@ public class GravitinoConnectorFactory452 extends GravitinoConnectorFactory { @Override protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector452(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector452(storedProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector452.java b/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector452.java index 3dd53e12d5..989ee44bd4 100644 --- a/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector452.java +++ b/trino-connector/trino-connector-452-468/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector452.java @@ -28,12 +28,14 @@ import io.trino.spi.connector.ConnectorSplitManager; import io.trino.spi.connector.SchemaTableName; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoSystemConnector452 extends GravitinoSystemConnector { public GravitinoSystemConnector452( - GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { - super(gravitinoStoredProcedureFactory); + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + super(gravitinoStoredProcedureFactory, systemTableFactory); } @Override @@ -43,11 +45,15 @@ public class GravitinoSystemConnector452 extends GravitinoSystemConnector { @Override protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider452(); + return new DatasourceProvider452(getSystemTableFactory()); } static class DatasourceProvider452 extends DatasourceProvider { + DatasourceProvider452(GravitinoSystemTableFactory systemTableFactory) { + super(systemTableFactory); + } + @Override protected ConnectorPageSource createPageSource(Page page) { return new SystemTablePageSource452(page); diff --git a/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory469.java b/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory469.java index 66b42c9ad0..7abe580823 100644 --- a/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory469.java +++ b/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory469.java @@ -22,6 +22,7 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoConnectorFactory469 extends GravitinoConnectorFactory { @@ -46,7 +47,8 @@ public class GravitinoConnectorFactory469 extends GravitinoConnectorFactory { @Override protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector469(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector469(storedProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector469.java b/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector469.java index b2d5483b27..4e25b3374f 100644 --- a/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector469.java +++ b/trino-connector/trino-connector-469-472/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector469.java @@ -28,12 +28,14 @@ import io.trino.spi.connector.ConnectorSplitManager; import io.trino.spi.connector.SchemaTableName; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoSystemConnector469 extends GravitinoSystemConnector { public GravitinoSystemConnector469( - GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { - super(gravitinoStoredProcedureFactory); + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + super(gravitinoStoredProcedureFactory, systemTableFactory); } @Override @@ -43,11 +45,15 @@ public class GravitinoSystemConnector469 extends GravitinoSystemConnector { @Override protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider469(); + return new DatasourceProvider469(getSystemTableFactory()); } static class DatasourceProvider469 extends DatasourceProvider { + DatasourceProvider469(GravitinoSystemTableFactory systemTableFactory) { + super(systemTableFactory); + } + @Override protected ConnectorPageSource createPageSource(Page page) { return new SystemTablePageSource469(page); diff --git a/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory478.java b/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory478.java index 9da0802a89..ec8adc60d8 100644 --- a/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory478.java +++ b/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnectorFactory478.java @@ -22,6 +22,7 @@ import org.apache.gravitino.client.GravitinoAdminClient; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoConnectorFactory478 extends GravitinoConnectorFactory { @@ -46,7 +47,8 @@ public class GravitinoConnectorFactory478 extends GravitinoConnectorFactory { @Override protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector478(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector478(storedProcedureFactory, systemTableFactory); } } diff --git a/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector478.java b/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector478.java index 208671e051..e3a78fffbc 100644 --- a/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector478.java +++ b/trino-connector/trino-connector-473-478/src/main/java/org/apache/gravitino/trino/connector/GravitinoSystemConnector478.java @@ -29,12 +29,14 @@ import io.trino.spi.connector.SchemaTableName; import io.trino.spi.connector.SourcePage; import org.apache.gravitino.trino.connector.system.GravitinoSystemConnector; import org.apache.gravitino.trino.connector.system.storedprocedure.GravitinoStoredProcedureFactory; +import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFactory; public class GravitinoSystemConnector478 extends GravitinoSystemConnector { public GravitinoSystemConnector478( - GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { - super(gravitinoStoredProcedureFactory); + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + super(gravitinoStoredProcedureFactory, systemTableFactory); } @Override @@ -44,11 +46,15 @@ public class GravitinoSystemConnector478 extends GravitinoSystemConnector { @Override protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider478(); + return new DatasourceProvider478(getSystemTableFactory()); } static class DatasourceProvider478 extends DatasourceProvider { + DatasourceProvider478(GravitinoSystemTableFactory systemTableFactory) { + super(systemTableFactory); + } + @Override protected ConnectorPageSource createPageSource(Page page) { return new SystemTablePageSource478(page); 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 917d9f3928..59b58e131b 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 @@ -186,7 +186,7 @@ public class GravitinoConnectorFactory implements ConnectorFactory { } GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory = new GravitinoStoredProcedureFactory(catalogConnectorManager, metalake); - return createSystemConnector(gravitinoStoredProcedureFactory); + return createSystemConnector(gravitinoStoredProcedureFactory, gravitinoSystemTableFactory); } } @@ -206,8 +206,9 @@ public class GravitinoConnectorFactory implements ConnectorFactory { } protected GravitinoSystemConnector createSystemConnector( - GravitinoStoredProcedureFactory storedProcedureFactory) { - return new GravitinoSystemConnector(storedProcedureFactory); + GravitinoStoredProcedureFactory storedProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { + return new GravitinoSystemConnector(storedProcedureFactory, systemTableFactory); } protected String getTrinoCatalogName(String metalakeName, String catalogName) { diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnector.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnector.java index 82c336f390..6ddec63d2f 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnector.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnector.java @@ -58,14 +58,28 @@ import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFac public class GravitinoSystemConnector implements Connector { private final GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory; + private final GravitinoSystemTableFactory systemTableFactory; /** * Constructs a new GravitinoSystemConnector. * * @param gravitinoStoredProcedureFactory the factory for creating stored procedures + * @param systemTableFactory the registry of system tables to expose */ - public GravitinoSystemConnector(GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory) { + public GravitinoSystemConnector( + GravitinoStoredProcedureFactory gravitinoStoredProcedureFactory, + GravitinoSystemTableFactory systemTableFactory) { this.gravitinoStoredProcedureFactory = gravitinoStoredProcedureFactory; + this.systemTableFactory = systemTableFactory; + } + + /** + * Retrieves the registry of system tables this connector exposes. + * + * @return the system table factory + */ + protected GravitinoSystemTableFactory getSystemTableFactory() { + return systemTableFactory; } @Override @@ -86,7 +100,7 @@ public class GravitinoSystemConnector implements Connector { } protected ConnectorMetadata createMetadata() { - return new GravitinoSystemConnectorMetadata(); + return new GravitinoSystemConnectorMetadata(systemTableFactory); } @Override @@ -106,7 +120,7 @@ public class GravitinoSystemConnector implements Connector { } protected ConnectorPageSourceProvider createPageSourceProvider() { - return new DatasourceProvider(); + return new DatasourceProvider(systemTableFactory); } /** The transaction handle for Gravitino system connector. */ @@ -118,6 +132,17 @@ public class GravitinoSystemConnector implements Connector { /** The datasource provider. */ public static class DatasourceProvider implements ConnectorPageSourceProvider { + private final GravitinoSystemTableFactory systemTableFactory; + + /** + * Constructs a new DatasourceProvider. + * + * @param systemTableFactory the registry the page data is read from + */ + public DatasourceProvider(GravitinoSystemTableFactory systemTableFactory) { + this.systemTableFactory = systemTableFactory; + } + @Override public ConnectorPageSource createPageSource( ConnectorTransactionHandle transaction, @@ -129,7 +154,7 @@ public class GravitinoSystemConnector implements Connector { SchemaTableName tableName = ((GravitinoSystemConnectorMetadata.SystemTableHandle) table).getName(); - Page page = GravitinoSystemTableFactory.loadPageData(tableName); + Page page = systemTableFactory.loadPageData(tableName); // Project the page down to the requested columns. Trino only expects the columns it asked // for, so handing it the whole row breaks any query that is not a SELECT *. diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnectorMetadata.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnectorMetadata.java index faeb7f35b7..1a95100401 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnectorMetadata.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/system/GravitinoSystemConnectorMetadata.java @@ -38,6 +38,17 @@ import org.apache.gravitino.trino.connector.system.table.GravitinoSystemTableFac /** An implementation of Apache Gravitino System connector Metadata */ public class GravitinoSystemConnectorMetadata implements ConnectorMetadata { + private final GravitinoSystemTableFactory systemTableFactory; + + /** + * Constructs a new GravitinoSystemConnectorMetadata. + * + * @param systemTableFactory the registry of system tables to expose + */ + public GravitinoSystemConnectorMetadata(GravitinoSystemTableFactory systemTableFactory) { + this.systemTableFactory = systemTableFactory; + } + @Override public List<String> listSchemaNames(ConnectorSession session) { return List.of(GravitinoSystemTable.SYSTEM_TABLE_SCHEMA_NAME); @@ -45,7 +56,7 @@ public class GravitinoSystemConnectorMetadata implements ConnectorMetadata { @Override public List<SchemaTableName> listTables(ConnectorSession session, Optional<String> schemaName) { - return GravitinoSystemTableFactory.SYSTEM_TABLES.keySet().stream().toList(); + return systemTableFactory.listTableNames(); } @Override @@ -54,16 +65,14 @@ public class GravitinoSystemConnectorMetadata implements ConnectorMetadata { SchemaTableName tableName, Optional<ConnectorTableVersion> startVersion, Optional<ConnectorTableVersion> endVersion) { - return GravitinoSystemTableFactory.SYSTEM_TABLES.get(tableName) != null - ? new SystemTableHandle(tableName) - : null; + return systemTableFactory.tableExists(tableName) ? new SystemTableHandle(tableName) : null; } @Override public ConnectorTableMetadata getTableMetadata( ConnectorSession session, ConnectorTableHandle table) { SchemaTableName tableName = ((SystemTableHandle) table).name; - return GravitinoSystemTableFactory.getTableMetaData(tableName); + return systemTableFactory.getTableMetaData(tableName); } @Override @@ -71,8 +80,7 @@ public class GravitinoSystemConnectorMetadata implements ConnectorMetadata { ConnectorSession session, ConnectorTableHandle tableHandle) { SchemaTableName tableName = ((SystemTableHandle) tableHandle).name; Map<String, ColumnHandle> columnHandles = new HashMap<>(); - List<ColumnMetadata> columns = - GravitinoSystemTableFactory.getTableMetaData(tableName).getColumns(); + List<ColumnMetadata> columns = systemTableFactory.getTableMetaData(tableName).getColumns(); for (int i = 0; i < columns.size(); i++) { columnHandles.put(columns.get(i).getName(), new SystemColumnHandle(i)); } @@ -83,7 +91,8 @@ public class GravitinoSystemConnectorMetadata implements ConnectorMetadata { public ColumnMetadata getColumnMetadata( ConnectorSession session, ConnectorTableHandle tableHandle, ColumnHandle columnHandle) { SchemaTableName tableName = ((SystemTableHandle) tableHandle).name; - return GravitinoSystemTableFactory.getTableMetaData(tableName) + return systemTableFactory + .getTableMetaData(tableName) .getColumns() .get(((SystemColumnHandle) columnHandle).index); } 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 832e89b78d..adc82731a1 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 @@ -23,6 +23,7 @@ import io.trino.spi.Page; import io.trino.spi.connector.ConnectorTableMetadata; import io.trino.spi.connector.SchemaTableName; import java.util.HashMap; +import java.util.List; import java.util.Map; import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; @@ -30,8 +31,11 @@ import org.apache.gravitino.trino.connector.catalog.CatalogConnectorManager; public class GravitinoSystemTableFactory { private final CatalogConnectorManager catalogConnectorManager; - /** Map of all registered system tables, keyed by their schema-qualified names. */ - public static final Map<SchemaTableName, GravitinoSystemTable> SYSTEM_TABLES = new HashMap<>(); + + // 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 + // registry would let any other connector in the JVM take it over. + private final Map<SchemaTableName, GravitinoSystemTable> systemTables = new HashMap<>(); /** * Constructs a new GravitinoSystemTableFactory. @@ -46,13 +50,13 @@ public class GravitinoSystemTableFactory { /** Register all the system tables */ private void registerSystemTables() { - SYSTEM_TABLES.put( + systemTables.put( GravitinoSystemTableCatalog.TABLE_NAME, new GravitinoSystemTableCatalog(catalogConnectorManager)); - SYSTEM_TABLES.put( + systemTables.put( GravitinoSystemTableCatalogStatus.TABLE_NAME, new GravitinoSystemTableCatalogStatus(catalogConnectorManager)); - SYSTEM_TABLES.put( + systemTables.put( GravitinoSystemTableLoadStatus.TABLE_NAME, new GravitinoSystemTableLoadStatus(catalogConnectorManager)); } @@ -64,9 +68,9 @@ public class GravitinoSystemTableFactory { * @return the page containing the table's data * @throws IllegalArgumentException if the table does not exist */ - public static Page loadPageData(SchemaTableName tableName) { - Preconditions.checkArgument(SYSTEM_TABLES.containsKey(tableName), "table does not exist"); - return SYSTEM_TABLES.get(tableName).loadPageData(); + public Page loadPageData(SchemaTableName tableName) { + Preconditions.checkArgument(systemTables.containsKey(tableName), "table does not exist"); + return systemTables.get(tableName).loadPageData(); } /** @@ -76,8 +80,27 @@ public class GravitinoSystemTableFactory { * @return the table metadata * @throws IllegalArgumentException if the table does not exist */ - public static ConnectorTableMetadata getTableMetaData(SchemaTableName tableName) { - Preconditions.checkArgument(SYSTEM_TABLES.containsKey(tableName), "table does not exist"); - return SYSTEM_TABLES.get(tableName).getTableMetaData(); + public ConnectorTableMetadata getTableMetaData(SchemaTableName tableName) { + Preconditions.checkArgument(systemTables.containsKey(tableName), "table does not exist"); + return systemTables.get(tableName).getTableMetaData(); + } + + /** + * Lists the names of every registered system table. + * + * @return the schema-qualified table names + */ + public List<SchemaTableName> listTableNames() { + return List.copyOf(systemTables.keySet()); + } + + /** + * Checks whether a system table is registered. + * + * @param tableName the schema-qualified name of the table + * @return true if the table exists, false otherwise + */ + public boolean tableExists(SchemaTableName tableName) { + return systemTables.containsKey(tableName); } } 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 43b60ecd7e..ae5ab372cd 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 @@ -122,6 +122,35 @@ public class TestGravitinoSystemStatusTables { assertEquals("{\"test\":\"Connection refused\"}", varchar(page, 5)); } + @Test + public void testEachFactoryOwnsItsTables() { + // The registry used to be static, so a second connector in the same JVM took it over and the + // system tables reported another manager's state. Each factory must now stand alone. + CatalogConnectorManager first = mock(CatalogConnectorManager.class); + when(first.getCatalogRegistrationStates()) + .thenReturn( + List.of( + CatalogRegistrationState.succeeded( + new GravitinoCatalog("prod", "memory", "memory", ImmutableMap.of(), 0L), + "memory"))); + CatalogConnectorManager second = mock(CatalogConnectorManager.class); + when(second.getCatalogRegistrationStates()).thenReturn(List.of()); + + GravitinoSystemTableFactory firstFactory = new GravitinoSystemTableFactory(first); + GravitinoSystemTableFactory secondFactory = new GravitinoSystemTableFactory(second); + + assertEquals( + 1, + firstFactory.loadPageData(GravitinoSystemTableCatalogStatus.TABLE_NAME).getPositionCount()); + assertEquals( + 0, + secondFactory + .loadPageData(GravitinoSystemTableCatalogStatus.TABLE_NAME) + .getPositionCount()); + assertTrue(firstFactory.tableExists(GravitinoSystemTableLoadStatus.TABLE_NAME)); + assertEquals(3, firstFactory.listTableNames().size()); + } + private static Page loadCatalogStatusPage(List<CatalogRegistrationState> states) { CatalogConnectorManager manager = mock(CatalogConnectorManager.class); when(manager.getCatalogRegistrationStates()).thenReturn(states);
