mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-28 20:09:14 +08:00
* [ISSUE #7493] Introduce a new event NettyEventType.ACTIVE for ChannelEventListener * introduce a new event NettyEventType.ACTIVE, * implement channelActive interface for NettyRemotingClient#NettyConnectManageHandler * add onChannelActive for ChannelEventListener interface. * Move send heartbeat to onChannelActive
This commit is contained in:
@@ -87,4 +87,9 @@ public class ClientHousekeepingService implements ChannelEventListener {
|
||||
this.brokerController.getConsumerManager().doChannelCloseEvent(remoteAddr, channel);
|
||||
this.brokerController.getBrokerStatsManager().incChannelIdleNum();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -159,14 +159,6 @@ public class MQClientInstance {
|
||||
private final ConcurrentMap<String, HashMap<Long, String>> brokerAddrTable = MQClientInstance.this.brokerAddrTable;
|
||||
@Override
|
||||
public void onChannelConnect(String remoteAddr, Channel channel) {
|
||||
for (Map.Entry<String, HashMap<Long, String>> addressEntry : brokerAddrTable.entrySet()) {
|
||||
for (String address : addressEntry.getValue().values()) {
|
||||
if (address.equals(remoteAddr)) {
|
||||
sendHeartbeatToAllBrokerWithLockV2(false);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -180,6 +172,18 @@ public class MQClientInstance {
|
||||
@Override
|
||||
public void onChannelIdle(String remoteAddr, Channel channel) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
for (Map.Entry<String, HashMap<Long, String>> addressEntry : brokerAddrTable.entrySet()) {
|
||||
for (String address : addressEntry.getValue().values()) {
|
||||
if (address.equals(remoteAddr)) {
|
||||
sendHeartbeatToAllBrokerWithLockV2(false);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
} else {
|
||||
channelEventListener = null;
|
||||
|
||||
+10
-1
@@ -49,6 +49,11 @@ public class ContainerClientHouseKeepingService implements ChannelEventListener
|
||||
onChannelOperation(CallbackCode.IDLE, remoteAddr, channel);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
onChannelOperation(CallbackCode.ACTIVE, remoteAddr, channel);
|
||||
}
|
||||
|
||||
private void onChannelOperation(CallbackCode callbackCode, String remoteAddr, Channel channel) {
|
||||
Collection<InnerBrokerController> masterBrokers = this.brokerContainer.getMasterBrokers();
|
||||
Collection<InnerSalveBrokerController> slaveBrokers = this.brokerContainer.getSlaveBrokers();
|
||||
@@ -103,6 +108,10 @@ public class ContainerClientHouseKeepingService implements ChannelEventListener
|
||||
/**
|
||||
* onChannelIdle
|
||||
*/
|
||||
IDLE
|
||||
IDLE,
|
||||
/**
|
||||
* onChannelActive
|
||||
*/
|
||||
ACTIVE
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,4 +48,9 @@ public class BrokerHousekeepingService implements ChannelEventListener {
|
||||
public void onChannelIdle(String remoteAddr, Channel channel) {
|
||||
this.controllerManager.getHeartbeatManager().onBrokerChannelClose(channel);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
+5
@@ -46,4 +46,9 @@ public class BrokerHousekeepingService implements ChannelEventListener {
|
||||
public void onChannelIdle(String remoteAddr, Channel channel) {
|
||||
this.namesrvController.getRouteInfoManager().onChannelDestroy(channel);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,5 +49,9 @@ public class ClientHousekeepingService implements ChannelEventListener {
|
||||
this.clientManagerActivity.doChannelCloseEvent(remoteAddr, channel);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onChannelActive(String remoteAddr, Channel channel) {
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -26,4 +26,6 @@ public interface ChannelEventListener {
|
||||
void onChannelException(final String remoteAddr, final Channel channel);
|
||||
|
||||
void onChannelIdle(final String remoteAddr, final Channel channel);
|
||||
|
||||
void onChannelActive(final String remoteAddr, final Channel channel);
|
||||
}
|
||||
|
||||
@@ -20,5 +20,6 @@ public enum NettyEventType {
|
||||
CONNECT,
|
||||
CLOSE,
|
||||
IDLE,
|
||||
EXCEPTION
|
||||
EXCEPTION,
|
||||
ACTIVE
|
||||
}
|
||||
|
||||
@@ -701,6 +701,9 @@ public abstract class NettyRemotingAbstract {
|
||||
case EXCEPTION:
|
||||
listener.onChannelException(event.getRemoteAddr(), event.getChannel());
|
||||
break;
|
||||
case ACTIVE:
|
||||
listener.onChannelActive(event.getRemoteAddr(), event.getChannel());
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
|
||||
|
||||
@@ -1106,6 +1106,17 @@ public class NettyRemotingClient extends NettyRemotingAbstract implements Remoti
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void channelActive(ChannelHandlerContext ctx) throws Exception {
|
||||
final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
|
||||
LOGGER.info("NETTY CLIENT PIPELINE: ACTIVE, {}", remoteAddress);
|
||||
super.channelActive(ctx);
|
||||
|
||||
if (NettyRemotingClient.this.channelEventListener != null) {
|
||||
NettyRemotingClient.this.putNettyEvent(new NettyEvent(NettyEventType.ACTIVE, remoteAddress, ctx.channel()));
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void disconnect(ChannelHandlerContext ctx, ChannelPromise promise) throws Exception {
|
||||
final String remoteAddress = RemotingHelper.parseChannelRemoteAddr(ctx.channel());
|
||||
|
||||
Reference in New Issue
Block a user