This is an automated email from the ASF dual-hosted git repository. shuwenwei pushed a commit to branch sync-generic-changes in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit dd5dd976980941fe86753c3e78c7f451b01cc4cb Author: shuwenwei <[email protected]> AuthorDate: Thu Sep 17 14:22:38 2026 +0800 [Information Schema] Add SPI extension points for information schema --- .../InformationSchemaContentSupplierFactory.java | 36 +++++++++---- .../AdditionalInformationSchemaProvider.java | 59 ++++++++++++++++++++++ ...dditionalInformationSchemaProviderRegistry.java | 43 ++++++++++++++++ .../DataNodeLocationSupplierFactory.java | 9 ++++ .../commons/schema/table/InformationSchema.java | 38 ++++++++++++++ 5 files changed, 175 insertions(+), 10 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/InformationSchemaContentSupplierFactory.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/InformationSchemaContentSupplierFactory.java index d39f9f396a8..f4fb59f7ca7 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/InformationSchemaContentSupplierFactory.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/source/relational/InformationSchemaContentSupplierFactory.java @@ -89,6 +89,8 @@ import org.apache.iotdb.db.queryengine.plan.execution.IQueryExecution; import org.apache.iotdb.db.queryengine.plan.execution.config.metadata.relational.ShowCreateViewTask; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanGraphPrinter; import org.apache.iotdb.db.queryengine.plan.relational.function.DataNodeTableBuiltinTableFunction; +import org.apache.iotdb.db.queryengine.plan.relational.information.AdditionalInformationSchemaProvider; +import org.apache.iotdb.db.queryengine.plan.relational.information.AdditionalInformationSchemaProviderRegistry; import org.apache.iotdb.db.queryengine.plan.relational.planner.node.InformationSchemaTableScanNode; import org.apache.iotdb.db.queryengine.plan.relational.planner.node.TableDiskUsageInformationSchemaTableScanNode; import org.apache.iotdb.db.queryengine.plan.relational.security.AccessControl; @@ -173,8 +175,17 @@ public class InformationSchemaContentSupplierFactory { final List<TSDataType> dataTypes, final UserEntity userEntity, final InformationSchemaTableScanNode node) { - String tableName = node.getQualifiedObjectName().getObjectName(); + final String tableName = node.getQualifiedObjectName().getObjectName(); try { + // Try additional providers first + for (final AdditionalInformationSchemaProvider provider : + AdditionalInformationSchemaProviderRegistry.getProviders()) { + final IInformationSchemaContentSupplier supplier = + provider.getContentSupplier(tableName, context, dataTypes, userEntity, node); + if (supplier != null) { + return supplier; + } + } switch (tableName) { case InformationSchema.QUERIES: return new QueriesSupplier(dataTypes, userEntity); @@ -384,14 +395,14 @@ public class InformationSchemaContentSupplierFactory { } } - private static class TableSupplier extends TsBlockSupplier { + public static class TableSupplier extends TsBlockSupplier { private final Iterator<Map.Entry<String, List<TTableInfo>>> dbIterator; private Iterator<TTableInfo> tableInfoIterator = null; - private TTableInfo currentTable; - private String dbName; + protected TTableInfo currentTable; + protected String dbName; private final UserEntity userEntity; - private TableSupplier(final List<TSDataType> dataTypes, final UserEntity userEntity) + protected TableSupplier(final List<TSDataType> dataTypes, final UserEntity userEntity) throws Exception { super(dataTypes); this.userEntity = userEntity; @@ -475,7 +486,7 @@ public class InformationSchemaContentSupplierFactory { } } - private static class ColumnSupplier extends TsBlockSupplier { + public static class ColumnSupplier extends TsBlockSupplier { private final Iterator<Map.Entry<String, Map<String, TableColumnDetailInfo>>> dbIterator; private Iterator<Map.Entry<String, TableColumnDetailInfo>> tableInfoIterator; private Iterator<TsTableColumnSchema> columnSchemaIterator; @@ -483,9 +494,10 @@ public class InformationSchemaContentSupplierFactory { private String tableName; private Set<String> preDeletedColumns; private Map<String, Byte> preAlteredColumns; + protected TsTableColumnSchema schema; private final UserEntity userEntity; - private ColumnSupplier(final List<TSDataType> dataTypes, final UserEntity userEntity) + protected ColumnSupplier(final List<TSDataType> dataTypes, final UserEntity userEntity) throws Exception { super(dataTypes); this.userEntity = userEntity; @@ -522,7 +534,6 @@ public class InformationSchemaContentSupplierFactory { @Override protected void constructLine() { - final TsTableColumnSchema schema = columnSchemaIterator.next(); columnBuilders[0].writeBinary(new Binary(dbName, TSFileConfig.STRING_CHARSET)); columnBuilders[1].writeBinary(new Binary(tableName, TSFileConfig.STRING_CHARSET)); columnBuilders[2].writeBinary( @@ -546,10 +557,14 @@ public class InformationSchemaContentSupplierFactory { columnBuilders[6].appendNull(); } resultBuilder.declarePosition(); + schema = null; } @Override public boolean hasNext() { + if (Objects.nonNull(schema)) { + return true; + } while (Objects.isNull(columnSchemaIterator) || !columnSchemaIterator.hasNext()) { while (Objects.isNull(tableInfoIterator) || !tableInfoIterator.hasNext()) { if (!dbIterator.hasNext()) { @@ -576,6 +591,7 @@ public class InformationSchemaContentSupplierFactory { } } } + schema = columnSchemaIterator.next(); return true; } } @@ -1660,13 +1676,13 @@ public class InformationSchemaContentSupplierFactory { } } - private abstract static class TsBlockSupplier implements IInformationSchemaContentSupplier { + public abstract static class TsBlockSupplier implements IInformationSchemaContentSupplier { protected final TsBlockBuilder resultBuilder; protected final ColumnBuilder[] columnBuilders; protected final AccessControl accessControl = AuthorityChecker.getAccessControl(); - private TsBlockSupplier(final List<TSDataType> dataTypes) { + protected TsBlockSupplier(final List<TSDataType> dataTypes) { this.resultBuilder = new TsBlockBuilder(dataTypes); this.columnBuilders = resultBuilder.getValueColumnBuilders(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProvider.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProvider.java new file mode 100644 index 00000000000..d0874676ca3 --- /dev/null +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProvider.java @@ -0,0 +1,59 @@ +/* + * 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.iotdb.db.queryengine.plan.relational.information; + +import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation; +import org.apache.iotdb.commons.audit.UserEntity; +import org.apache.iotdb.db.queryengine.execution.operator.OperatorContext; +import org.apache.iotdb.db.queryengine.execution.operator.source.relational.InformationSchemaContentSupplierFactory.IInformationSchemaContentSupplier; +import org.apache.iotdb.db.queryengine.plan.relational.planner.node.InformationSchemaTableScanNode; + +import org.apache.tsfile.enums.TSDataType; + +import java.util.List; + +/** Provides additional content and execution locations for information schema tables. */ +public interface AdditionalInformationSchemaProvider { + + /** + * Returns the content supplier for the table, or {@code null} if the table is not handled. + * + * @throws Exception if the content supplier cannot be created. + */ + default IInformationSchemaContentSupplier getContentSupplier( + final String tableName, + final OperatorContext context, + final List<TSDataType> dataTypes, + final UserEntity userEntity, + final InformationSchemaTableScanNode node) + throws Exception { + return null; + } + + /** Returns the execution location for the table, or {@code null} if the table is not handled. */ + default List<TDataNodeLocation> getTableLocation(final String tableName) { + return null; + } + + enum InformationSchemaTableLocation { + LOCAL_DATA_NODE, + ALL_READABLE_DATA_NODES + } +} diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProviderRegistry.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProviderRegistry.java new file mode 100644 index 00000000000..0f3f22136c3 --- /dev/null +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/information/AdditionalInformationSchemaProviderRegistry.java @@ -0,0 +1,43 @@ +/* + * 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.iotdb.db.queryengine.plan.relational.information; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.ServiceLoader; + +public class AdditionalInformationSchemaProviderRegistry { + + private static final List<AdditionalInformationSchemaProvider> PROVIDERS = new ArrayList<>(); + + static { + for (final AdditionalInformationSchemaProvider provider : + ServiceLoader.load(AdditionalInformationSchemaProvider.class)) { + PROVIDERS.add(provider); + } + } + + private AdditionalInformationSchemaProviderRegistry() {} + + public static List<AdditionalInformationSchemaProvider> getProviders() { + return Collections.unmodifiableList(PROVIDERS); + } +} diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/DataNodeLocationSupplierFactory.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/DataNodeLocationSupplierFactory.java index 629518313e2..b4d8e8d7a43 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/DataNodeLocationSupplierFactory.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/DataNodeLocationSupplierFactory.java @@ -33,6 +33,8 @@ import org.apache.iotdb.db.protocol.client.ConfigNodeClient; import org.apache.iotdb.db.protocol.client.ConfigNodeClientManager; import org.apache.iotdb.db.protocol.client.ConfigNodeInfo; import org.apache.iotdb.db.queryengine.common.DataNodeEndPoints; +import org.apache.iotdb.db.queryengine.plan.relational.information.AdditionalInformationSchemaProvider; +import org.apache.iotdb.db.queryengine.plan.relational.information.AdditionalInformationSchemaProviderRegistry; import org.apache.iotdb.rpc.TSStatusCode; import org.apache.thrift.TException; @@ -127,6 +129,13 @@ public class DataNodeLocationSupplierFactory { @Override public List<TDataNodeLocation> getDataNodeLocations(final String tableName) { + for (final AdditionalInformationSchemaProvider provider : + AdditionalInformationSchemaProviderRegistry.getProviders()) { + final List<TDataNodeLocation> location = provider.getTableLocation(tableName); + if (location != null) { + return location; + } + } switch (tableName) { case InformationSchema.QUERIES: case InformationSchema.TABLE_DISK_USAGE: diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/schema/table/InformationSchema.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/schema/table/InformationSchema.java index 462d9da8eab..92f8950a386 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/schema/table/InformationSchema.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/schema/table/InformationSchema.java @@ -32,7 +32,10 @@ import java.util.HashMap; import java.util.HashSet; import java.util.Locale; import java.util.Map; +import java.util.ServiceLoader; import java.util.Set; +import java.util.function.BiConsumer; +import java.util.function.Function; public class InformationSchema { public static final String INFORMATION_DATABASE = "information_schema"; @@ -494,6 +497,41 @@ public class InformationSchema { tablesThatSupportPushDownLimitOffset.add(TABLE_DISK_USAGE); } + // ==================== SPI extension point ==================== + + static { + for (final InformationSchemaExtension extension : + ServiceLoader.load(InformationSchemaExtension.class)) { + extension.registerAdditionalInformationSchemaTables(schemaTables::put); + extension.enhanceInformationSchemaTables(schemaTables::get); + } + } + + /** Extends the information schema by adding new tables or enhancing existing tables. */ + public interface InformationSchemaExtension { + + /** + * Registers extension-specific information schema tables that do not exist in the built-in + * schema. + * + * @param registerFunc accepts the table name and its schema definition + */ + default void registerAdditionalInformationSchemaTables( + final BiConsumer<String, TsTable> registerFunc) { + // Do nothing by default + } + + /** + * Enhances built-in information schema tables, for example by adding extension-specific + * columns. The provider returns the mutable schema definition of the requested table. + * + * @param tableProvider returns the schema definition for a built-in table name + */ + default void enhanceInformationSchemaTables(final Function<String, TsTable> tableProvider) { + // Do nothing by default + } + } + public static Map<String, TsTable> getSchemaTables() { return schemaTables; }
