Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
apache
Shardingsphere
提交
308f09f1
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 搜索 >>
未验证
提交
308f09f1
编写于
5月 21, 2020
作者:
Z
Zhang Yonglun
提交者:
GitHub
5月 21, 2020
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
#3556, carry PR#5194 to dev-4.x (#5733)
上级
8ad915c2
变更
5
隐藏空白更改
内联
并排
Showing
5 changed file
with
37 addition
and
10 deletion
+37
-10
sharding-proxy/sharding-proxy-backend/src/main/java/org/apache/shardingsphere/shardingproxy/backend/communication/jdbc/execute/JDBCExecuteEngine.java
...backend/communication/jdbc/execute/JDBCExecuteEngine.java
+17
-4
sharding-proxy/sharding-proxy-backend/src/main/java/org/apache/shardingsphere/shardingproxy/backend/response/update/UpdateResponse.java
...shardingproxy/backend/response/update/UpdateResponse.java
+5
-0
sharding-proxy/sharding-proxy-frontend/sharding-proxy-frontend-postgresql/src/main/java/org/apache/shardingsphere/shardingproxy/frontend/postgresql/command/query/binary/bind/PostgreSQLComBindExecutor.java
.../command/query/binary/bind/PostgreSQLComBindExecutor.java
+1
-1
sharding-proxy/sharding-proxy-frontend/sharding-proxy-frontend-postgresql/src/main/java/org/apache/shardingsphere/shardingproxy/frontend/postgresql/command/query/text/PostgreSQLComQueryExecutor.java
...gresql/command/query/text/PostgreSQLComQueryExecutor.java
+1
-1
shardingsphere-database-protocol/shardingsphere-database-protocol-postgresql/src/main/java/org/apache/shardingsphere/database/protocol/postgresql/packet/generic/PostgreSQLCommandCompletePacket.java
...resql/packet/generic/PostgreSQLCommandCompletePacket.java
+13
-4
未找到文件。
sharding-proxy/sharding-proxy-backend/src/main/java/org/apache/shardingsphere/shardingproxy/backend/communication/jdbc/execute/JDBCExecuteEngine.java
浏览文件 @
308f09f1
...
...
@@ -18,7 +18,6 @@
package
org.apache.shardingsphere.shardingproxy.backend.communication.jdbc.execute
;
import
lombok.Getter
;
import
org.apache.shardingsphere.underlying.executor.context.ExecutionContext
;
import
org.apache.shardingsphere.sharding.execute.sql.StatementExecuteUnit
;
import
org.apache.shardingsphere.sharding.execute.sql.execute.SQLExecuteTemplate
;
import
org.apache.shardingsphere.sharding.execute.sql.execute.threadlocal.ExecutorExceptionHandler
;
...
...
@@ -36,8 +35,11 @@ import org.apache.shardingsphere.shardingproxy.backend.response.query.QueryRespo
import
org.apache.shardingsphere.shardingproxy.backend.response.update.UpdateResponse
;
import
org.apache.shardingsphere.shardingproxy.context.ShardingProxyContext
;
import
org.apache.shardingsphere.sql.parser.binder.statement.SQLStatementContext
;
import
org.apache.shardingsphere.sql.parser.sql.statement.dml.DeleteStatement
;
import
org.apache.shardingsphere.sql.parser.sql.statement.dml.InsertStatement
;
import
org.apache.shardingsphere.sql.parser.sql.statement.dml.UpdateStatement
;
import
org.apache.shardingsphere.underlying.common.config.properties.ConfigurationPropertyKey
;
import
org.apache.shardingsphere.underlying.executor.context.ExecutionContext
;
import
org.apache.shardingsphere.underlying.executor.engine.InputGroup
;
import
java.sql.SQLException
;
...
...
@@ -75,12 +77,23 @@ public final class JDBCExecuteEngine implements SQLExecuteEngine {
boolean
isExceptionThrown
=
ExecutorExceptionHandler
.
isExceptionThrown
();
Collection
<
InputGroup
<
StatementExecuteUnit
>>
inputGroups
=
sqlExecutePrepareTemplate
.
getExecuteUnitGroups
(
executionContext
.
getExecutionUnits
(),
new
ProxyJDBCExecutePrepareCallback
(
backendConnection
,
jdbcExecutorWrapper
,
isReturnGeneratedKeys
));
Collection
<
ExecuteResponse
>
executeResponses
=
sqlExecuteTemplate
.
execute
((
Collection
)
inputGroups
,
Collection
<
ExecuteResponse
>
executeResponses
=
sqlExecuteTemplate
.
execute
((
Collection
)
inputGroups
,
new
ProxySQLExecuteCallback
(
sqlStatementContext
,
backendConnection
,
jdbcExecutorWrapper
,
isExceptionThrown
,
isReturnGeneratedKeys
,
true
),
new
ProxySQLExecuteCallback
(
sqlStatementContext
,
backendConnection
,
jdbcExecutorWrapper
,
isExceptionThrown
,
isReturnGeneratedKeys
,
false
));
ExecuteResponse
executeResponse
=
executeResponses
.
iterator
().
next
();
return
executeResponse
instanceof
ExecuteQueryResponse
?
getExecuteQueryResponse
(((
ExecuteQueryResponse
)
executeResponse
).
getQueryHeaders
(),
executeResponses
)
:
new
UpdateResponse
(
executeResponses
);
if
(
executeResponse
instanceof
ExecuteQueryResponse
)
{
return
getExecuteQueryResponse
(((
ExecuteQueryResponse
)
executeResponse
).
getQueryHeaders
(),
executeResponses
);
}
else
{
UpdateResponse
updateResponse
=
new
UpdateResponse
(
executeResponses
);
if
(
sqlStatementContext
.
getSqlStatement
()
instanceof
InsertStatement
)
{
updateResponse
.
setType
(
"INSERT"
);
}
else
if
(
sqlStatementContext
.
getSqlStatement
()
instanceof
DeleteStatement
)
{
updateResponse
.
setType
(
"DELETE"
);
}
else
if
(
sqlStatementContext
.
getSqlStatement
()
instanceof
UpdateStatement
)
{
updateResponse
.
setType
(
"UPDATE"
);
}
return
updateResponse
;
}
}
private
BackendResponse
getExecuteQueryResponse
(
final
List
<
QueryHeader
>
queryHeaders
,
final
Collection
<
ExecuteResponse
>
executeResponses
)
{
...
...
sharding-proxy/sharding-proxy-backend/src/main/java/org/apache/shardingsphere/shardingproxy/backend/response/update/UpdateResponse.java
浏览文件 @
308f09f1
...
...
@@ -18,6 +18,7 @@
package
org.apache.shardingsphere.shardingproxy.backend.response.update
;
import
lombok.Getter
;
import
lombok.Setter
;
import
org.apache.shardingsphere.shardingproxy.backend.communication.jdbc.execute.response.ExecuteResponse
;
import
org.apache.shardingsphere.shardingproxy.backend.communication.jdbc.execute.response.ExecuteUpdateResponse
;
import
org.apache.shardingsphere.shardingproxy.backend.response.BackendResponse
;
...
...
@@ -40,6 +41,10 @@ public final class UpdateResponse implements BackendResponse {
@Getter
private
long
updateCount
;
@Getter
@Setter
private
String
type
;
public
UpdateResponse
()
{
this
(
Collections
.
emptyList
());
}
...
...
sharding-proxy/sharding-proxy-frontend/sharding-proxy-frontend-postgresql/src/main/java/org/apache/shardingsphere/shardingproxy/frontend/postgresql/command/query/binary/bind/PostgreSQLComBindExecutor.java
浏览文件 @
308f09f1
...
...
@@ -103,7 +103,7 @@ public final class PostgreSQLComBindExecutor implements QueryCommandExecutor {
}
private
PostgreSQLCommandCompletePacket
createUpdatePacket
(
final
UpdateResponse
updateResponse
)
{
return
new
PostgreSQLCommandCompletePacket
();
return
new
PostgreSQLCommandCompletePacket
(
updateResponse
.
getType
(),
updateResponse
.
getUpdateCount
()
);
}
private
Optional
<
PostgreSQLRowDescriptionPacket
>
createQueryPacket
(
final
QueryResponse
queryResponse
)
{
...
...
sharding-proxy/sharding-proxy-frontend/sharding-proxy-frontend-postgresql/src/main/java/org/apache/shardingsphere/shardingproxy/frontend/postgresql/command/query/text/PostgreSQLComQueryExecutor.java
浏览文件 @
308f09f1
...
...
@@ -88,7 +88,7 @@ public final class PostgreSQLComQueryExecutor implements QueryCommandExecutor {
}
private
PostgreSQLCommandCompletePacket
createUpdatePacket
(
final
UpdateResponse
updateResponse
)
{
return
new
PostgreSQLCommandCompletePacket
();
return
new
PostgreSQLCommandCompletePacket
(
updateResponse
.
getType
(),
updateResponse
.
getUpdateCount
()
);
}
private
Optional
<
PostgreSQLRowDescriptionPacket
>
createQueryPacket
(
final
QueryResponse
queryResponse
)
{
...
...
shardingsphere-database-protocol/shardingsphere-database-protocol-postgresql/src/main/java/org/apache/shardingsphere/database/protocol/postgresql/packet/generic/PostgreSQLCommandCompletePacket.java
浏览文件 @
308f09f1
...
...
@@ -30,13 +30,22 @@ public final class PostgreSQLCommandCompletePacket implements PostgreSQLPacket {
@Getter
private
final
char
messageType
=
PostgreSQLCommandPacketType
.
COMMAND_COMPLETE
.
getValue
();
private
final
String
sqlCommand
=
""
;
private
final
String
sqlCommand
;
private
final
int
rowCount
=
0
;
private
final
long
rowCount
;
public
PostgreSQLCommandCompletePacket
()
{
sqlCommand
=
""
;
rowCount
=
0
;
}
public
PostgreSQLCommandCompletePacket
(
final
String
sqlCommand
,
final
long
rowCount
)
{
this
.
sqlCommand
=
sqlCommand
;
this
.
rowCount
=
rowCount
;
}
@Override
public
void
write
(
final
PostgreSQLPacketPayload
payload
)
{
// TODO payload.writeStringNul(sqlCommand + " " + rowCount);
payload
.
writeStringNul
(
""
);
payload
.
writeStringNul
(
sqlCommand
+
" "
+
rowCount
);
}
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录