diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/transport/NettyRemotingClientAbstract.java b/remoting/src/main/java/org/apache/rocketmq/remoting/transport/NettyRemotingClientAbstract.java index 5a0d9aed61dfc464fe53fc19a1728b2b300c4fae..306253b4cc0133182caa197f83fe6b80b6310b46 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/transport/NettyRemotingClientAbstract.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/transport/NettyRemotingClientAbstract.java @@ -41,6 +41,7 @@ import org.apache.rocketmq.logging.InternalLoggerFactory; import org.apache.rocketmq.remoting.RemotingChannel; import org.apache.rocketmq.remoting.common.RemotingHelper; import org.apache.rocketmq.remoting.common.RemotingUtil; +import org.apache.rocketmq.remoting.netty.NettyChannelHandlerContextImpl; import org.apache.rocketmq.remoting.netty.NettyChannelImpl; import org.apache.rocketmq.remoting.netty.NettyEvent; import org.apache.rocketmq.remoting.netty.NettyEventType; @@ -79,9 +80,9 @@ public abstract class NettyRemotingClientAbstract extends NettyRemotingAbstract if (remotingChannel instanceof NettyChannelImpl) { channel = ((NettyChannelImpl) remotingChannel).getChannel(); } -// if (remotingChannel instanceof NettyChannelHandlerContextImpl) { -// channel = ((NettyChannelHandlerContextImpl) remotingChannel).getChannelHandlerContext().channel(); -// } + if (remotingChannel instanceof NettyChannelHandlerContextImpl) { + channel = ((NettyChannelHandlerContextImpl) remotingChannel).getChannelHandlerContext().channel(); + } closeChannel(addr, channel); } diff --git a/snode/src/main/java/org/apache/rocketmq/snode/SnodeController.java b/snode/src/main/java/org/apache/rocketmq/snode/SnodeController.java index d3a7c0c9353c79ae58582702d994a0b151eda7e3..0f424a83a602f8ad3e80cad1c07db0c161f45e96 100644 --- a/snode/src/main/java/org/apache/rocketmq/snode/SnodeController.java +++ b/snode/src/main/java/org/apache/rocketmq/snode/SnodeController.java @@ -167,10 +167,6 @@ public class SnodeController { log.info("Set user specified name server address: {}", this.snodeConfig.getNamesrvAddr()); } -// this.producerManager = new ProducerManager(); - -// ConsumerIdsChangeListener consumerIdsChangeListener = new DefaultConsumerIdsChangeListener(this); -// this.consumerManager = new ConsumerManager(consumerIdsChangeListener); this.subscriptionGroupManager = new SubscriptionGroupManager(this); this.consumerOffsetManager = new ConsumerOffsetManager(this); this.consumerManageProcessor = new ConsumerManageProcessor(this);