This is an automated email from the ASF dual-hosted git repository.
menghaoranss 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 c215ada3e5b Fix single table loading with mixed protocol and storage
database types (#38854)
c215ada3e5b is described below
commit c215ada3e5bf9dfbd02f25f9a710ca1fb68c8015
Author: Haoran Meng <[email protected]>
AuthorDate: Tue Jun 16 12:07:27 2026 +0800
Fix single table loading with mixed protocol and storage database types
(#38854)
---
.../handler/update/LoadSingleTableExecutor.java | 32 ++++++++++++++++++----
.../update/LoadSingleTableExecutorTest.java | 27 +++++++++++++++++-
2 files changed, 52 insertions(+), 7 deletions(-)
diff --git
a/kernel/single/distsql/handler/src/main/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutor.java
b/kernel/single/distsql/handler/src/main/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutor.java
index 47ec1f7c302..0a627bdefc4 100644
---
a/kernel/single/distsql/handler/src/main/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutor.java
+++
b/kernel/single/distsql/handler/src/main/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutor.java
@@ -19,6 +19,7 @@ package
org.apache.shardingsphere.single.distsql.handler.update;
import lombok.Setter;
import
org.apache.shardingsphere.database.connector.core.metadata.database.metadata.DialectDatabaseMetaData;
+import org.apache.shardingsphere.database.connector.core.type.DatabaseType;
import
org.apache.shardingsphere.database.connector.core.type.DatabaseTypeRegistry;
import
org.apache.shardingsphere.database.exception.core.exception.syntax.table.TableExistsException;
import
org.apache.shardingsphere.distsql.handler.engine.update.rdl.rule.spi.database.type.DatabaseRuleCreateExecutor;
@@ -63,7 +64,7 @@ public final class LoadSingleTableExecutor implements
DatabaseRuleCreateExecutor
ShardingSpherePreconditions.checkNotEmpty(database.getResourceMetaData().getStorageUnits(),
() -> new EmptyStorageUnitException(database.getName()));
database.checkStorageUnitsExisted(storageUnitNames);
}
- String defaultSchemaName = new
DatabaseTypeRegistry(database.getProtocolType()).getDefaultSchemaName(database.getName());
+ String defaultSchemaName = database.getDefaultSchemaName();
checkShouldNotExistLogicTables(sqlStatement, defaultSchemaName);
if (!storageUnitNames.isEmpty()) {
checkShouldExistActualTables(sqlStatement, storageUnitNames,
defaultSchemaName);
@@ -107,23 +108,30 @@ public final class LoadSingleTableExecutor implements
DatabaseRuleCreateExecutor
Collection<String> invalidDataSources =
storageUnitNames.stream().filter(each ->
!aggregatedDataSourceMap.containsKey(each)).collect(Collectors.toList());
ShardingSpherePreconditions.checkState(invalidDataSources.isEmpty(),
() -> new InvalidStorageUnitStatusException(String.format("`%s` is invalid,
please use `%s`",
String.join(",", invalidDataSources), String.join(",",
aggregatedDataSourceMap.keySet()))));
- Map<String, Map<String, Collection<String>>> actualTableNodes =
getActualTableNodes(storageUnitNames, aggregatedDataSourceMap);
+ Map<String, DatabaseType> storageUnitDatabaseTypes =
getStorageUnitDatabaseTypes(storageUnitNames, aggregatedDataSourceMap);
+ Map<String, Map<String, Collection<String>>> actualTableNodes =
getActualTableNodes(storageUnitNames, aggregatedDataSourceMap,
storageUnitDatabaseTypes);
for (SingleTableSegment each : sqlStatement.getTables()) {
String tableName = each.getTableName();
if (!SingleTableConstants.ASTERISK.equals(tableName)) {
String storageUnitName = each.getStorageUnitName();
- String schemaName = each.getSchemaName().isPresent() ?
each.getSchemaName().get() : defaultSchemaName;
-
ShardingSpherePreconditions.checkState(actualTableNodes.containsKey(storageUnitName)
&& actualTableNodes.get(storageUnitName).get(schemaName).contains(tableName),
+
ShardingSpherePreconditions.checkState(actualTableNodes.containsKey(storageUnitName),
+ () -> new TableNotFoundException(tableName,
storageUnitName));
+ DatabaseType storageUnitDatabaseType =
storageUnitDatabaseTypes.get(storageUnitName);
+ String actualSchemaName =
isSameDatabaseType(database.getProtocolType(), storageUnitDatabaseType) ?
defaultSchemaName
+ : new
DatabaseTypeRegistry(storageUnitDatabaseType).formatIdentifierPattern(defaultSchemaName);
+ String schemaName = each.getSchemaName().isPresent() ?
each.getSchemaName().get() : actualSchemaName;
+
ShardingSpherePreconditions.checkState(actualTableNodes.get(storageUnitName).get(schemaName).contains(tableName),
() -> new TableNotFoundException(tableName,
storageUnitName));
}
}
}
- private Map<String, Map<String, Collection<String>>>
getActualTableNodes(final Collection<String> storageUnitNames, final
Map<String, DataSource> aggregatedDataSourceMap) {
+ private Map<String, Map<String, Collection<String>>>
getActualTableNodes(final Collection<String> storageUnitNames, final
Map<String, DataSource> aggregatedDataSourceMap,
+
final Map<String, DatabaseType> storageUnitDatabaseTypes) {
Map<String, Map<String, Collection<String>>> result = new
LinkedHashMap<>(storageUnitNames.size(), 1F);
for (String each : storageUnitNames) {
DataSource dataSource = aggregatedDataSourceMap.get(each);
- Map<String, Collection<String>> schemaTableNames =
SingleTableDataNodeLoader.loadSchemaTableNames(database.getName(),
DatabaseTypeEngine.getStorageType(dataSource),
+ Map<String, Collection<String>> schemaTableNames =
SingleTableDataNodeLoader.loadSchemaTableNames(database.getName(),
storageUnitDatabaseTypes.get(each),
dataSource, each, Collections.emptySet(),
Collections.emptySet());
if (!schemaTableNames.isEmpty()) {
result.put(each, schemaTableNames);
@@ -132,6 +140,18 @@ public final class LoadSingleTableExecutor implements
DatabaseRuleCreateExecutor
return result;
}
+ private Map<String, DatabaseType> getStorageUnitDatabaseTypes(final
Collection<String> storageUnitNames, final Map<String, DataSource>
aggregatedDataSourceMap) {
+ Map<String, DatabaseType> result = new
LinkedHashMap<>(storageUnitNames.size(), 1F);
+ for (String each : storageUnitNames) {
+ result.put(each,
DatabaseTypeEngine.getStorageType(aggregatedDataSourceMap.get(each)));
+ }
+ return result;
+ }
+
+ private boolean isSameDatabaseType(final DatabaseType protocolType, final
DatabaseType storageUnitDatabaseType) {
+ return
protocolType.getType().equalsIgnoreCase(storageUnitDatabaseType.getType());
+ }
+
@Override
public SingleRuleConfiguration buildToBeCreatedRuleConfiguration(final
LoadSingleTableStatement sqlStatement) {
SingleRuleConfiguration result = new SingleRuleConfiguration();
diff --git
a/kernel/single/distsql/handler/src/test/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutorTest.java
b/kernel/single/distsql/handler/src/test/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutorTest.java
index af91ff41705..a5bd95a4a35 100644
---
a/kernel/single/distsql/handler/src/test/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutorTest.java
+++
b/kernel/single/distsql/handler/src/test/java/org/apache/shardingsphere/single/distsql/handler/update/LoadSingleTableExecutorTest.java
@@ -22,6 +22,7 @@ import
org.apache.shardingsphere.database.connector.core.type.DatabaseType;
import
org.apache.shardingsphere.database.connector.core.type.DatabaseTypeRegistry;
import
org.apache.shardingsphere.database.exception.core.exception.syntax.table.TableExistsException;
import
org.apache.shardingsphere.distsql.handler.engine.update.rdl.rule.spi.database.DatabaseRuleDefinitionExecutor;
+import org.apache.shardingsphere.infra.database.DatabaseTypeEngine;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.TableNotFoundException;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.datanode.InvalidDataNodeFormatException;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.resource.storageunit.EmptyStorageUnitException;
@@ -73,7 +74,7 @@ import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.when;
@ExtendWith(AutoMockExtension.class)
-@StaticMockSettings({SingleTableDataNodeLoader.class,
PhysicalDataSourceAggregator.class})
+@StaticMockSettings({DatabaseTypeEngine.class,
SingleTableDataNodeLoader.class, PhysicalDataSourceAggregator.class})
@MockitoSettings(strictness = Strictness.LENIENT)
class LoadSingleTableExecutorTest {
@@ -88,6 +89,8 @@ class LoadSingleTableExecutorTest {
void setUp() {
when(database.getName()).thenReturn("foo_db");
when(database.getProtocolType()).thenReturn(databaseType);
+ when(database.getDefaultSchemaName()).thenReturn("foo_db");
+
when(DatabaseTypeEngine.getStorageType(any(DataSource.class))).thenReturn(databaseType);
executor.setDatabase(database);
}
@@ -97,6 +100,9 @@ class LoadSingleTableExecutorTest {
final boolean
tableExists, final Class<? extends RuntimeException> expectedException) {
prepareStorageUnits();
prepareSchema(tableExists, schemaSupported ? "foo_schema" : "foo_db");
+ if (schemaSupported) {
+ when(database.getDefaultSchemaName()).thenReturn("foo_schema");
+ }
LoadSingleTableStatement sqlStatement = new
LoadSingleTableStatement(Collections.singleton(tableSegment));
if (schemaSupported) {
try (MockedConstruction<DatabaseTypeRegistry> ignored =
mockSchemaSupportedDatabaseTypeRegistry()) {
@@ -127,11 +133,21 @@ class LoadSingleTableExecutorTest {
prepareActualTableValidationScenario(Collections.singletonMap("foo_ds", new
MockedDataSource()), Collections.singletonMap("foo_schema",
Collections.singleton("foo_tbl")));
LoadSingleTableStatement sqlStatement = new
LoadSingleTableStatement(Arrays.asList(new SingleTableSegment("foo_ds",
"foo_schema", "foo_tbl"), new SingleTableSegment("*", "*")));
try (MockedConstruction<DatabaseTypeRegistry> ignored =
mockSchemaSupportedDatabaseTypeRegistry()) {
+ when(database.getDefaultSchemaName()).thenReturn("foo_schema");
prepareSchema(false, "foo_schema");
assertDoesNotThrow(() -> executor.checkBeforeUpdate(sqlStatement));
}
}
+ @Test
+ void assertCheckBeforeUpdateWithDifferentProtocolAndStorageTypes() {
+
when(DatabaseTypeEngine.getStorageType(any(DataSource.class))).thenReturn(TypedSPILoader.getService(DatabaseType.class,
"MySQL"));
+
prepareActualTableValidationScenario(Collections.singletonMap("foo_ds", new
MockedDataSource()), Collections.singletonMap("FOO_DB",
Collections.singleton("foo_tbl")));
+ try (MockedConstruction<DatabaseTypeRegistry> ignored =
mockStorageDatabaseTypeRegistry()) {
+ assertDoesNotThrow(() -> executor.checkBeforeUpdate(new
LoadSingleTableStatement(Collections.singletonList(new
SingleTableSegment("foo_ds", "foo_tbl")))));
+ }
+ }
+
@Test
void assertCheckBeforeUpdateWithAllTablesPattern() {
prepareSchema(false, "foo_db");
@@ -190,6 +206,15 @@ class LoadSingleTableExecutorTest {
});
}
+ private MockedConstruction<DatabaseTypeRegistry>
mockStorageDatabaseTypeRegistry() {
+ DialectDatabaseMetaData dialectDatabaseMetaData =
mock(DialectDatabaseMetaData.class, RETURNS_DEEP_STUBS);
+ return mockConstruction(DatabaseTypeRegistry.class, (mock, context) ->
{
+ when(mock.getDefaultSchemaName("foo_db")).thenReturn("foo_db");
+ when(mock.formatIdentifierPattern("foo_db")).thenReturn("FOO_DB");
+
when(mock.getDialectDatabaseMetaData()).thenReturn(dialectDatabaseMetaData);
+ });
+ }
+
private static Stream<Arguments>
assertCheckBeforeUpdateWithPreValidationFailureArguments() {
return Stream.of(
Arguments.of("schema unsupported rejects schema name", false,
new SingleTableSegment("foo_ds", "foo_schema", "foo_tbl"), false,
InvalidDataNodeFormatException.class),