Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
apache
Shardingsphere
提交
86c5a512
Shardingsphere
项目概览
apache
/
Shardingsphere
通知
56
Star
3
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
Shardingsphere
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
未验证
提交
86c5a512
编写于
11月 15, 2020
作者:
L
Liang Zhang
提交者:
GitHub
11月 15, 2020
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Rename refreshSchema (#8166)
上级
674d5980
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
25 addition
and
22 deletion
+25
-22
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/AbstractStatementExecutor.java
...dingsphere/driver/executor/AbstractStatementExecutor.java
+10
-11
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/PreparedStatementExecutor.java
...dingsphere/driver/executor/PreparedStatementExecutor.java
+1
-1
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/StatementExecutor.java
...che/shardingsphere/driver/executor/StatementExecutor.java
+1
-1
shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/jdbc/JDBCDatabaseCommunicationEngine.java
...d/communication/jdbc/JDBCDatabaseCommunicationEngine.java
+13
-9
未找到文件。
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/AbstractStatementExecutor.java
浏览文件 @
86c5a512
...
...
@@ -33,7 +33,6 @@ import org.apache.shardingsphere.infra.metadata.schema.builder.SchemaBuilderMate
import
org.apache.shardingsphere.infra.metadata.schema.refresher.SchemaRefresher
;
import
org.apache.shardingsphere.infra.metadata.schema.refresher.SchemaRefresherFactory
;
import
org.apache.shardingsphere.infra.metadata.schema.refresher.spi.SchemaChangedNotifier
;
import
org.apache.shardingsphere.infra.route.context.RouteMapper
;
import
org.apache.shardingsphere.infra.route.context.RouteUnit
;
import
org.apache.shardingsphere.infra.rule.ShardingSphereRule
;
import
org.apache.shardingsphere.infra.rule.type.DataNodeContainedRule
;
...
...
@@ -76,30 +75,30 @@ public abstract class AbstractStatementExecutor {
}
@SuppressWarnings
({
"unchecked"
,
"rawtypes"
})
protected
final
void
refresh
TableMetaDat
a
(
final
ShardingSphereMetaData
metaData
,
final
SQLStatement
sqlStatement
,
final
Collection
<
RouteUnit
>
routeUnits
)
throws
SQLException
{
protected
final
void
refresh
Schem
a
(
final
ShardingSphereMetaData
metaData
,
final
SQLStatement
sqlStatement
,
final
Collection
<
RouteUnit
>
routeUnits
)
throws
SQLException
{
if
(
null
==
sqlStatement
)
{
return
;
}
Optional
<
SchemaRefresher
>
schemaRefresher
=
SchemaRefresherFactory
.
newInstance
(
sqlStatement
);
if
(
schemaRefresher
.
isPresent
())
{
Collection
<
String
>
routeDataSourceNames
=
routeUnits
.
stream
().
map
(
RouteUnit:
:
getDataSourceMapper
).
map
(
RouteMapper:
:
getLogicName
).
collect
(
Collectors
.
toList
());
schemaRefresher
.
get
().
refresh
(
metaData
.
getSchema
(),
routeDataSourceNames
,
sqlStatement
,
new
SchemaBuilderMaterials
(
metaDataContexts
.
getDatabaseType
(),
dataSourceMap
,
metaData
.
getRuleMetaData
().
getRules
(),
metaDataContexts
.
getProps
())
);
notify
PersistSchema
(
DefaultSchema
.
LOGIC_NAME
,
metaData
.
getSchema
());
Collection
<
String
>
routeDataSourceNames
=
routeUnits
.
stream
().
map
(
each
->
each
.
getDataSourceMapper
().
getLogicName
()
).
collect
(
Collectors
.
toList
());
SchemaBuilderMaterials
materials
=
new
SchemaBuilderMaterials
(
metaDataContexts
.
getDatabaseType
(),
dataSourceMap
,
metaData
.
getRuleMetaData
().
getRules
(),
metaDataContexts
.
getProps
());
schemaRefresher
.
get
().
refresh
(
metaData
.
getSchema
(),
routeDataSourceNames
,
sqlStatement
,
materials
);
notify
SchemaChanged
(
DefaultSchema
.
LOGIC_NAME
,
metaData
.
getSchema
());
}
}
private
void
notifySchemaChanged
(
final
String
schemaName
,
final
ShardingSphereSchema
schema
)
{
OrderedSPIRegistry
.
getRegisteredServices
(
Collections
.
singletonList
(
schema
),
SchemaChangedNotifier
.
class
).
values
().
forEach
(
each
->
each
.
notify
(
schemaName
,
schema
));
}
protected
final
boolean
executeAndRefreshMetaData
(
final
Collection
<
InputGroup
<
StatementExecuteUnit
>>
inputGroups
,
final
SQLStatement
sqlStatement
,
final
Collection
<
RouteUnit
>
routeUnits
,
final
SQLExecutorCallback
<
Boolean
>
sqlExecutorCallback
)
throws
SQLException
{
List
<
Boolean
>
result
=
sqlExecutor
.
execute
(
inputGroups
,
sqlExecutorCallback
);
refresh
TableMetaDat
a
(
metaDataContexts
.
getDefaultMetaData
(),
sqlStatement
,
routeUnits
);
refresh
Schem
a
(
metaDataContexts
.
getDefaultMetaData
(),
sqlStatement
,
routeUnits
);
return
null
!=
result
&&
!
result
.
isEmpty
()
&&
null
!=
result
.
get
(
0
)
&&
result
.
get
(
0
);
}
private
void
notifyPersistSchema
(
final
String
schemaName
,
final
ShardingSphereSchema
schema
)
{
OrderedSPIRegistry
.
getRegisteredServices
(
Collections
.
singletonList
(
schema
),
SchemaChangedNotifier
.
class
).
values
().
forEach
(
each
->
each
.
notify
(
schemaName
,
schema
));
}
/**
* Execute SQL.
*
...
...
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/PreparedStatementExecutor.java
浏览文件 @
86c5a512
...
...
@@ -81,7 +81,7 @@ public final class PreparedStatementExecutor extends AbstractStatementExecutor {
boolean
isExceptionThrown
=
ExecutorExceptionHandler
.
isExceptionThrown
();
SQLExecutorCallback
<
Integer
>
sqlExecutorCallback
=
createDefaultSQLExecutorCallbackWithInteger
(
isExceptionThrown
);
List
<
Integer
>
results
=
getSqlExecutor
().
execute
(
inputGroups
,
sqlExecutorCallback
);
refresh
TableMetaDat
a
(
getMetaDataContexts
().
getDefaultMetaData
(),
sqlStatementContext
.
getSqlStatement
(),
routeUnits
);
refresh
Schem
a
(
getMetaDataContexts
().
getDefaultMetaData
(),
sqlStatementContext
.
getSqlStatement
(),
routeUnits
);
return
isNeedAccumulate
(
getMetaDataContexts
().
getDefaultMetaData
().
getRuleMetaData
().
getRules
().
stream
().
filter
(
rule
->
rule
instanceof
DataNodeContainedRule
).
collect
(
Collectors
.
toList
()),
sqlStatementContext
)
?
accumulate
(
results
)
:
results
.
get
(
0
);
}
...
...
shardingsphere-jdbc/shardingsphere-jdbc-core/src/main/java/org/apache/shardingsphere/driver/executor/StatementExecutor.java
浏览文件 @
86c5a512
...
...
@@ -132,7 +132,7 @@ public final class StatementExecutor extends AbstractStatementExecutor {
}
};
List
<
Integer
>
results
=
getSqlExecutor
().
execute
(
inputGroups
,
sqlExecutorCallback
);
refresh
TableMetaDat
a
(
getMetaDataContexts
().
getDefaultMetaData
(),
sqlStatementContext
.
getSqlStatement
(),
routeUnits
);
refresh
Schem
a
(
getMetaDataContexts
().
getDefaultMetaData
(),
sqlStatementContext
.
getSqlStatement
(),
routeUnits
);
if
(
isNeedAccumulate
(
getMetaDataContexts
().
getDefaultMetaData
().
getRuleMetaData
().
getRules
().
stream
().
filter
(
rule
->
rule
instanceof
DataNodeContainedRule
).
collect
(
Collectors
.
toList
()),
sqlStatementContext
))
{
return
accumulate
(
results
);
...
...
shardingsphere-proxy/shardingsphere-proxy-backend/src/main/java/org/apache/shardingsphere/proxy/backend/communication/jdbc/JDBCDatabaseCommunicationEngine.java
浏览文件 @
86c5a512
...
...
@@ -18,8 +18,6 @@
package
org.apache.shardingsphere.proxy.backend.communication.jdbc
;
import
lombok.RequiredArgsConstructor
;
import
org.apache.shardingsphere.governance.core.event.GovernanceEventBus
;
import
org.apache.shardingsphere.governance.core.event.model.schema.SchemaPersistEvent
;
import
org.apache.shardingsphere.infra.binder.LogicSQL
;
import
org.apache.shardingsphere.infra.binder.statement.SQLStatementContext
;
import
org.apache.shardingsphere.infra.config.properties.ConfigurationPropertyKey
;
...
...
@@ -31,12 +29,13 @@ import org.apache.shardingsphere.infra.executor.sql.raw.execute.result.query.Que
import
org.apache.shardingsphere.infra.merge.MergeEngine
;
import
org.apache.shardingsphere.infra.merge.result.MergedResult
;
import
org.apache.shardingsphere.infra.metadata.ShardingSphereMetaData
;
import
org.apache.shardingsphere.infra.metadata.schema.ShardingSphereSchema
;
import
org.apache.shardingsphere.infra.metadata.schema.builder.SchemaBuilderMaterials
;
import
org.apache.shardingsphere.infra.metadata.schema.refresher.SchemaRefresher
;
import
org.apache.shardingsphere.infra.metadata.schema.refresher.SchemaRefresherFactory
;
import
org.apache.shardingsphere.infra.route.context.RouteMapper
;
import
org.apache.shardingsphere.infra.route.context.RouteUnit
;
import
org.apache.shardingsphere.infra.metadata.schema.refresher.spi.SchemaChangedNotifier
;
import
org.apache.shardingsphere.infra.rule.type.DataNodeContainedRule
;
import
org.apache.shardingsphere.infra.spi.ordered.OrderedSPIRegistry
;
import
org.apache.shardingsphere.proxy.backend.communication.DatabaseCommunicationEngine
;
import
org.apache.shardingsphere.proxy.backend.communication.jdbc.execute.SQLExecuteEngine
;
import
org.apache.shardingsphere.proxy.backend.context.ProxyContext
;
...
...
@@ -49,6 +48,7 @@ import org.apache.shardingsphere.sql.parser.sql.common.statement.SQLStatement;
import
java.sql.SQLException
;
import
java.util.ArrayList
;
import
java.util.Collection
;
import
java.util.Collections
;
import
java.util.List
;
import
java.util.Optional
;
import
java.util.stream.Collectors
;
...
...
@@ -90,26 +90,30 @@ public final class JDBCDatabaseCommunicationEngine implements DatabaseCommunicat
}
sqlExecuteEngine
.
checkExecutePrerequisites
(
executionContext
);
response
=
sqlExecuteEngine
.
execute
(
executionContext
);
Collection
<
String
>
routeDataSourceNames
=
executionContext
.
getRouteContext
().
getRouteUnits
().
stream
()
.
map
(
RouteUnit:
:
getDataSourceMapper
).
map
(
RouteMapper:
:
getLogicName
).
collect
(
Collectors
.
toList
());
refreshTableMetaData
(
executionContext
.
getSqlStatementContext
().
getSqlStatement
(),
routeDataSourceNames
);
refreshSchema
(
executionContext
);
return
merge
(
executionContext
.
getSqlStatementContext
());
}
@SuppressWarnings
({
"unchecked"
,
"rawtypes"
})
private
void
refreshTableMetaData
(
final
SQLStatement
sqlStatement
,
final
Collection
<
String
>
routeDataSourceNames
)
throws
SQLException
{
private
void
refreshSchema
(
final
ExecutionContext
executionContext
)
throws
SQLException
{
SQLStatement
sqlStatement
=
executionContext
.
getSqlStatementContext
().
getSqlStatement
();
if
(
null
==
sqlStatement
)
{
return
;
}
Optional
<
SchemaRefresher
>
schemaRefresher
=
SchemaRefresherFactory
.
newInstance
(
sqlStatement
);
if
(
schemaRefresher
.
isPresent
())
{
Collection
<
String
>
routeDataSourceNames
=
executionContext
.
getRouteContext
().
getRouteUnits
().
stream
().
map
(
each
->
each
.
getDataSourceMapper
().
getLogicName
()).
collect
(
Collectors
.
toList
());
SchemaBuilderMaterials
materials
=
new
SchemaBuilderMaterials
(
ProxyContext
.
getInstance
().
getMetaDataContexts
().
getDatabaseType
(),
metaData
.
getResource
().
getDataSources
(),
metaData
.
getRuleMetaData
().
getRules
(),
ProxyContext
.
getInstance
().
getMetaDataContexts
().
getProps
());
schemaRefresher
.
get
().
refresh
(
metaData
.
getSchema
(),
routeDataSourceNames
,
sqlStatement
,
materials
);
GovernanceEventBus
.
getInstance
().
post
(
new
SchemaPersistEvent
(
metaData
.
getName
(),
metaData
.
getSchema
()
));
notifySchemaChanged
(
metaData
.
getName
(),
metaData
.
getSchema
(
));
}
}
private
void
notifySchemaChanged
(
final
String
schemaName
,
final
ShardingSphereSchema
schema
)
{
OrderedSPIRegistry
.
getRegisteredServices
(
Collections
.
singletonList
(
schema
),
SchemaChangedNotifier
.
class
).
values
().
forEach
(
each
->
each
.
notify
(
schemaName
,
schema
));
}
private
BackendResponse
merge
(
final
SQLStatementContext
<?>
sqlStatementContext
)
throws
SQLException
{
if
(
response
instanceof
UpdateResponse
)
{
mergeUpdateCount
(
sqlStatementContext
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录