This is an automated email from the ASF dual-hosted git repository.
menghaoran pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere.git
The following commit(s) were added to refs/heads/master by this push:
new 6b7aa50 fix CalciteContextException when execute federation query
after create new tables (#14298)
6b7aa50 is described below
commit 6b7aa50f27a0e3d10a91ef7f2c0d0ac532aee566
Author: Zhengqiang Duan <[email protected]>
AuthorDate: Fri Dec 24 20:14:40 2021 +0800
fix CalciteContextException when execute federation query after create new
tables (#14298)
* fix CalciteContextException when execute federation query after create
new tables
* refresh OptimizerPlannerContext in ContextManager
* fix unit test
---
.../context/refresher/MetaDataRefreshEngine.java | 10 +++++++--
.../infra/context/refresher/MetaDataRefresher.java | 6 +++++-
.../type/AlterIndexStatementSchemaRefresher.java | 6 ++++--
.../type/AlterTableStatementSchemaRefresher.java | 22 ++++++++++++-------
.../type/CreateIndexStatementSchemaRefresher.java | 6 ++++--
.../type/CreateTableStatementSchemaRefresher.java | 8 +++++--
.../type/CreateViewStatementSchemaRefresher.java | 6 ++++--
.../type/DropIndexStatementSchemaRefresher.java | 6 ++++--
.../type/DropTableStatementSchemaRefresher.java | 8 +++++--
.../type/DropViewStatementSchemaRefresher.java | 6 ++++--
.../federation/executor/FederationExecutor.java | 7 ++++--
.../customized/CustomizedFilterableExecutor.java | 6 ++++--
.../original/OriginalFilterableExecutor.java | 20 ++++++++++-------
.../table/FilterableTableScanExecutor.java | 25 ++++++++++------------
.../table/FilterableTableScanExecutorContext.java} | 21 +++++++++++-------
.../optimizer/context/OptimizerContext.java | 3 ---
.../optimizer/context/OptimizerContextFactory.java | 2 +-
.../context/planner/OptimizerPlannerContext.java | 2 --
.../planner/OptimizerPlannerContextFactory.java | 16 ++++++++++++++
.../driver/executor/DriverJDBCExecutor.java | 3 ++-
.../statement/ShardingSpherePreparedStatement.java | 2 +-
.../core/statement/ShardingSphereStatement.java | 2 +-
.../mode/manager/ContextManager.java | 18 ++++++++++------
.../ClusterContextManagerCoordinatorTest.java | 13 +++++++++--
.../communication/DatabaseCommunicationEngine.java | 3 ++-
25 files changed, 150 insertions(+), 77 deletions(-)
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefreshEngine.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefreshEngine.java
index 17fbd23..8d52580 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefreshEngine.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefreshEngine.java
@@ -19,6 +19,7 @@ package org.apache.shardingsphere.infra.context.refresher;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.infra.metadata.mapper.SQLStatementEventMapper;
@@ -29,6 +30,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.SQLStatement;
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
import java.util.Optional;
/**
@@ -44,11 +46,15 @@ public final class MetaDataRefreshEngine {
private final FederationSchemaMetaData federationMetaData;
+ private final Map<String, OptimizerPlannerContext> optimizerPlanners;
+
private final ConfigurationProperties props;
- public MetaDataRefreshEngine(final ShardingSphereMetaData schemaMetaData,
final FederationSchemaMetaData federationMetaData, final
ConfigurationProperties props) {
+ public MetaDataRefreshEngine(final ShardingSphereMetaData schemaMetaData,
final FederationSchemaMetaData federationMetaData,
+ final Map<String, OptimizerPlannerContext>
optimizerPlanners, final ConfigurationProperties props) {
this.schemaMetaData = schemaMetaData;
this.federationMetaData = federationMetaData;
+ this.optimizerPlanners = optimizerPlanners;
this.props = props;
}
@@ -62,7 +68,7 @@ public final class MetaDataRefreshEngine {
public void refresh(final SQLStatement sqlStatement, final
Collection<String> logicDataSourceNames) throws SQLException {
Optional<MetaDataRefresher> schemaRefresher =
TypedSPIRegistry.findRegisteredService(MetaDataRefresher.class,
sqlStatement.getClass().getSuperclass().getCanonicalName(), null);
if (schemaRefresher.isPresent()) {
- schemaRefresher.get().refresh(schemaMetaData, federationMetaData,
logicDataSourceNames, sqlStatement, props);
+ schemaRefresher.get().refresh(schemaMetaData, federationMetaData,
optimizerPlanners, logicDataSourceNames, sqlStatement, props);
}
Optional<SQLStatementEventMapper> sqlStatementEventMapper =
SQLStatementEventMapperFactory.newInstance(sqlStatement);
if (sqlStatementEventMapper.isPresent()) {
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefresher.java
index 7c21d25..a25a9ce 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/MetaDataRefresher.java
@@ -18,6 +18,7 @@
package org.apache.shardingsphere.infra.context.refresher;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.spi.typed.TypedSPI;
@@ -25,6 +26,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.SQLStatement;
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
/**
* ShardingSphere schema refresher.
@@ -38,10 +40,12 @@ public interface MetaDataRefresher<T extends SQLStatement>
extends TypedSPI {
*
* @param schemaMetaData schema meta data
* @param schema federation schema meta data
+ * @param optimizerPlanners optimizer planners
* @param logicDataSourceNames route data source names
* @param sqlStatement SQL statement
* @param props configuration properties
* @throws SQLException SQL exception
*/
- void refresh(ShardingSphereMetaData schemaMetaData,
FederationSchemaMetaData schema, Collection<String> logicDataSourceNames, T
sqlStatement, ConfigurationProperties props) throws SQLException;
+ void refresh(ShardingSphereMetaData schemaMetaData,
FederationSchemaMetaData schema, Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ Collection<String> logicDataSourceNames, T sqlStatement,
ConfigurationProperties props) throws SQLException;
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterIndexStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterIndexStatementSchemaRefresher.java
index 287a81c..5919ba2 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterIndexStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterIndexStatementSchemaRefresher.java
@@ -21,6 +21,7 @@ import com.google.common.base.Preconditions;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.infra.metadata.schema.ShardingSphereSchema;
@@ -33,6 +34,7 @@ import
org.apache.shardingsphere.sql.parser.sql.dialect.handler.ddl.AlterIndexSt
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
import java.util.Optional;
/**
@@ -41,8 +43,8 @@ import java.util.Optional;
public final class AlterIndexStatementSchemaRefresher implements
MetaDataRefresher<AlterIndexStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final AlterIndexStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
AlterIndexStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
Optional<IndexSegment> renameIndex =
AlterIndexStatementHandler.getRenameIndexSegment(sqlStatement);
if (!sqlStatement.getIndex().isPresent() || !renameIndex.isPresent()) {
return;
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterTableStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterTableStatementSchemaRefresher.java
index 4b6c07d..caac27c 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterTableStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/AlterTableStatementSchemaRefresher.java
@@ -20,6 +20,8 @@ package
org.apache.shardingsphere.infra.context.refresher.type;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContextFactory;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.builder.SchemaBuilderMaterials;
@@ -33,6 +35,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.AlterTableS
import java.sql.SQLException;
import java.util.Collection;
import java.util.Collections;
+import java.util.Map;
import java.util.Optional;
/**
@@ -41,28 +44,30 @@ import java.util.Optional;
public final class AlterTableStatementSchemaRefresher implements
MetaDataRefresher<AlterTableStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final AlterTableStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
AlterTableStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
String tableName =
sqlStatement.getTable().getTableName().getIdentifier().getValue();
if (sqlStatement.getRenameTable().isPresent()) {
- putTableMetaData(schemaMetaData, schema, logicDataSourceNames,
sqlStatement.getRenameTable().get().getTableName().getIdentifier().getValue(),
props);
- removeTableMetaData(schemaMetaData, schema, tableName);
+ putTableMetaData(schemaMetaData, schema, optimizerPlanners,
logicDataSourceNames,
sqlStatement.getRenameTable().get().getTableName().getIdentifier().getValue(),
props);
+ removeTableMetaData(schemaMetaData, schema, optimizerPlanners,
tableName);
} else {
- putTableMetaData(schemaMetaData, schema, logicDataSourceNames,
tableName, props);
+ putTableMetaData(schemaMetaData, schema, optimizerPlanners,
logicDataSourceNames, tableName, props);
}
SchemaAlteredEvent event = new
SchemaAlteredEvent(schemaMetaData.getName());
event.getAlteredTables().add(schemaMetaData.getSchema().get(tableName));
ShardingSphereEventBus.getInstance().post(event);
}
- private void removeTableMetaData(final ShardingSphereMetaData
schemaMetaData, final FederationSchemaMetaData schema, final String tableName) {
+ private void removeTableMetaData(final ShardingSphereMetaData
schemaMetaData, final FederationSchemaMetaData schema,
+ final Map<String,
OptimizerPlannerContext> optimizerPlanners, final String tableName) {
schemaMetaData.getSchema().remove(tableName);
schemaMetaData.getRuleMetaData().findRules(MutableDataNodeRule.class).forEach(each
-> each.remove(tableName));
schema.remove(tableName);
+ optimizerPlanners.put(schema.getName(),
OptimizerPlannerContextFactory.create(schema));
}
- private void putTableMetaData(final ShardingSphereMetaData schemaMetaData,
final FederationSchemaMetaData schema, final Collection<String>
logicDataSourceNames, final String tableName,
- final ConfigurationProperties props) throws
SQLException {
+ private void putTableMetaData(final ShardingSphereMetaData schemaMetaData,
final FederationSchemaMetaData schema, final Map<String,
OptimizerPlannerContext> optimizerPlanners,
+ final Collection<String>
logicDataSourceNames, final String tableName, final ConfigurationProperties
props) throws SQLException {
if (!containsInDataNodeContainedRule(tableName, schemaMetaData)) {
schemaMetaData.getRuleMetaData().findRules(MutableDataNodeRule.class).forEach(each
-> each.put(tableName, logicDataSourceNames.iterator().next()));
}
@@ -72,6 +77,7 @@ public final class AlterTableStatementSchemaRefresher
implements MetaDataRefresh
actualTableMetaData.ifPresent(tableMetaData -> {
schemaMetaData.getSchema().put(tableName, tableMetaData);
schema.put(tableMetaData);
+ optimizerPlanners.put(schema.getName(),
OptimizerPlannerContextFactory.create(schema));
});
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateIndexStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateIndexStatementSchemaRefresher.java
index 6286e57..c468ee8 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateIndexStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateIndexStatementSchemaRefresher.java
@@ -21,6 +21,7 @@ import com.google.common.base.Strings;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.builder.util.IndexMetaDataUtil;
@@ -30,6 +31,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.CreateIndex
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
/**
* Schema refresher for create index statement.
@@ -37,8 +39,8 @@ import java.util.Collection;
public final class CreateIndexStatementSchemaRefresher implements
MetaDataRefresher<CreateIndexStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final CreateIndexStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
CreateIndexStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
String indexName = null != sqlStatement.getIndex() ?
sqlStatement.getIndex().getIdentifier().getValue() :
IndexMetaDataUtil.getGeneratedLogicIndexName(sqlStatement.getColumns());
if (Strings.isNullOrEmpty(indexName)) {
return;
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateTableStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateTableStatementSchemaRefresher.java
index 65c05f5..c5f357f 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateTableStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateTableStatementSchemaRefresher.java
@@ -20,6 +20,8 @@ package
org.apache.shardingsphere.infra.context.refresher.type;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContextFactory;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.builder.SchemaBuilderMaterials;
@@ -33,6 +35,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.CreateTable
import java.sql.SQLException;
import java.util.Collection;
import java.util.Collections;
+import java.util.Map;
import java.util.Optional;
/**
@@ -41,8 +44,8 @@ import java.util.Optional;
public final class CreateTableStatementSchemaRefresher implements
MetaDataRefresher<CreateTableStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final CreateTableStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
CreateTableStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
String tableName =
sqlStatement.getTable().getTableName().getIdentifier().getValue();
if (!containsInDataNodeContainedRule(tableName, schemaMetaData)) {
schemaMetaData.getRuleMetaData().findRules(MutableDataNodeRule.class).forEach(each
-> each.put(tableName, logicDataSourceNames.iterator().next()));
@@ -53,6 +56,7 @@ public final class CreateTableStatementSchemaRefresher
implements MetaDataRefres
actualTableMetaData.ifPresent(tableMetaData -> {
schemaMetaData.getSchema().put(tableName, tableMetaData);
schema.put(tableMetaData);
+ optimizerPlanners.put(schema.getName(),
OptimizerPlannerContextFactory.create(schema));
SchemaAlteredEvent event = new
SchemaAlteredEvent(schemaMetaData.getName());
event.getAlteredTables().add(tableMetaData);
ShardingSphereEventBus.getInstance().post(event);
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateViewStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateViewStatementSchemaRefresher.java
index 6207157..ec521ad 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateViewStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/CreateViewStatementSchemaRefresher.java
@@ -20,6 +20,7 @@ package
org.apache.shardingsphere.infra.context.refresher.type;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.event.SchemaAlteredEvent;
@@ -30,6 +31,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.CreateViewS
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
/**
* Schema refresher for create view statement.
@@ -37,8 +39,8 @@ import java.util.Collection;
public final class CreateViewStatementSchemaRefresher implements
MetaDataRefresher<CreateViewStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final CreateViewStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
CreateViewStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
String viewName =
sqlStatement.getView().getTableName().getIdentifier().getValue();
TableMetaData tableMetaData = new TableMetaData();
schemaMetaData.getSchema().put(viewName, tableMetaData);
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropIndexStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropIndexStatementSchemaRefresher.java
index 4106ae6..db1c82c 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropIndexStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropIndexStatementSchemaRefresher.java
@@ -22,6 +22,7 @@ import com.google.common.base.Strings;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.infra.metadata.schema.ShardingSphereSchema;
@@ -34,6 +35,7 @@ import
org.apache.shardingsphere.sql.parser.sql.dialect.handler.ddl.DropIndexSta
import java.sql.SQLException;
import java.util.Collection;
import java.util.LinkedList;
+import java.util.Map;
import java.util.Optional;
import java.util.stream.Collectors;
@@ -43,8 +45,8 @@ import java.util.stream.Collectors;
public final class DropIndexStatementSchemaRefresher implements
MetaDataRefresher<DropIndexStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final DropIndexStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
DropIndexStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
Collection<String> indexNames = getIndexNames(sqlStatement);
Optional<SimpleTableSegment> simpleTableSegment =
DropIndexStatementHandler.getSimpleTableSegment(sqlStatement);
String tableName = simpleTableSegment.map(tableSegment ->
tableSegment.getTableName().getIdentifier().getValue()).orElse("");
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropTableStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropTableStatementSchemaRefresher.java
index d9f989f..fa895ba 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropTableStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropTableStatementSchemaRefresher.java
@@ -20,6 +20,8 @@ package
org.apache.shardingsphere.infra.context.refresher.type;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContextFactory;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.event.SchemaAlteredEvent;
@@ -29,6 +31,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.DropTableSt
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
/**
* Schema refresher for drop table statement.
@@ -36,12 +39,13 @@ import java.util.Collection;
public final class DropTableStatementSchemaRefresher implements
MetaDataRefresher<DropTableStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final DropTableStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
DropTableStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
SchemaAlteredEvent event = new
SchemaAlteredEvent(schemaMetaData.getName());
sqlStatement.getTables().forEach(each -> {
schemaMetaData.getSchema().remove(each.getTableName().getIdentifier().getValue());
schema.remove(each.getTableName().getIdentifier().getValue());
+ optimizerPlanners.put(schema.getName(),
OptimizerPlannerContextFactory.create(schema));
event.getDroppedTables().add(each.getTableName().getIdentifier().getValue());
});
Collection<MutableDataNodeRule> rules =
schemaMetaData.getRuleMetaData().findRules(MutableDataNodeRule.class);
diff --git
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropViewStatementSchemaRefresher.java
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropViewStatementSchemaRefresher.java
index 0dec356..8e40218 100644
---
a/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropViewStatementSchemaRefresher.java
+++
b/shardingsphere-infra/shardingsphere-infra-context/src/main/java/org/apache/shardingsphere/infra/context/refresher/type/DropViewStatementSchemaRefresher.java
@@ -20,6 +20,7 @@ package
org.apache.shardingsphere.infra.context.refresher.type;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.context.refresher.MetaDataRefresher;
import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.schema.event.SchemaAlteredEvent;
@@ -29,6 +30,7 @@ import
org.apache.shardingsphere.sql.parser.sql.common.statement.ddl.DropViewSta
import java.sql.SQLException;
import java.util.Collection;
+import java.util.Map;
/**
* Schema refresher for drop view statement.
@@ -36,8 +38,8 @@ import java.util.Collection;
public final class DropViewStatementSchemaRefresher implements
MetaDataRefresher<DropViewStatement> {
@Override
- public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Collection<String> logicDataSourceNames,
final DropViewStatement sqlStatement,
- final ConfigurationProperties props) throws
SQLException {
+ public void refresh(final ShardingSphereMetaData schemaMetaData, final
FederationSchemaMetaData schema, final Map<String, OptimizerPlannerContext>
optimizerPlanners,
+ final Collection<String> logicDataSourceNames, final
DropViewStatement sqlStatement, final ConfigurationProperties props) throws
SQLException {
SchemaAlteredEvent event = new
SchemaAlteredEvent(schemaMetaData.getName());
sqlStatement.getViews().forEach(each -> {
schemaMetaData.getSchema().remove(each.getTableName().getIdentifier().getValue());
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/FederationExecutor.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/FederationExecutor.java
index b3f74c1..ea8df15 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/FederationExecutor.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/FederationExecutor.java
@@ -22,10 +22,12 @@ import
org.apache.shardingsphere.infra.executor.sql.execute.engine.driver.jdbc.J
import
org.apache.shardingsphere.infra.executor.sql.execute.engine.driver.jdbc.JDBCExecutorCallback;
import
org.apache.shardingsphere.infra.executor.sql.execute.result.ExecuteResult;
import
org.apache.shardingsphere.infra.executor.sql.prepare.driver.DriverExecutionPrepareEngine;
+import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
+import java.util.Map;
/**
* Federation executor.
@@ -38,11 +40,12 @@ public interface FederationExecutor extends AutoCloseable {
* @param prepareEngine prepare engine
* @param callback callback
* @param logicSQL logic SQL
+ * @param metaDataMap meta data map
* @return result set
* @throws SQLException SQL exception
*/
- ResultSet executeQuery(DriverExecutionPrepareEngine<JDBCExecutionUnit,
Connection> prepareEngine,
- JDBCExecutorCallback<? extends ExecuteResult>
callback, LogicSQL logicSQL) throws SQLException;
+ ResultSet executeQuery(DriverExecutionPrepareEngine<JDBCExecutionUnit,
Connection> prepareEngine, JDBCExecutorCallback<? extends ExecuteResult>
callback,
+ LogicSQL logicSQL, Map<String,
ShardingSphereMetaData> metaDataMap) throws SQLException;
/**
* Get result set.
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/customized/CustomizedFilterableExecutor.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/customized/CustomizedFilterableExecutor.java
index 558491b..c9ded4e 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/customized/CustomizedFilterableExecutor.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/customized/CustomizedFilterableExecutor.java
@@ -31,11 +31,13 @@ import
org.apache.shardingsphere.infra.executor.sql.prepare.driver.DriverExecuti
import org.apache.shardingsphere.infra.federation.executor.FederationExecutor;
import
org.apache.shardingsphere.infra.federation.optimizer.ShardingSphereOptimizer;
import
org.apache.shardingsphere.infra.federation.optimizer.context.OptimizerContext;
+import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.sql.parser.sql.common.statement.SQLStatement;
import java.sql.Connection;
import java.sql.ResultSet;
import java.sql.SQLException;
+import java.util.Map;
/**
* Customized filterable executor.
@@ -52,8 +54,8 @@ public final class CustomizedFilterableExecutor implements
FederationExecutor {
}
@Override
- public ResultSet executeQuery(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
- final JDBCExecutorCallback<? extends
ExecuteResult> callback, final LogicSQL logicSQL) throws SQLException {
+ public ResultSet executeQuery(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
final JDBCExecutorCallback<? extends ExecuteResult> callback,
+ final LogicSQL logicSQL, final Map<String,
ShardingSphereMetaData> metaDataMap) throws SQLException {
// TODO
return null;
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/OriginalFilterableExecutor.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/OriginalFilterableExecutor.java
index 70e8506..3ab6a8f 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/OriginalFilterableExecutor.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/OriginalFilterableExecutor.java
@@ -28,7 +28,9 @@ import
org.apache.shardingsphere.infra.executor.sql.execute.result.ExecuteResult
import
org.apache.shardingsphere.infra.executor.sql.prepare.driver.DriverExecutionPrepareEngine;
import org.apache.shardingsphere.infra.federation.executor.FederationExecutor;
import
org.apache.shardingsphere.infra.federation.executor.original.table.FilterableTableScanExecutor;
+import
org.apache.shardingsphere.infra.federation.executor.original.table.FilterableTableScanExecutorContext;
import
org.apache.shardingsphere.infra.federation.optimizer.context.OptimizerContext;
+import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.sql.parser.sql.common.util.SQLUtil;
import java.sql.Connection;
@@ -38,6 +40,7 @@ import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.List;
+import java.util.Map;
/**
* Original filterable executor.
@@ -68,24 +71,25 @@ public final class OriginalFilterableExecutor implements
FederationExecutor {
}
@Override
- public ResultSet executeQuery(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
- final JDBCExecutorCallback<? extends
ExecuteResult> callback, final LogicSQL logicSQL) throws SQLException {
- PreparedStatement preparedStatement = createConnection(prepareEngine,
callback,
logicSQL.getParameters()).prepareStatement(SQLUtil.trimSemicolon(logicSQL.getSql()));
+ public ResultSet executeQuery(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
final JDBCExecutorCallback<? extends ExecuteResult> callback,
+ final LogicSQL logicSQL, final Map<String,
ShardingSphereMetaData> metaDataMap) throws SQLException {
+ PreparedStatement preparedStatement = createConnection(prepareEngine,
callback, logicSQL.getParameters(),
metaDataMap).prepareStatement(SQLUtil.trimSemicolon(logicSQL.getSql()));
setParameters(preparedStatement, logicSQL.getParameters());
this.statement = preparedStatement;
return preparedStatement.executeQuery();
}
- private Connection createConnection(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
- final JDBCExecutorCallback<? extends
ExecuteResult> callback, final List<Object> parameters) throws SQLException {
+ private Connection createConnection(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
final JDBCExecutorCallback<? extends ExecuteResult> callback,
+ final List<Object> parameters, final
Map<String, ShardingSphereMetaData> metaDataMap) throws SQLException {
Connection result = DriverManager.getConnection(CONNECTION_URL,
optimizerContext.getParserContexts().get(schemaName).getDialectProps());
- addSchema(result.unwrap(CalciteConnection.class), prepareEngine,
callback, parameters);
+ addSchema(result.unwrap(CalciteConnection.class), prepareEngine,
callback, parameters, metaDataMap);
return result;
}
private void addSchema(final CalciteConnection connection, final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
- final JDBCExecutorCallback<? extends ExecuteResult>
callback, final List<Object> parameters) throws SQLException {
- FilterableTableScanExecutor executor = new
FilterableTableScanExecutor(prepareEngine, jdbcExecutor, callback, props,
optimizerContext, schemaName, parameters);
+ final JDBCExecutorCallback<? extends ExecuteResult>
callback, final List<Object> parameters, final Map<String,
ShardingSphereMetaData> metaDataMap) throws SQLException {
+ FilterableTableScanExecutorContext executorContext = new
FilterableTableScanExecutorContext(schemaName, parameters, props, metaDataMap);
+ FilterableTableScanExecutor executor = new
FilterableTableScanExecutor(prepareEngine, jdbcExecutor, callback,
optimizerContext, executorContext);
FilterableSchema schema = new
FilterableSchema(optimizerContext.getFederationMetaData().getSchemas().get(schemaName),
executor);
connection.getRootSchema().add(schemaName, schema);
connection.setSchema(schemaName);
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutor.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutor.java
index 622b2c5..1f63f04 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutor.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutor.java
@@ -105,24 +105,18 @@ public final class FilterableTableScanExecutor {
private final JDBCExecutorCallback<? extends ExecuteResult> callback;
- private final ConfigurationProperties props;
-
private final OptimizerContext optimizerContext;
- private final String schemaName;
-
- private final List<Object> parameters;
+ private final FilterableTableScanExecutorContext executorContext;
- public FilterableTableScanExecutor(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
final JDBCExecutor jdbcExecutor,
- final JDBCExecutorCallback<? extends
ExecuteResult> callback, final ConfigurationProperties props,
- final OptimizerContext
optimizerContext, final String schemaName, final List<Object> parameters) {
+ public FilterableTableScanExecutor(final
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection> prepareEngine,
+ final JDBCExecutor jdbcExecutor, final
JDBCExecutorCallback<? extends ExecuteResult> callback,
+ final OptimizerContext
optimizerContext, final FilterableTableScanExecutorContext executorContext) {
this.jdbcExecutor = jdbcExecutor;
this.callback = callback;
this.prepareEngine = prepareEngine;
- this.props = props;
this.optimizerContext = optimizerContext;
- this.schemaName = schemaName;
- this.parameters = parameters;
+ this.executorContext = executorContext;
}
/**
@@ -133,12 +127,14 @@ public final class FilterableTableScanExecutor {
* @return query results
*/
public Enumerable<Object[]> execute(final FederationTableMetaData
tableMetaData, final FilterableTableScanContext scanContext) {
+ String schemaName = executorContext.getSchemaName();
DatabaseType databaseType =
DatabaseTypeRegistry.getTrunkDatabaseType(optimizerContext.getParserContexts().get(schemaName).getDatabaseType().getName());
SqlString sqlString = createSQLString(tableMetaData, scanContext,
databaseType);
// TODO replace sql parse with sql convert
SQLStatement sqlStatement = new
SQLStatementParserEngine(databaseType.getName(),
optimizerContext.getSqlParserRule()).parse(sqlString.getSql(), false);
- LogicSQL logicSQL = createLogicSQL(optimizerContext.getMetaDataMap(),
sqlString.getSql(), getParameters(sqlString.getDynamicParameters()),
sqlStatement);
- ShardingSphereMetaData metaData =
optimizerContext.getMetaDataMap().get(schemaName);
+ LogicSQL logicSQL = createLogicSQL(executorContext.getMetaDataMap(),
sqlString.getSql(), getParameters(sqlString.getDynamicParameters()),
sqlStatement);
+ ShardingSphereMetaData metaData =
executorContext.getMetaDataMap().get(schemaName);
+ ConfigurationProperties props = executorContext.getProps();
ExecutionContext context = new
KernelProcessor().generateExecutionContext(logicSQL, metaData, props);
try {
ExecutionGroupContext<JDBCExecutionUnit> executionGroupContext =
prepareEngine.prepare(context.getRouteContext(), context.getExecutionUnits());
@@ -187,12 +183,13 @@ public final class FilterableTableScanExecutor {
}
List<Object> result = new ArrayList<>();
for (Integer each : parameterIndices) {
- result.add(parameters.get(each));
+ result.add(executorContext.getParameters().get(each));
}
return result;
}
private RelNode createRelNode(final FederationTableMetaData tableMetaData,
final FilterableTableScanContext scanContext) {
+ String schemaName = executorContext.getSchemaName();
RelOptCluster relOptCluster =
optimizerContext.getPlannerContexts().get(schemaName).getConverter().getCluster();
RelOptSchema relOptSchema = (RelOptSchema)
optimizerContext.getPlannerContexts().get(schemaName).getValidator().getCatalogReader();
RelBuilder builder =
RelFactories.LOGICAL_BUILDER.create(relOptCluster,
relOptSchema).scan(tableMetaData.getName()).filter(scanContext.getFilters());
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutorContext.java
similarity index 61%
copy from
shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
copy to
shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutorContext.java
index f2d2aac..15a3db4 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-executor/src/main/java/org/apache/shardingsphere/infra/federation/executor/original/table/FilterableTableScanExecutorContext.java
@@ -15,23 +15,28 @@
* limitations under the License.
*/
-package org.apache.shardingsphere.infra.federation.optimizer.context.planner;
+package org.apache.shardingsphere.infra.federation.executor.original.table;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
-import org.apache.calcite.sql.validate.SqlValidator;
-import org.apache.calcite.sql2rel.SqlToRelConverter;
+import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
+import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
+
+import java.util.List;
+import java.util.Map;
/**
- * Optimize planner context.
+ * Filterable table scan executor context.
*/
@RequiredArgsConstructor
@Getter
-public final class OptimizerPlannerContext {
+public final class FilterableTableScanExecutorContext {
- private final SqlValidator validator;
+ private final String schemaName;
- private final SqlToRelConverter converter;
+ private final List<Object> parameters;
+
+ private final ConfigurationProperties props;
- // TODO refresh validator and converter after federation schema changed
+ private final Map<String, ShardingSphereMetaData> metaDataMap;
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContext.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContext.java
index d892f42..b557aff 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContext.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContext.java
@@ -22,7 +22,6 @@ import lombok.RequiredArgsConstructor;
import
org.apache.shardingsphere.infra.federation.optimizer.context.parser.OptimizerParserContext;
import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContext;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationMetaData;
-import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import org.apache.shardingsphere.parser.rule.SQLParserRule;
import java.util.Map;
@@ -38,8 +37,6 @@ public final class OptimizerContext {
private final FederationMetaData federationMetaData;
- private final Map<String, ShardingSphereMetaData> metaDataMap;
-
private final Map<String, OptimizerParserContext> parserContexts;
private final Map<String, OptimizerPlannerContext> plannerContexts;
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContextFactory.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContextFactory.java
index c8b15a1..ee73dbd 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContextFactory.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/OptimizerContextFactory.java
@@ -48,6 +48,6 @@ public final class OptimizerContextFactory {
Map<String, OptimizerParserContext> parserContexts =
OptimizerParserContextFactory.create(metaDataMap);
Map<String, OptimizerPlannerContext> plannerContexts =
OptimizerPlannerContextFactory.create(federationMetaData);
SQLParserRule sqlParserRule =
globalRuleMetaData.findSingleRule(SQLParserRule.class).orElse(null);
- return new OptimizerContext(sqlParserRule, federationMetaData,
metaDataMap, parserContexts, plannerContexts);
+ return new OptimizerContext(sqlParserRule, federationMetaData,
parserContexts, plannerContexts);
}
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
index f2d2aac..9e27a59 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContext.java
@@ -32,6 +32,4 @@ public final class OptimizerPlannerContext {
private final SqlValidator validator;
private final SqlToRelConverter converter;
-
- // TODO refresh validator and converter after federation schema changed
}
diff --git
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContextFactory.java
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContextFactory.java
index fda7b2a..0681730 100644
---
a/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContextFactory.java
+++
b/shardingsphere-infra/shardingsphere-infra-federation/shardingsphere-infra-federation-optimizer/src/main/java/org/apache/shardingsphere/infra/federation/optimizer/context/planner/OptimizerPlannerContextFactory.java
@@ -74,6 +74,22 @@ public final class OptimizerPlannerContextFactory {
return result;
}
+ /**
+ * Create optimizer planner context.
+ *
+ * @param schemaMetaData federation schema meta data
+ * @return created optimizer planner context
+ */
+ public static OptimizerPlannerContext create(final
FederationSchemaMetaData schemaMetaData) {
+ FederationSchema federationSchema = new
FederationSchema(schemaMetaData);
+ CalciteConnectionConfig connectionConfig = new
CalciteConnectionConfigImpl(createConnectionProperties());
+ RelDataTypeFactory relDataTypeFactory = new JavaTypeFactoryImpl();
+ CalciteCatalogReader catalogReader =
createCatalogReader(schemaMetaData.getName(), federationSchema,
relDataTypeFactory, connectionConfig);
+ SqlValidator validator = createValidator(catalogReader,
relDataTypeFactory, connectionConfig);
+ SqlToRelConverter converter = createConverter(catalogReader,
validator, relDataTypeFactory);
+ return new OptimizerPlannerContext(validator, converter);
+ }
+
private static Properties createConnectionProperties() {
Properties result = new Properties();
result.setProperty(CalciteConnectionProperty.TIME_ZONE.camelName(),
"UTC");
diff --git
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/DriverJDBCExecutor.java
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/DriverJDBCExecutor.java
index b04833b..5c38820 100644
---
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/DriverJDBCExecutor.java
+++
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/DriverJDBCExecutor.java
@@ -58,7 +58,8 @@ public final class DriverJDBCExecutor {
this.metaDataContexts = metaDataContexts;
this.jdbcExecutor = jdbcExecutor;
metadataRefreshEngine = new
MetaDataRefreshEngine(metaDataContexts.getMetaData(schemaName),
-
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName),
metaDataContexts.getProps());
+
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName),
+ metaDataContexts.getOptimizerContext().getPlannerContexts(),
metaDataContexts.getProps());
}
/**
diff --git
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSpherePreparedStatement.java
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSpherePreparedStatement.java
index 6f7d295..6c3b28f 100644
---
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSpherePreparedStatement.java
+++
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSpherePreparedStatement.java
@@ -223,7 +223,7 @@ public final class ShardingSpherePreparedStatement extends
AbstractPreparedState
private ResultSet executeFederationQuery(final LogicSQL logicSQL) throws
SQLException {
PreparedStatementExecuteQueryCallback callback = new
PreparedStatementExecuteQueryCallback(metaDataContexts.getMetaData(connection.getSchema()).getResource().getDatabaseType(),
sqlStatement,
SQLExecutorExceptionHandler.isExceptionThrown());
- return
executor.getFederationExecutor().executeQuery(createDriverExecutionPrepareEngine(),
callback, logicSQL);
+ return
executor.getFederationExecutor().executeQuery(createDriverExecutionPrepareEngine(),
callback, logicSQL, metaDataContexts.getMetaDataMap());
}
private DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection>
createDriverExecutionPrepareEngine() {
diff --git
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSphereStatement.java
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSphereStatement.java
index 0872bcf..9d1d64c 100644
---
a/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSphereStatement.java
+++
b/shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/jdbc/core/statement/ShardingSphereStatement.java
@@ -168,7 +168,7 @@ public final class ShardingSphereStatement extends
AbstractStatementAdapter {
private ResultSet executeFederationQuery(final LogicSQL logicSQL) throws
SQLException {
StatementExecuteQueryCallback callback = new
StatementExecuteQueryCallback(metaDataContexts.getMetaData(connection.getSchema()).getResource().getDatabaseType(),
executionContext.getSqlStatementContext().getSqlStatement(),
SQLExecutorExceptionHandler.isExceptionThrown());
- return
executor.getFederationExecutor().executeQuery(createDriverExecutionPrepareEngine(),
callback, logicSQL);
+ return
executor.getFederationExecutor().executeQuery(createDriverExecutionPrepareEngine(),
callback, logicSQL, metaDataContexts.getMetaDataMap());
}
private DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection>
createDriverExecutionPrepareEngine() {
diff --git
a/shardingsphere-mode/shardingsphere-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/ContextManager.java
b/shardingsphere-mode/shardingsphere-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/ContextManager.java
index b9a0634..e19b106 100644
---
a/shardingsphere-mode/shardingsphere-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/ContextManager.java
+++
b/shardingsphere-mode/shardingsphere-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/ContextManager.java
@@ -25,6 +25,7 @@ import
org.apache.shardingsphere.infra.config.datasource.DataSourceConfiguration
import org.apache.shardingsphere.infra.config.datasource.DataSourceConverter;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import org.apache.shardingsphere.infra.database.type.DatabaseType;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.planner.OptimizerPlannerContextFactory;
import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.resource.ShardingSphereResource;
@@ -125,8 +126,9 @@ public final class ContextManager implements AutoCloseable {
return;
}
MetaDataContexts newMetaDataContexts =
buildNewMetaDataContext(schemaName);
-
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().put(schemaName,
-
newMetaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName));
+ FederationSchemaMetaData schemaMetaData =
newMetaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName);
+
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().put(schemaName,
schemaMetaData);
+
metaDataContexts.getOptimizerContext().getPlannerContexts().put(schemaName,
OptimizerPlannerContextFactory.create(schemaMetaData));
metaDataContexts.getMetaDataMap().put(schemaName,
newMetaDataContexts.getMetaData(schemaName));
metaDataContexts.getMetaDataPersistService().ifPresent(optional ->
optional.getSchemaMetaDataService().persist(schemaName));
}
@@ -233,8 +235,9 @@ public final class ContextManager implements AutoCloseable {
metaDataContexts.getMetaData(schemaName).getRuleMetaData(),
schema);
Map<String, ShardingSphereMetaData> kernelMetaDataMap = new
HashMap<>(metaDataContexts.getMetaDataMap());
kernelMetaDataMap.put(schemaName, kernelMetaData);
-
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().put(schemaName,
- new FederationSchemaMetaData(schemaName, schema.getTables()));
+ FederationSchemaMetaData schemaMetaData = new
FederationSchemaMetaData(schemaName, schema.getTables());
+
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().put(schemaName,
schemaMetaData);
+
metaDataContexts.getOptimizerContext().getPlannerContexts().put(schemaName,
OptimizerPlannerContextFactory.create(schemaMetaData));
renewMetaDataContexts(rebuildMetaDataContexts(kernelMetaDataMap));
}
@@ -246,13 +249,16 @@ public final class ContextManager implements
AutoCloseable {
* @param deletedTable deleted table
*/
public void alterSchema(final String schemaName, final TableMetaData
changedTableMetaData, final String deletedTable) {
+ FederationSchemaMetaData schemaMetaData =
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName);
if (null != changedTableMetaData) {
metaDataContexts.getMetaData(schemaName).getSchema().put(changedTableMetaData.getName(),
changedTableMetaData);
-
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName).put(changedTableMetaData);
+ schemaMetaData.put(changedTableMetaData);
+
metaDataContexts.getOptimizerContext().getPlannerContexts().put(schemaName,
OptimizerPlannerContextFactory.create(schemaMetaData));
}
if (null != deletedTable) {
metaDataContexts.getMetaData(schemaName).getSchema().remove(deletedTable);
-
metaDataContexts.getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName).remove(deletedTable);
+ schemaMetaData.remove(deletedTable);
+
metaDataContexts.getOptimizerContext().getPlannerContexts().put(schemaName,
OptimizerPlannerContextFactory.create(schemaMetaData));
}
}
diff --git
a/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/test/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/ClusterContextManagerCoordinatorTest.java
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/test/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/ClusterContextManagerCoordinatorTest.java
index e83c841..5fec6be 100644
---
a/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/test/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/ClusterContextManagerCoordinatorTest.java
+++
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/test/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/ClusterContextManagerCoordinatorTest.java
@@ -29,6 +29,8 @@ import
org.apache.shardingsphere.infra.config.mode.PersistRepositoryConfiguratio
import
org.apache.shardingsphere.infra.config.properties.ConfigurationProperties;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationPropertyKey;
import org.apache.shardingsphere.infra.executor.kernel.ExecutorEngine;
+import
org.apache.shardingsphere.infra.federation.optimizer.context.OptimizerContext;
+import
org.apache.shardingsphere.infra.federation.optimizer.metadata.FederationSchemaMetaData;
import org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData;
import
org.apache.shardingsphere.infra.metadata.resource.ShardingSphereResource;
import
org.apache.shardingsphere.infra.metadata.rule.ShardingSphereRuleMetaData;
@@ -36,7 +38,6 @@ import
org.apache.shardingsphere.infra.metadata.schema.QualifiedSchema;
import org.apache.shardingsphere.infra.metadata.schema.ShardingSphereSchema;
import org.apache.shardingsphere.infra.metadata.schema.model.TableMetaData;
import org.apache.shardingsphere.infra.metadata.user.ShardingSphereUser;
-import
org.apache.shardingsphere.infra.federation.optimizer.context.OptimizerContext;
import org.apache.shardingsphere.infra.rule.ShardingSphereRule;
import org.apache.shardingsphere.mode.manager.ContextManager;
import
org.apache.shardingsphere.mode.manager.cluster.ClusterContextManagerBuilder;
@@ -110,7 +111,7 @@ public final class ClusterContextManagerCoordinatorTest {
contextManager = builder.build(configuration, new HashMap<>(), new
HashMap<>(), new LinkedList<>(), new Properties(), false, null, null);
contextManager.renewMetaDataContexts(new
MetaDataContexts(contextManager.getMetaDataContexts().getMetaDataPersistService().get(),
createMetaDataMap(), globalRuleMetaData,
mock(ExecutorEngine.class),
- new ConfigurationProperties(new Properties()),
mock(OptimizerContext.class, RETURNS_DEEP_STUBS)));
+ new ConfigurationProperties(new Properties()),
createOptimizerContext()));
contextManager.renewTransactionContexts(mock(TransactionContexts.class,
RETURNS_DEEP_STUBS));
coordinator = new
ClusterContextManagerCoordinator(metaDataPersistService, contextManager);
}
@@ -239,4 +240,12 @@ public final class ClusterContextManagerCoordinatorTest {
when(metaData.getRuleMetaData().getConfigurations()).thenReturn(Collections.emptyList());
return Collections.singletonMap("schema", metaData);
}
+
+ private OptimizerContext createOptimizerContext() {
+ OptimizerContext result = mock(OptimizerContext.class,
RETURNS_DEEP_STUBS);
+ Map<String, FederationSchemaMetaData> schemas = new HashMap<>(1, 1);
+ schemas.put("schema", new FederationSchemaMetaData("schema",
Collections.emptyMap()));
+ when(result.getFederationMetaData().getSchemas()).thenReturn(schemas);
+ return result;
+ }
}
diff --git
a/shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/DatabaseCommunicationEngine.java
b/shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/DatabaseCommunicationEngine.java
index e3b4948..966dd88 100644
---
a/shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/DatabaseCommunicationEngine.java
+++
b/shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/DatabaseCommunicationEngine.java
@@ -116,6 +116,7 @@ public final class DatabaseCommunicationEngine {
String schemaName =
backendConnection.getConnectionSession().getSchemaName();
metadataRefreshEngine = new MetaDataRefreshEngine(metaData,
ProxyContext.getInstance().getContextManager().getMetaDataContexts().getOptimizerContext().getFederationMetaData().getSchemas().get(schemaName),
+
ProxyContext.getInstance().getContextManager().getMetaDataContexts().getOptimizerContext().getPlannerContexts(),
ProxyContext.getInstance().getContextManager().getMetaDataContexts().getProps());
MetaDataContexts metaDataContexts =
ProxyContext.getInstance().getContextManager().getMetaDataContexts();
federationExecutor = FederationExecutorFactory.newInstance(schemaName,
metaDataContexts.getOptimizerContext(),
@@ -172,7 +173,7 @@ public final class DatabaseCommunicationEngine {
logicSQL.getSqlStatementContext().getSqlStatement(), this,
isReturnGeneratedKeys, SQLExecutorExceptionHandler.isExceptionThrown(), true);
backendConnection.setFederationExecutor(federationExecutor);
DriverExecutionPrepareEngine<JDBCExecutionUnit, Connection>
prepareEngine = createDriverExecutionPrepareEngine(isReturnGeneratedKeys,
metaDataContexts);
- return federationExecutor.executeQuery(prepareEngine, callback,
logicSQL);
+ return federationExecutor.executeQuery(prepareEngine, callback,
logicSQL, metaDataContexts.getMetaDataMap());
}
private ResponseHeader processExecuteFederation(final ResultSet resultSet,
final MetaDataContexts metaDataContexts) throws SQLException {