Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
s920243400
Rocketmq
提交
1dde4fe7
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看板
提交
1dde4fe7
编写于
2月 19, 2019
作者:
D
duhenglucky
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Polish rebalance process in real push mode
上级
3c264c8f
变更
4
显示空白变更内容
内联
并排
Showing
4 changed file
with
6 addition
and
10 deletion
+6
-10
client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalancePushImpl.java
...ache/rocketmq/client/impl/consumer/RebalancePushImpl.java
+1
-6
snode/src/main/java/org/apache/rocketmq/snode/client/impl/SubscriptionManagerImpl.java
...e/rocketmq/snode/client/impl/SubscriptionManagerImpl.java
+4
-1
snode/src/main/java/org/apache/rocketmq/snode/processor/HeartbeatProcessor.java
...g/apache/rocketmq/snode/processor/HeartbeatProcessor.java
+1
-2
snode/src/main/java/org/apache/rocketmq/snode/processor/SendMessageProcessor.java
...apache/rocketmq/snode/processor/SendMessageProcessor.java
+0
-1
未找到文件。
client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalancePushImpl.java
浏览文件 @
1dde4fe7
...
...
@@ -16,7 +16,6 @@
*/
package
org.apache.rocketmq.client.impl.consumer
;
import
java.util.HashSet
;
import
java.util.List
;
import
java.util.Set
;
import
java.util.concurrent.TimeUnit
;
...
...
@@ -59,11 +58,7 @@ public class RebalancePushImpl extends RebalanceImpl {
log
.
info
(
"{} Rebalance changed, also update version: {}, {}"
,
topic
,
subscriptionData
.
getSubVersion
(),
newVersion
);
subscriptionData
.
setSubVersion
(
newVersion
);
Set
<
MessageQueue
>
queueIdSet
=
new
HashSet
<
MessageQueue
>();
for
(
MessageQueue
messageQueue
:
mqAll
)
{
queueIdSet
.
add
(
messageQueue
);
}
subscriptionData
.
setMessageQueueSet
(
queueIdSet
);
subscriptionData
.
setMessageQueueSet
(
mqDivided
);
int
currentQueueCount
=
this
.
processQueueTable
.
size
();
if
(
currentQueueCount
!=
0
)
{
int
pullThresholdForTopic
=
this
.
defaultMQPushConsumerImpl
.
getDefaultMQPushConsumer
().
getPullThresholdForTopic
();
...
...
snode/src/main/java/org/apache/rocketmq/snode/client/impl/SubscriptionManagerImpl.java
浏览文件 @
1dde4fe7
...
...
@@ -44,6 +44,7 @@ public class SubscriptionManagerImpl implements SubscriptionManager {
@Override
public
void
registerPushSession
(
Set
<
SubscriptionData
>
subscriptionDataSet
,
RemotingChannel
remotingChannel
,
String
groupId
)
{
log
.
debug
(
"Before ConsumerGroup: {} RemotingChannel: {} subscription: {}"
,
groupId
,
remotingChannel
.
remoteAddress
(),
subscriptionDataSet
);
Set
<
MessageQueue
>
prevSubSet
=
this
.
clientSubscriptionTable
.
get
(
remotingChannel
);
Set
<
MessageQueue
>
keySet
=
new
HashSet
<>();
for
(
SubscriptionData
subscriptionData
:
subscriptionDataSet
)
{
...
...
@@ -56,7 +57,7 @@ public class SubscriptionManagerImpl implements SubscriptionManager {
Set
<
RemotingChannel
>
prev
=
pushTable
.
putIfAbsent
(
messageQueue
,
clientSet
);
clientSet
=
prev
!=
null
?
prev
:
clientSet
;
}
log
.
debug
(
"Register push session message queue: {}, group: {} remoting: {}"
,
messageQueue
,
groupId
,
remotingChannel
.
remoteAddress
());
log
.
info
(
"Register push session message queue: {}, group: {} remoting: {}"
,
messageQueue
,
groupId
,
remotingChannel
.
remoteAddress
());
clientSet
.
add
(
remotingChannel
);
}
}
...
...
@@ -64,6 +65,7 @@ public class SubscriptionManagerImpl implements SubscriptionManager {
if
(
keySet
.
size
()
>
0
)
{
this
.
clientSubscriptionTable
.
put
(
remotingChannel
,
keySet
);
}
if
(
prevSubSet
!=
null
)
{
for
(
MessageQueue
messageQueue
:
prevSubSet
)
{
if
(!
keySet
.
contains
(
messageQueue
))
{
...
...
@@ -75,6 +77,7 @@ public class SubscriptionManagerImpl implements SubscriptionManager {
}
}
}
log
.
debug
(
"After ConsumerGroup: {} RemotingChannel: {} subscription: {}"
,
groupId
,
remotingChannel
.
remoteAddress
(),
this
.
clientSubscriptionTable
.
get
(
remotingChannel
));
}
@Override
...
...
snode/src/main/java/org/apache/rocketmq/snode/processor/HeartbeatProcessor.java
浏览文件 @
1dde4fe7
...
...
@@ -71,7 +71,7 @@ public class HeartbeatProcessor implements RequestProcessor {
private
RemotingCommand
register
(
RemotingChannel
remotingChannel
,
RemotingCommand
request
)
{
HeartbeatData
heartbeatData
=
HeartbeatData
.
decode
(
request
.
getBody
(),
HeartbeatData
.
class
);
log
.
info
(
"heartbeatData: {}"
,
heartbeatData
);
log
.
debug
(
"heartbeatData: {}"
,
heartbeatData
);
Channel
channel
=
null
;
Attribute
<
Client
>
clientAttribute
=
null
;
if
(
remotingChannel
instanceof
NettyChannelHandlerContextImpl
)
{
...
...
@@ -85,7 +85,6 @@ public class HeartbeatProcessor implements RequestProcessor {
client
.
setClientRole
(
ClientRole
.
Producer
);
this
.
snodeController
.
getProducerManager
().
register
(
producerData
.
getGroupName
(),
client
);
}
Set
<
String
>
groupSet
=
new
HashSet
<>();
for
(
ConsumerData
consumerData
:
heartbeatData
.
getConsumerDataSet
())
{
client
.
setClientRole
(
ClientRole
.
Consumer
);
...
...
snode/src/main/java/org/apache/rocketmq/snode/processor/SendMessageProcessor.java
浏览文件 @
1dde4fe7
...
...
@@ -99,7 +99,6 @@ public class SendMessageProcessor implements RequestProcessor {
remotingChannel
.
reply
(
data
);
this
.
snodeController
.
getMetricsService
().
recordRequestSize
(
stringBuffer
.
toString
(),
request
.
getBody
().
length
);
if
(
data
.
getCode
()
==
ResponseCode
.
SUCCESS
&&
isNeedPush
)
{
log
.
info
(
"Send message response: {}"
,
data
);
this
.
snodeController
.
getPushService
().
pushMessage
(
sendMessageRequestHeader
,
message
,
data
);
}
}
else
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录