Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
s920243400
Rocketmq
提交
6381e6e1
R
Rocketmq
项目概览
s920243400
/
Rocketmq
与 Fork 源项目一致
Fork自
Apache RocketMQ / Rocketmq
通知
1
Star
1
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
R
Rocketmq
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
提交
6381e6e1
编写于
11月 16, 2021
作者:
D
dongeforever
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Rename the rpc header
上级
86b3ff7b
变更
13
隐藏空白更改
内联
并排
Showing
13 changed file
with
90 addition
and
59 deletion
+90
-59
common/src/main/java/org/apache/rocketmq/common/protocol/header/GetEarliestMsgStoretimeResponseHeader.java
...rotocol/header/GetEarliestMsgStoretimeResponseHeader.java
+2
-1
common/src/main/java/org/apache/rocketmq/common/protocol/header/GetMinOffsetResponseHeader.java
...mq/common/protocol/header/GetMinOffsetResponseHeader.java
+2
-1
common/src/main/java/org/apache/rocketmq/common/protocol/header/PullMessageResponseHeader.java
...tmq/common/protocol/header/PullMessageResponseHeader.java
+2
-1
common/src/main/java/org/apache/rocketmq/common/protocol/header/SearchOffsetResponseHeader.java
...mq/common/protocol/header/SearchOffsetResponseHeader.java
+2
-1
common/src/main/java/org/apache/rocketmq/common/rpc/RequestBuilder.java
...n/java/org/apache/rocketmq/common/rpc/RequestBuilder.java
+3
-3
common/src/main/java/org/apache/rocketmq/common/rpc/RpcClientImpl.java
...in/java/org/apache/rocketmq/common/rpc/RpcClientImpl.java
+13
-19
common/src/main/java/org/apache/rocketmq/common/rpc/RpcClientUtils.java
...n/java/org/apache/rocketmq/common/rpc/RpcClientUtils.java
+2
-2
common/src/main/java/org/apache/rocketmq/common/rpc/RpcHeader.java
...c/main/java/org/apache/rocketmq/common/rpc/RpcHeader.java
+46
-0
common/src/main/java/org/apache/rocketmq/common/rpc/RpcRequest.java
.../main/java/org/apache/rocketmq/common/rpc/RpcRequest.java
+3
-3
common/src/main/java/org/apache/rocketmq/common/rpc/RpcRequestHeader.java
...java/org/apache/rocketmq/common/rpc/RpcRequestHeader.java
+1
-14
common/src/main/java/org/apache/rocketmq/common/rpc/RpcResponse.java
...main/java/org/apache/rocketmq/common/rpc/RpcResponse.java
+11
-11
common/src/main/java/org/apache/rocketmq/common/rpc/TopicQueueRequestHeader.java
...g/apache/rocketmq/common/rpc/TopicQueueRequestHeader.java
+1
-1
remoting/src/main/java/org/apache/rocketmq/remoting/protocol/RemotingCommand.java
...rg/apache/rocketmq/remoting/protocol/RemotingCommand.java
+2
-2
未找到文件。
common/src/main/java/org/apache/rocketmq/common/protocol/header/GetEarliestMsgStoretimeResponseHeader.java
浏览文件 @
6381e6e1
...
...
@@ -20,11 +20,12 @@
*/
package
org.apache.rocketmq.common.protocol.header
;
import
org.apache.rocketmq.common.rpc.RpcHeader
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
import
org.apache.rocketmq.remoting.annotation.CFNotNull
;
import
org.apache.rocketmq.remoting.exception.RemotingCommandException
;
public
class
GetEarliestMsgStoretimeResponseHeader
implements
CommandCustom
Header
{
public
class
GetEarliestMsgStoretimeResponseHeader
extends
Rpc
Header
{
@CFNotNull
private
Long
timestamp
;
...
...
common/src/main/java/org/apache/rocketmq/common/protocol/header/GetMinOffsetResponseHeader.java
浏览文件 @
6381e6e1
...
...
@@ -20,11 +20,12 @@
*/
package
org.apache.rocketmq.common.protocol.header
;
import
org.apache.rocketmq.common.rpc.RpcHeader
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
import
org.apache.rocketmq.remoting.annotation.CFNotNull
;
import
org.apache.rocketmq.remoting.exception.RemotingCommandException
;
public
class
GetMinOffsetResponseHeader
implements
CommandCustom
Header
{
public
class
GetMinOffsetResponseHeader
extends
Rpc
Header
{
@CFNotNull
private
Long
offset
;
...
...
common/src/main/java/org/apache/rocketmq/common/protocol/header/PullMessageResponseHeader.java
浏览文件 @
6381e6e1
...
...
@@ -20,12 +20,13 @@
*/
package
org.apache.rocketmq.common.protocol.header
;
import
org.apache.rocketmq.common.rpc.RpcHeader
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
import
org.apache.rocketmq.remoting.annotation.CFNotNull
;
import
org.apache.rocketmq.remoting.annotation.CFNullable
;
import
org.apache.rocketmq.remoting.exception.RemotingCommandException
;
public
class
PullMessageResponseHeader
implements
CommandCustom
Header
{
public
class
PullMessageResponseHeader
extends
Rpc
Header
{
@CFNotNull
private
Long
suggestWhichBrokerId
;
@CFNotNull
...
...
common/src/main/java/org/apache/rocketmq/common/protocol/header/SearchOffsetResponseHeader.java
浏览文件 @
6381e6e1
...
...
@@ -20,11 +20,12 @@
*/
package
org.apache.rocketmq.common.protocol.header
;
import
org.apache.rocketmq.common.rpc.RpcHeader
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
import
org.apache.rocketmq.remoting.annotation.CFNotNull
;
import
org.apache.rocketmq.remoting.exception.RemotingCommandException
;
public
class
SearchOffsetResponseHeader
implements
CommandCustom
Header
{
public
class
SearchOffsetResponseHeader
extends
Rpc
Header
{
@CFNotNull
private
Long
offset
;
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/RequestBuilder.java
浏览文件 @
6381e6e1
...
...
@@ -14,17 +14,17 @@ public class RequestBuilder {
requestCodeMap
.
put
(
RequestCode
.
PULL_MESSAGE
,
PullMessageRequestHeader
.
class
);
}
public
static
CommonRpc
Header
buildCommonRpcHeader
(
int
requestCode
,
String
destBrokerName
)
{
public
static
RpcRequest
Header
buildCommonRpcHeader
(
int
requestCode
,
String
destBrokerName
)
{
return
buildCommonRpcHeader
(
requestCode
,
null
,
destBrokerName
);
}
public
static
CommonRpc
Header
buildCommonRpcHeader
(
int
requestCode
,
Boolean
oneway
,
String
destBrokerName
)
{
public
static
RpcRequest
Header
buildCommonRpcHeader
(
int
requestCode
,
Boolean
oneway
,
String
destBrokerName
)
{
Class
requestHeaderClass
=
requestCodeMap
.
get
(
requestCode
);
if
(
requestHeaderClass
==
null
)
{
throw
new
UnsupportedOperationException
(
"unknown "
+
requestCode
);
}
try
{
CommonRpcHeader
requestHeader
=
(
CommonRpc
Header
)
requestHeaderClass
.
newInstance
();
RpcRequestHeader
requestHeader
=
(
RpcRequest
Header
)
requestHeaderClass
.
newInstance
();
requestHeader
.
setCode
(
requestCode
);
requestHeader
.
setOneway
(
oneway
);
requestHeader
.
setBname
(
destBrokerName
);
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/RpcClientImpl.java
浏览文件 @
6381e6e1
package
org.apache.rocketmq.common.rpc
;
import
com.google.common.util.concurrent.Futures
;
import
com.sun.org.apache.xpath.internal.functions.FuncPosition
;
import
io.netty.util.concurrent.DefaultPromise
;
import
io.netty.util.concurrent.EventExecutor
;
import
io.netty.util.concurrent.ImmediateEventExecutor
;
import
io.netty.util.concurrent.Promise
;
import
org.apache.rocketmq.common.message.MessageQueue
;
...
...
@@ -20,10 +16,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand;
import
java.util.ArrayList
;
import
java.util.List
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.Executors
;
import
java.util.concurrent.Future
;
import
java.util.concurrent.ThreadFactory
;
public
class
RpcClientImpl
implements
RpcClient
{
...
...
@@ -68,12 +61,12 @@ public class RpcClientImpl implements RpcClient {
String
addr
=
getBrokerAddrByNameOrException
(
request
.
getHeader
().
bname
);
Promise
<
RpcResponse
>
rpcResponsePromise
=
null
;
try
{
switch
(
request
.
getCode
())
{
switch
(
request
.
get
Header
().
get
Code
())
{
case
RequestCode
.
PULL_MESSAGE
:
rpcResponsePromise
=
handlePullMessage
(
addr
,
request
,
timeoutMs
);
break
;
default
:
throw
new
RpcException
(
ResponseCode
.
REQUEST_CODE_NOT_SUPPORTED
,
"Unknown request code "
+
request
.
getCode
());
throw
new
RpcException
(
ResponseCode
.
REQUEST_CODE_NOT_SUPPORTED
,
"Unknown request code "
+
request
.
get
Header
().
get
Code
());
}
}
catch
(
RpcException
rpcException
)
{
throw
rpcException
;
...
...
@@ -135,7 +128,8 @@ public class RpcClientImpl implements RpcClient {
case
ResponseCode
.
PULL_OFFSET_MOVED
:
PullMessageResponseHeader
responseHeader
=
(
PullMessageResponseHeader
)
responseCommand
.
decodeCommandCustomHeader
(
PullMessageResponseHeader
.
class
);
rpcResponsePromise
.
setSuccess
(
new
RpcResponse
(
responseCommand
.
getCode
(),
responseHeader
,
responseCommand
.
getBody
()));
responseHeader
.
setCode
(
responseCommand
.
getCode
());
rpcResponsePromise
.
setSuccess
(
new
RpcResponse
(
responseHeader
,
responseCommand
.
getBody
()));
default
:
RpcResponse
rpcResponse
=
new
RpcResponse
(
new
RpcException
(
responseCommand
.
getCode
(),
"unexpected remote response code"
));
rpcResponsePromise
.
setSuccess
(
rpcResponse
);
...
...
@@ -162,11 +156,11 @@ public class RpcClientImpl implements RpcClient {
case
ResponseCode
.
SUCCESS
:
{
SearchOffsetResponseHeader
responseHeader
=
(
SearchOffsetResponseHeader
)
responseCommand
.
decodeCommandCustomHeader
(
SearchOffsetResponseHeader
.
class
);
return
new
RpcResponse
(
responseCommand
.
getCode
(),
responseHeader
,
responseCommand
.
getBody
());
responseHeader
.
setCode
(
responseCommand
.
getCode
());
return
new
RpcResponse
(
responseHeader
,
responseCommand
.
getBody
());
}
default
:{
RpcResponse
rpcResponse
=
new
RpcResponse
(
responseCommand
.
getCode
(),
null
,
null
);
rpcResponse
.
setException
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
RpcResponse
rpcResponse
=
new
RpcResponse
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
return
rpcResponse
;
}
}
...
...
@@ -183,11 +177,11 @@ public class RpcClientImpl implements RpcClient {
case
ResponseCode
.
SUCCESS
:
{
GetMinOffsetResponseHeader
responseHeader
=
(
GetMinOffsetResponseHeader
)
responseCommand
.
decodeCommandCustomHeader
(
GetMinOffsetResponseHeader
.
class
);
return
new
RpcResponse
(
responseCommand
.
getCode
(),
responseHeader
,
responseCommand
.
getBody
());
responseHeader
.
setCode
(
responseCommand
.
getCode
());
return
new
RpcResponse
(
responseHeader
,
responseCommand
.
getBody
());
}
default
:{
RpcResponse
rpcResponse
=
new
RpcResponse
(
responseCommand
.
getCode
(),
null
,
null
);
rpcResponse
.
setException
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
RpcResponse
rpcResponse
=
new
RpcResponse
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
return
rpcResponse
;
}
}
...
...
@@ -204,12 +198,12 @@ public class RpcClientImpl implements RpcClient {
case
ResponseCode
.
SUCCESS
:
{
GetEarliestMsgStoretimeResponseHeader
responseHeader
=
(
GetEarliestMsgStoretimeResponseHeader
)
responseCommand
.
decodeCommandCustomHeader
(
GetEarliestMsgStoretimeResponseHeader
.
class
);
return
new
RpcResponse
(
responseCommand
.
getCode
(),
responseHeader
,
responseCommand
.
getBody
());
responseHeader
.
setCode
(
responseCommand
.
getCode
());
return
new
RpcResponse
(
responseHeader
,
responseCommand
.
getBody
());
}
default
:{
RpcResponse
rpcResponse
=
new
RpcResponse
(
responseCommand
.
getCode
(),
null
,
null
);
rpcResponse
.
setException
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
RpcResponse
rpcResponse
=
new
RpcResponse
(
new
RpcException
(
responseCommand
.
getCode
(),
"unknown remote error"
));
return
rpcResponse
;
}
}
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/RpcClientUtils.java
浏览文件 @
6381e6e1
...
...
@@ -8,13 +8,13 @@ import java.nio.ByteBuffer;
public
class
RpcClientUtils
{
public
static
RemotingCommand
createCommandForRpcRequest
(
RpcRequest
rpcRequest
)
{
RemotingCommand
cmd
=
RemotingCommand
.
createRequestCommand
(
rpcRequest
.
getCode
(),
rpcRequest
.
getHeader
());
RemotingCommand
cmd
=
RemotingCommand
.
createRequestCommand
WithHeader
(
rpcRequest
.
getHeader
()
.
getCode
(),
rpcRequest
.
getHeader
());
cmd
.
setBody
(
encodeBody
(
rpcRequest
.
getBody
()));
return
cmd
;
}
public
static
RemotingCommand
createCommandForRpcResponse
(
RpcResponse
rpcResponse
)
{
RemotingCommand
cmd
=
RemotingCommand
.
createResponseCommand
(
rpcResponse
.
getCode
(),
rpcResponse
.
getHeader
());
RemotingCommand
cmd
=
RemotingCommand
.
createResponseCommand
WithHeader
(
rpcResponse
.
getCode
(),
rpcResponse
.
getHeader
());
cmd
.
setRemark
(
rpcResponse
.
getException
()
==
null
?
""
:
rpcResponse
.
getException
().
getMessage
());
cmd
.
setBody
(
encodeBody
(
rpcResponse
.
getBody
()));
return
cmd
;
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/RpcHeader.java
0 → 100644
浏览文件 @
6381e6e1
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package
org.apache.rocketmq.common.rpc
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
import
org.apache.rocketmq.remoting.exception.RemotingCommandException
;
public
class
RpcHeader
implements
CommandCustomHeader
{
protected
int
code
;
public
RpcHeader
()
{
}
public
RpcHeader
(
int
code
)
{
this
.
code
=
code
;
}
public
int
getCode
()
{
return
code
;
}
public
void
setCode
(
int
code
)
{
this
.
code
=
code
;
}
@Override
public
void
checkFields
()
throws
RemotingCommandException
{
}
}
common/src/main/java/org/apache/rocketmq/common/rpc/RpcRequest.java
浏览文件 @
6381e6e1
...
...
@@ -17,15 +17,15 @@
package
org.apache.rocketmq.common.rpc
;
public
class
RpcRequest
{
private
CommonRpc
Header
header
;
private
RpcRequest
Header
header
;
private
Object
body
;
public
RpcRequest
(
CommonRpc
Header
header
,
Object
body
)
{
public
RpcRequest
(
RpcRequest
Header
header
,
Object
body
)
{
this
.
header
=
header
;
this
.
body
=
body
;
}
public
CommonRpc
Header
getHeader
()
{
public
RpcRequest
Header
getHeader
()
{
return
header
;
}
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/
CommonRpc
Header.java
→
common/src/main/java/org/apache/rocketmq/common/rpc/
RpcRequest
Header.java
浏览文件 @
6381e6e1
...
...
@@ -16,16 +16,11 @@
*/
package
org.apache.rocketmq.common.rpc
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
public
abstract
class
CommonRpcHeader
implements
CommandCustomHeader
{
public
abstract
class
RpcRequestHeader
extends
RpcHeader
{
//the namespace name
protected
String
namespace
;
//if the data has been namespaced
protected
Boolean
namespaced
;
protected
int
code
;
//the abstract remote addr name, usually the physical broker name
protected
String
bname
;
...
...
@@ -62,12 +57,4 @@ public abstract class CommonRpcHeader implements CommandCustomHeader {
public
void
setNamespaced
(
Boolean
namespaced
)
{
this
.
namespaced
=
namespaced
;
}
public
int
getCode
()
{
return
code
;
}
public
void
setCode
(
int
code
)
{
this
.
code
=
code
;
}
}
common/src/main/java/org/apache/rocketmq/common/rpc/RpcResponse.java
浏览文件 @
6381e6e1
...
...
@@ -16,11 +16,8 @@
*/
package
org.apache.rocketmq.common.rpc
;
import
org.apache.rocketmq.remoting.CommandCustomHeader
;
public
class
RpcResponse
{
private
int
code
;
private
CommandCustomHeader
header
;
private
RpcHeader
header
;
private
Object
body
;
public
RpcException
exception
;
...
...
@@ -28,29 +25,32 @@ public class RpcResponse {
}
public
RpcResponse
(
int
code
,
CommandCustomHeader
header
,
byte
[]
body
)
{
this
.
code
=
code
;
public
RpcResponse
(
RpcHeader
header
,
byte
[]
body
)
{
this
.
header
=
header
;
this
.
body
=
body
;
}
public
RpcResponse
(
RpcException
rpcException
)
{
this
.
code
=
rpcException
.
getErrorCode
(
);
this
.
header
=
new
RpcHeader
(
rpcException
.
getErrorCode
()
);
this
.
exception
=
rpcException
;
}
public
int
getCode
()
{
return
code
;
public
RpcHeader
getHeader
()
{
return
header
;
}
public
CommandCustomHeader
getHeader
(
)
{
return
header
;
public
void
setHeader
(
RpcHeader
header
)
{
this
.
header
=
header
;
}
public
Object
getBody
()
{
return
body
;
}
public
void
setBody
(
Object
body
)
{
this
.
body
=
body
;
}
public
RpcException
getException
()
{
return
exception
;
}
...
...
common/src/main/java/org/apache/rocketmq/common/rpc/TopicQueueRequestHeader.java
浏览文件 @
6381e6e1
...
...
@@ -16,7 +16,7 @@
*/
package
org.apache.rocketmq.common.rpc
;
public
abstract
class
TopicQueueRequestHeader
extends
CommonRpc
Header
{
public
abstract
class
TopicQueueRequestHeader
extends
RpcRequest
Header
{
//Physical or Logical
protected
Boolean
physical
;
...
...
remoting/src/main/java/org/apache/rocketmq/remoting/protocol/RemotingCommand.java
浏览文件 @
6381e6e1
...
...
@@ -87,7 +87,7 @@ public class RemotingCommand {
}
public
static
RemotingCommand
createRequestCommand
(
int
code
,
CommandCustomHeader
customHeader
)
{
public
static
RemotingCommand
createRequestCommand
WithHeader
(
int
code
,
CommandCustomHeader
customHeader
)
{
RemotingCommand
cmd
=
new
RemotingCommand
();
cmd
.
setCode
(
code
);
cmd
.
customHeader
=
customHeader
;
...
...
@@ -95,7 +95,7 @@ public class RemotingCommand {
return
cmd
;
}
public
static
RemotingCommand
createResponseCommand
(
int
code
,
CommandCustomHeader
customHeader
)
{
public
static
RemotingCommand
createResponseCommand
WithHeader
(
int
code
,
CommandCustomHeader
customHeader
)
{
RemotingCommand
cmd
=
new
RemotingCommand
();
cmd
.
setCode
(
code
);
cmd
.
markResponseType
();
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录