mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
This commit is contained in:
+1
-1
@@ -172,7 +172,7 @@ public class DefaultPullMessageResultHandler implements PullMessageResultHandler
|
||||
}
|
||||
|
||||
private boolean channelIsWritable(Channel channel, PullMessageRequestHeader requestHeader) {
|
||||
if (this.brokerController.getBrokerConfig().isNetWorkFlowController()) {
|
||||
if (this.brokerController.getBrokerConfig().isEnableNetWorkFlowControl()) {
|
||||
if (!channel.isWritable()) {
|
||||
log.warn("channel {} not writable ,cid {}", channel.remoteAddress(), requestHeader.getConsumerGroup());
|
||||
return false;
|
||||
|
||||
-2
@@ -83,10 +83,8 @@ public class PullMessageProcessorTest {
|
||||
SubscriptionGroupManager subscriptionGroupManager = new SubscriptionGroupManager(brokerController);
|
||||
pullMessageProcessor = new PullMessageProcessor(brokerController);
|
||||
Channel mockChannel = mock(Channel.class);
|
||||
when(mockChannel.isWritable()).thenReturn(true);
|
||||
when(mockChannel.remoteAddress()).thenReturn(new InetSocketAddress(1024));
|
||||
when(handlerContext.channel()).thenReturn(mockChannel);
|
||||
when(handlerContext.channel().isWritable()).thenReturn(true);
|
||||
when(brokerController.getSubscriptionGroupManager()).thenReturn(subscriptionGroupManager);
|
||||
brokerController.getTopicConfigManager().getTopicConfigTable().put(topic, new TopicConfig());
|
||||
clientChannelInfo = new ClientChannelInfo(mockChannel);
|
||||
|
||||
@@ -191,7 +191,7 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
*/
|
||||
private long brokerNotActiveTimeoutMillis = 10 * 1000;
|
||||
|
||||
private boolean netWorkFlowController = true;
|
||||
private boolean enableNetWorkFlowControl = false;
|
||||
|
||||
private int popPollingSize = 1024;
|
||||
private int popPollingMapSize = 100000;
|
||||
@@ -1156,12 +1156,12 @@ public class BrokerConfig extends BrokerIdentity {
|
||||
this.brokerNotActiveTimeoutMillis = brokerNotActiveTimeoutMillis;
|
||||
}
|
||||
|
||||
public boolean isNetWorkFlowController() {
|
||||
return netWorkFlowController;
|
||||
public boolean isEnableNetWorkFlowControl() {
|
||||
return enableNetWorkFlowControl;
|
||||
}
|
||||
|
||||
public void setNetWorkFlowController(boolean netWorkFlowController) {
|
||||
this.netWorkFlowController = netWorkFlowController;
|
||||
public void setEnableNetWorkFlowControl(boolean enableNetWorkFlowControl) {
|
||||
this.enableNetWorkFlowControl = enableNetWorkFlowControl;
|
||||
}
|
||||
|
||||
public boolean isRealTimeNotifyConsumerChange() {
|
||||
|
||||
Reference in New Issue
Block a user