From 35f06a4fd24ece1a5a13c4bc0dccb04c6a838796 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Fri, 15 Apr 2022 16:58:59 +0800 Subject: [PATCH] [ISSUE #3949] v2 client manager --- .../broker/client/ConsumerGroupEvent.java | 6 +- .../broker/client/ConsumerGroupInfo.java | 9 +- .../broker/client/ConsumerManager.java | 48 +++++-- .../DefaultConsumerIdsChangeListener.java | 2 + ...tener.java => ProducerChangeListener.java} | 4 +- .../broker/client/ProducerGroupEvent.java | 28 ++++ .../broker/client/ProducerManager.java | 27 ++-- .../broker/client/ConsumerManagerTest.java | 130 ++++++++++++++++++ .../broker/client/ProducerManagerTest.java | 45 +++++- .../proxy/channel/ChannelManager.java | 48 ++----- .../grpc/v2/service/GrpcClientManager.java | 4 + .../grpc/v2/service/LocalGrpcService.java | 44 ++++-- .../service/cluster/ForwardClientService.java | 56 ++++++-- .../grpc/v2/service/cluster/RouteService.java | 1 - 14 files changed, 366 insertions(+), 86 deletions(-) rename broker/src/main/java/org/apache/rocketmq/broker/client/{ProducerGroupOfflineListener.java => ProducerChangeListener.java} (86%) create mode 100644 broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupEvent.java create mode 100644 broker/src/test/java/org/apache/rocketmq/broker/client/ConsumerManagerTest.java diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupEvent.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupEvent.java index 717fb7085e..2318edb5f3 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupEvent.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupEvent.java @@ -29,5 +29,9 @@ public enum ConsumerGroupEvent { /** * The group of consumer is registered. */ - REGISTER + REGISTER, + /** + * The client of this consumer is unregistered. + */ + CLIENT_UNREGISTER } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupInfo.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupInfo.java index 09e1241518..638c522feb 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupInfo.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerGroupInfo.java @@ -98,23 +98,24 @@ public class ConsumerGroupInfo { return result; } - public void unregisterChannel(final ClientChannelInfo clientChannelInfo) { + public boolean unregisterChannel(final ClientChannelInfo clientChannelInfo) { ClientChannelInfo old = this.channelInfoTable.remove(clientChannelInfo.getChannel()); if (old != null) { log.info("unregister a consumer[{}] from consumerGroupInfo {}", this.groupName, old.toString()); + return true; } + return false; } - public boolean doChannelCloseEvent(final String remoteAddr, final Channel channel) { + public ClientChannelInfo doChannelCloseEvent(final String remoteAddr, final Channel channel) { final ClientChannelInfo info = this.channelInfoTable.remove(channel); if (info != null) { log.warn( "NETTY EVENT: remove not active channel[{}] from ConsumerGroupInfo groupChannelTable, consumer group: {}", info.toString(), groupName); - return true; } - return false; + return info; } /** diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerManager.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerManager.java index b3bee7cdc0..2f0f9a6789 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ConsumerManager.java @@ -18,12 +18,14 @@ package org.apache.rocketmq.broker.client; import java.util.HashSet; import java.util.Iterator; +import java.util.List; import java.util.Map.Entry; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import io.netty.channel.Channel; +import java.util.concurrent.CopyOnWriteArrayList; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; @@ -40,16 +42,17 @@ public class ConsumerManager { private static final long CHANNEL_EXPIRED_TIMEOUT = 1000 * 120; private final ConcurrentMap consumerTable = new ConcurrentHashMap(1024); - private final ConsumerIdsChangeListener consumerIdsChangeListener; + private final List consumerIdsChangeListenerList = new CopyOnWriteArrayList<>(); protected final BrokerStatsManager brokerStatsManager; public ConsumerManager(final ConsumerIdsChangeListener consumerIdsChangeListener) { - this.consumerIdsChangeListener = consumerIdsChangeListener; + this.consumerIdsChangeListenerList.add(consumerIdsChangeListener); this.brokerStatsManager = null; } - public ConsumerManager(final ConsumerIdsChangeListener consumerIdsChangeListener, final BrokerStatsManager brokerStatsManager) { - this.consumerIdsChangeListener = consumerIdsChangeListener; + public ConsumerManager(final ConsumerIdsChangeListener consumerIdsChangeListener, + final BrokerStatsManager brokerStatsManager) { + this.consumerIdsChangeListenerList.add(consumerIdsChangeListener); this.brokerStatsManager = brokerStatsManager; } @@ -93,18 +96,19 @@ public class ConsumerManager { while (it.hasNext()) { Entry next = it.next(); ConsumerGroupInfo info = next.getValue(); - removed = info.doChannelCloseEvent(remoteAddr, channel); - if (removed) { + ClientChannelInfo clientChannelInfo = info.doChannelCloseEvent(remoteAddr, channel); + if (clientChannelInfo != null) { + callConsumerIdsChangeListener(ConsumerGroupEvent.CLIENT_UNREGISTER, next.getKey(), clientChannelInfo); if (info.getChannelInfoTable().isEmpty()) { ConsumerGroupInfo remove = this.consumerTable.remove(next.getKey()); if (remove != null) { LOGGER.info("unregister consumer ok, no any connection, and remove consumer group, {}", next.getKey()); - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.UNREGISTER, next.getKey()); + callConsumerIdsChangeListener(ConsumerGroupEvent.UNREGISTER, next.getKey()); } } - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.CHANGE, next.getKey(), info.getAllChannel()); + callConsumerIdsChangeListener(ConsumerGroupEvent.CHANGE, next.getKey(), info.getAllChannel()); } } return removed; @@ -128,14 +132,14 @@ public class ConsumerManager { if (r1 || r2) { if (isNotifyConsumerIdsChangedEnable) { - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.CHANGE, group, consumerGroupInfo.getAllChannel()); + callConsumerIdsChangeListener(ConsumerGroupEvent.CHANGE, group, consumerGroupInfo.getAllChannel()); } } if (null != this.brokerStatsManager) { this.brokerStatsManager.incConsumerRegisterTime((int) (System.currentTimeMillis() - start)); } - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.REGISTER, group, subList); + callConsumerIdsChangeListener(ConsumerGroupEvent.REGISTER, group, subList); return r1 || r2; } @@ -144,17 +148,20 @@ public class ConsumerManager { boolean isNotifyConsumerIdsChangedEnable) { ConsumerGroupInfo consumerGroupInfo = this.consumerTable.get(group); if (null != consumerGroupInfo) { - consumerGroupInfo.unregisterChannel(clientChannelInfo); + boolean removed = consumerGroupInfo.unregisterChannel(clientChannelInfo); + if (removed) { + callConsumerIdsChangeListener(ConsumerGroupEvent.CLIENT_UNREGISTER, group, clientChannelInfo); + } if (consumerGroupInfo.getChannelInfoTable().isEmpty()) { ConsumerGroupInfo remove = this.consumerTable.remove(group); if (remove != null) { LOGGER.info("unregister consumer ok, no any connection, and remove consumer group, {}", group); - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.UNREGISTER, group); + callConsumerIdsChangeListener(ConsumerGroupEvent.UNREGISTER, group); } } if (isNotifyConsumerIdsChangedEnable) { - this.consumerIdsChangeListener.handle(ConsumerGroupEvent.CHANGE, group, consumerGroupInfo.getAllChannel()); + callConsumerIdsChangeListener(ConsumerGroupEvent.CHANGE, group, consumerGroupInfo.getAllChannel()); } } } @@ -177,6 +184,7 @@ public class ConsumerManager { LOGGER.warn( "SCAN: remove expired channel from ConsumerManager consumerTable. channel={}, consumerGroup={}", RemotingHelper.parseChannelRemoteAddr(clientChannelInfo.getChannel()), group); + callConsumerIdsChangeListener(ConsumerGroupEvent.CLIENT_UNREGISTER, group, clientChannelInfo); RemotingUtil.closeChannel(clientChannelInfo.getChannel()); itChannel.remove(); } @@ -204,4 +212,18 @@ public class ConsumerManager { } return groups; } + + public void appendConsumerIdsChangeListener(ConsumerIdsChangeListener listener) { + consumerIdsChangeListenerList.add(listener); + } + + protected void callConsumerIdsChangeListener(ConsumerGroupEvent event, String group, Object... args) { + for (ConsumerIdsChangeListener listener : consumerIdsChangeListenerList) { + try { + listener.handle(event, group, args); + } catch (Throwable t) { + LOGGER.error("err when call consumerIdsChangeListener", t); + } + } + } } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java b/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java index 8e6e667dd3..72ccc8f16f 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/DefaultConsumerIdsChangeListener.java @@ -90,6 +90,8 @@ public class DefaultConsumerIdsChangeListener implements ConsumerIdsChangeListen Collection subscriptionDataList = (Collection) args[0]; this.brokerController.getConsumerFilterManager().register(group, subscriptionDataList); break; + case CLIENT_UNREGISTER: + break; default: throw new RuntimeException("Unknown event " + event); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupOfflineListener.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerChangeListener.java similarity index 86% rename from broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupOfflineListener.java rename to broker/src/main/java/org/apache/rocketmq/broker/client/ProducerChangeListener.java index 1107f0d9ed..576faf8a81 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupOfflineListener.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerChangeListener.java @@ -16,7 +16,7 @@ */ package org.apache.rocketmq.broker.client; -public interface ProducerGroupOfflineListener { +public interface ProducerChangeListener { - void onOffline(String group); + void handle(ProducerGroupEvent event, String group, ClientChannelInfo clientChannelInfo); } diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupEvent.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupEvent.java new file mode 100644 index 0000000000..cbf27ce61e --- /dev/null +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerGroupEvent.java @@ -0,0 +1,28 @@ +/* + * 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.broker.client; + +public enum ProducerGroupEvent { + /** + * The group of producer is unregistered. + */ + GROUP_UNREGISTER, + /** + * The client of this producer is unregistered. + */ + CLIENT_UNREGISTER +} diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerManager.java b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerManager.java index 09da153a96..41df13f8f5 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/ProducerManager.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; import org.apache.rocketmq.broker.util.PositiveAtomicCounter; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.logging.InternalLogger; @@ -40,7 +41,7 @@ public class ProducerManager { private final ConcurrentHashMap clientChannelTable = new ConcurrentHashMap<>(); protected final BrokerStatsManager brokerStatsManager; private PositiveAtomicCounter positiveAtomicCounter = new PositiveAtomicCounter(); - private volatile ProducerGroupOfflineListener producerGroupOfflineListener; + private final List producerChangeListenerList = new CopyOnWriteArrayList<>(); public ProducerManager() { this.brokerStatsManager = null; @@ -85,6 +86,7 @@ public class ProducerManager { log.warn( "ProducerManager#scanNotActiveChannel: remove expired channel[{}] from ProducerManager groupChannelTable, producer group name: {}", RemotingHelper.parseChannelRemoteAddr(info.getChannel()), group); + callProducerChangeListener(ProducerGroupEvent.CLIENT_UNREGISTER, group, info); RemotingUtil.closeChannel(info.getChannel()); } } @@ -92,7 +94,7 @@ public class ProducerManager { if (chlMap.isEmpty()) { log.warn("SCAN: remove expired channel from ProducerManager groupChannelTable, all clear, group={}", group); iterator.remove(); - this.notifyProducerOffline(group); + callProducerChangeListener(ProducerGroupEvent.GROUP_UNREGISTER, group, null); } } } @@ -113,11 +115,12 @@ public class ProducerManager { log.info( "NETTY EVENT: remove channel[{}][{}] from ProducerManager groupChannelTable, producer group: {}", clientChannelInfo.toString(), remoteAddr, group); + callProducerChangeListener(ProducerGroupEvent.CLIENT_UNREGISTER, group, clientChannelInfo); if (clientChannelInfoTable.isEmpty()) { ConcurrentHashMap oldGroupTable = this.groupChannelTable.remove(group); if (oldGroupTable != null) { log.info("unregister a producer group[{}] from groupChannelTable", group); - this.notifyProducerOffline(group); + callProducerChangeListener(ProducerGroupEvent.GROUP_UNREGISTER, group, null); } } } @@ -158,11 +161,12 @@ public class ProducerManager { if (old != null) { log.info("unregister a producer[{}] from groupChannelTable {}", group, clientChannelInfo.toString()); + callProducerChangeListener(ProducerGroupEvent.CLIENT_UNREGISTER, group, clientChannelInfo); } if (channelTable.isEmpty()) { this.groupChannelTable.remove(group); - this.notifyProducerOffline(group); + callProducerChangeListener(ProducerGroupEvent.GROUP_UNREGISTER, group, null); log.info("unregister a producer group[{}] from groupChannelTable", group); } } @@ -212,13 +216,18 @@ public class ProducerManager { return clientChannelTable.get(clientId); } - public void notifyProducerOffline(String group) { - if (this.producerGroupOfflineListener != null) { - this.producerGroupOfflineListener.onOffline(group); + private void callProducerChangeListener(ProducerGroupEvent event, String group, + ClientChannelInfo clientChannelInfo) { + for (ProducerChangeListener listener : producerChangeListenerList) { + try { + listener.handle(event, group, clientChannelInfo); + } catch (Throwable t) { + log.error("err when call producerChangeListener", t); + } } } - public void setProducerOfflineListener(ProducerGroupOfflineListener producerGroupOfflineListener) { - this.producerGroupOfflineListener = producerGroupOfflineListener; + public void appendProducerChangeListener(ProducerChangeListener producerChangeListener) { + producerChangeListenerList.add(producerChangeListener); } } diff --git a/broker/src/test/java/org/apache/rocketmq/broker/client/ConsumerManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/client/ConsumerManagerTest.java new file mode 100644 index 0000000000..f6dbd3e973 --- /dev/null +++ b/broker/src/test/java/org/apache/rocketmq/broker/client/ConsumerManagerTest.java @@ -0,0 +1,130 @@ +package org.apache.rocketmq.broker.client; + +import io.netty.channel.Channel; +import io.netty.channel.ChannelFuture; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import org.apache.rocketmq.common.consumer.ConsumeFromWhere; +import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; +import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; +import org.apache.rocketmq.remoting.protocol.LanguageCode; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnitRunner; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.class) +public class ConsumerManagerTest { + private ConsumerManager consumerManager; + private String group = "FooBar"; + private String clientId = "clientId"; + private ClientChannelInfo clientInfo; + private Map> groupEventListMap = new HashMap<>(); + + @Mock + private Channel channel; + + @Before + public void init() { + clientInfo = new ClientChannelInfo(channel, clientId, LanguageCode.JAVA, 0); + + consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() { + @Override + public void handle(ConsumerGroupEvent event, String group, Object... args) { + groupEventListMap.compute(event, (eventKey, dataListVal) -> { + if (dataListVal == null) { + dataListVal = new ArrayList<>(); + } + dataListVal.add(new ConsumerIdsChangeListenerData(event, group, args)); + return dataListVal; + }); + } + + @Override + public void shutdown() { + + } + }); + } + + private static class ConsumerIdsChangeListenerData { + private ConsumerGroupEvent event; + private String group; + private Object[] args; + + public ConsumerIdsChangeListenerData(ConsumerGroupEvent event, String group, Object[] args) { + this.event = event; + this.group = group; + this.args = args; + } + } + + @Test + public void testClientUnregisterEventInDoChannelCloseEvent() { + assertThat(consumerManager.registerConsumer( + group, + clientInfo, + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + new HashSet<>(), + false + )).isTrue(); + + consumerManager.doChannelCloseEvent("remoteAddr", channel); + + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).size()).isEqualTo(1); + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]).isInstanceOf(ClientChannelInfo.class); + ClientChannelInfo clientChannelInfo = (ClientChannelInfo) groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]; + assertThat(clientChannelInfo).isSameAs(clientInfo); + } + + @Test + public void testClientUnregisterEventInUnregisterConsumer() { + assertThat(consumerManager.registerConsumer( + group, + clientInfo, + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + new HashSet<>(), + false + )).isTrue(); + + consumerManager.unregisterConsumer(group, clientInfo, false); + + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).size()).isEqualTo(1); + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]).isInstanceOf(ClientChannelInfo.class); + ClientChannelInfo clientChannelInfo = (ClientChannelInfo) groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]; + assertThat(clientChannelInfo).isSameAs(clientInfo); + } + + @Test + public void testClientUnregisterEventInScanNotActiveChannel() { + assertThat(consumerManager.registerConsumer( + group, + clientInfo, + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + new HashSet<>(), + false + )).isTrue(); + clientInfo.setLastUpdateTimestamp(0); + when(channel.close()).thenReturn(mock(ChannelFuture.class)); + + consumerManager.scanNotActiveChannel(); + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).size()).isEqualTo(1); + assertThat(groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]).isInstanceOf(ClientChannelInfo.class); + ClientChannelInfo clientChannelInfo = (ClientChannelInfo) groupEventListMap.get(ConsumerGroupEvent.CLIENT_UNREGISTER).get(0).args[0]; + assertThat(clientChannelInfo).isSameAs(clientInfo); + } +} \ No newline at end of file diff --git a/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java index 3d05d39ef8..fd76312941 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/client/ProducerManagerTest.java @@ -52,7 +52,19 @@ public class ProducerManagerTest { public void scanNotActiveChannel() throws Exception { producerManager.registerProducer(group, clientInfo); AtomicReference groupRef = new AtomicReference<>(); - producerManager.setProducerOfflineListener(groupRef::set); + AtomicReference clientChannelInfoRef = new AtomicReference<>(); + producerManager.appendProducerChangeListener((event, group, clientChannelInfo) -> { + switch (event) { + case GROUP_UNREGISTER: + groupRef.set(group); + break; + case CLIENT_UNREGISTER: + clientChannelInfoRef.set(clientChannelInfo); + break; + default: + break; + } + }); assertThat(producerManager.getGroupChannelTable().get(group).get(channel)).isNotNull(); assertThat(producerManager.findChannel("clientId")).isNotNull(); Field field = ProducerManager.class.getDeclaredField("CHANNEL_EXPIRED_TIMEOUT"); @@ -63,6 +75,7 @@ public class ProducerManagerTest { producerManager.scanNotActiveChannel(); assertThat(producerManager.getGroupChannelTable().get(group)).isNull(); assertThat(groupRef.get()).isEqualTo(group); + assertThat(clientChannelInfoRef.get()).isSameAs(clientInfo); assertThat(producerManager.findChannel("clientId")).isNull(); } @@ -70,12 +83,25 @@ public class ProducerManagerTest { public void doChannelCloseEvent() throws Exception { producerManager.registerProducer(group, clientInfo); AtomicReference groupRef = new AtomicReference<>(); - producerManager.setProducerOfflineListener(groupRef::set); + AtomicReference clientChannelInfoRef = new AtomicReference<>(); + producerManager.appendProducerChangeListener((event, group, clientChannelInfo) -> { + switch (event) { + case GROUP_UNREGISTER: + groupRef.set(group); + break; + case CLIENT_UNREGISTER: + clientChannelInfoRef.set(clientChannelInfo); + break; + default: + break; + } + }); assertThat(producerManager.getGroupChannelTable().get(group).get(channel)).isNotNull(); assertThat(producerManager.findChannel("clientId")).isNotNull(); producerManager.doChannelCloseEvent("127.0.0.1", channel); assertThat(producerManager.getGroupChannelTable().get(group)).isNull(); assertThat(groupRef.get()).isEqualTo(group); + assertThat(clientChannelInfoRef.get()).isSameAs(clientInfo); assertThat(producerManager.findChannel("clientId")).isNull(); } @@ -94,7 +120,19 @@ public class ProducerManagerTest { public void unregisterProducer() throws Exception { producerManager.registerProducer(group, clientInfo); AtomicReference groupRef = new AtomicReference<>(); - producerManager.setProducerOfflineListener(groupRef::set); + AtomicReference clientChannelInfoRef = new AtomicReference<>(); + producerManager.appendProducerChangeListener((event, group, clientChannelInfo) -> { + switch (event) { + case GROUP_UNREGISTER: + groupRef.set(group); + break; + case CLIENT_UNREGISTER: + clientChannelInfoRef.set(clientChannelInfo); + break; + default: + break; + } + }); Map channelMap = producerManager.getGroupChannelTable().get(group); assertThat(channelMap).isNotNull(); assertThat(channelMap.get(channel)).isEqualTo(clientInfo); @@ -105,6 +143,7 @@ public class ProducerManagerTest { channelMap = producerManager.getGroupChannelTable().get(group); channel1 = producerManager.findChannel("clientId"); assertThat(groupRef.get()).isEqualTo(group); + assertThat(clientChannelInfoRef.get()).isSameAs(clientInfo); assertThat(channelMap).isNull(); assertThat(channel1).isNull(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java index 5adc5e0b80..aae433f54c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java @@ -20,25 +20,22 @@ package org.apache.rocketmq.proxy.channel; import io.grpc.Context; import java.util.ArrayList; import java.util.Collections; -import java.util.Iterator; import java.util.List; -import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.function.Supplier; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.constant.LoggerName; -import org.apache.rocketmq.proxy.common.Cleaner; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.grpc.v1.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; +import org.apache.rocketmq.proxy.grpc.v1.adapter.channel.GrpcClientChannel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class ChannelManager { private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); - private final ConcurrentMap clientIdChannelMap = new ConcurrentHashMap<>(); + private final ConcurrentMap clientIdChannelMap = new ConcurrentHashMap<>(); private final ConcurrentMap/* clientId */> groupClientIdMap = new ConcurrentHashMap<>(); public SimpleChannel createChannel() { @@ -123,35 +120,20 @@ public class ChannelManager { return new ArrayList<>(groupClientIdMap.get(group)); } - /** - * Scan and remove inactive mocking channels; Scan and clean expired requests; - */ - public void scanAndCleanChannels() { - try { - Iterator> iterator = clientIdChannelMap.entrySet().iterator(); - while (iterator.hasNext()) { - Map.Entry entry = iterator.next(); - if (!entry.getValue().isActive()) { - iterator.remove(); - if (entry.getValue() instanceof GrpcClientChannel) { - GrpcClientChannel grpcClientChannel = (GrpcClientChannel) entry.getValue(); - groupClientIdMap.computeIfPresent(grpcClientChannel.getGroup(), (group, clientIds) -> { - clientIds.remove(grpcClientChannel.getClientId()); - if (clientIds.isEmpty()) { - return null; - } - return clientIds; - }); - } - } else { - if (entry.getValue() instanceof Cleaner) { - Cleaner cleaner = (Cleaner) entry.getValue(); - cleaner.clean(); - } + public void onClientOffline(String clientId) { + SimpleChannel simpleChannel = clientIdChannelMap.remove(clientId); + if (simpleChannel == null) { + return; + } + if (simpleChannel instanceof GrpcClientChannel) { + GrpcClientChannel grpcClientChannel = (GrpcClientChannel) simpleChannel; + groupClientIdMap.computeIfPresent(grpcClientChannel.getGroup(), (group, clientIds) -> { + clientIds.remove(grpcClientChannel.getClientId()); + if (clientIds.isEmpty()) { + return null; } - } - } catch (Throwable e) { - log.error("Unexpected exception", e); + return clientIds; + }); } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java index e3705c8532..f949af8ecf 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java @@ -39,4 +39,8 @@ public class GrpcClientManager { public void updateClientSettings(String clientId, ClientSettings clientSettings) { CLIENT_SETTINGS_MAP.put(clientId, clientSettings); } + + public ClientSettings removeClientSettings(String clientId) { + return CLIENT_SETTINGS_MAP.remove(clientId); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java index 72a1701116..73998b4283 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java @@ -65,6 +65,11 @@ import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.broker.client.ClientChannelInfo; +import org.apache.rocketmq.broker.client.ConsumerGroupEvent; +import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; +import org.apache.rocketmq.broker.client.ProducerChangeListener; +import org.apache.rocketmq.broker.client.ProducerGroupEvent; import org.apache.rocketmq.common.MQVersion; import org.apache.rocketmq.common.ThreadFactoryImpl; import org.apache.rocketmq.common.constant.LoggerName; @@ -91,7 +96,6 @@ import org.apache.rocketmq.proxy.channel.SimpleChannel; import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.common.DelayPolicy; -import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.common.TelemetryCommandManager; import org.apache.rocketmq.proxy.common.TelemetryCommandRecord; import org.apache.rocketmq.proxy.connector.ConnectorManager; @@ -144,8 +148,11 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo this.grpcClientManager = new GrpcClientManager(); this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager); this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel()); + + this.brokerController.getConsumerManager().appendConsumerIdsChangeListener(new ConsumerIdsChangeListenerImpl()); + this.brokerController.getProducerManager().appendProducerChangeListener(new ProducerChangeListenerImpl()); + this.appendStartAndShutdown(connectorManager); - this.appendStartAndShutdown(new LocalGrpcServiceStartAndShutdown()); } @Override @@ -615,17 +622,36 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo }; } - private class LocalGrpcServiceStartAndShutdown implements StartAndShutdown { - @Override public void start() throws Exception { - LocalGrpcService.this.scheduledExecutorService.scheduleWithFixedDelay(LocalGrpcService.this::scanAndCleanChannels, 5, 5, TimeUnit.MINUTES); + protected class ConsumerIdsChangeListenerImpl implements ConsumerIdsChangeListener { + + @Override + public void handle(ConsumerGroupEvent event, String group, Object... args) { + if (event == ConsumerGroupEvent.CLIENT_UNREGISTER) { + if (args == null || args.length < 1) { + return; + } + if (args[0] instanceof ClientChannelInfo) { + ClientChannelInfo clientChannelInfo = (ClientChannelInfo) args[0]; + channelManager.onClientOffline(clientChannelInfo.getClientId()); + grpcClientManager.removeClientSettings(clientChannelInfo.getClientId()); + } + } } - @Override public void shutdown() throws Exception { - LocalGrpcService.this.scheduledExecutorService.shutdown(); + @Override + public void shutdown() { + } } - private void scanAndCleanChannels() { - this.channelManager.scanAndCleanChannels(); + protected class ProducerChangeListenerImpl implements ProducerChangeListener { + + @Override + public void handle(ProducerGroupEvent event, String group, ClientChannelInfo clientChannelInfo) { + if (event == ProducerGroupEvent.CLIENT_UNREGISTER) { + channelManager.onClientOffline(clientChannelInfo.getClientId()); + grpcClientManager.removeClientSettings(clientChannelInfo.getClientId()); + } + } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java index 770294eefb..cea463d2f1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java @@ -39,13 +39,15 @@ import org.apache.rocketmq.broker.client.ClientChannelInfo; import org.apache.rocketmq.broker.client.ConsumerGroupEvent; import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.broker.client.ConsumerManager; +import org.apache.rocketmq.broker.client.ProducerChangeListener; +import org.apache.rocketmq.broker.client.ProducerGroupEvent; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.MQVersion; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.proxy.channel.ChannelManager; -import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.common.TelemetryCommandManager; +import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException; @@ -82,17 +84,49 @@ public class ForwardClientService extends BaseService { this.grpcClientManager = grpcClientManager; this.telemetryCommandManager = telemetryCommandManager; - this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() { - @Override - public void handle(ConsumerGroupEvent event, String group, Object... args) { - } - - @Override - public void shutdown() { - } - }); + this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListenerImpl()); this.producerManager = new ProducerManager(); - this.producerManager.setProducerOfflineListener(connectorManager.getTransactionHeartbeatRegisterService()::onProducerGroupOffline); + this.producerManager.appendProducerChangeListener(new ProducerChangeListenerImpl()); + } + + protected class ConsumerIdsChangeListenerImpl implements ConsumerIdsChangeListener { + + @Override + public void handle(ConsumerGroupEvent event, String group, Object... args) { + if (event == ConsumerGroupEvent.CLIENT_UNREGISTER) { + if (args == null || args.length < 1) { + return; + } + if (args[0] instanceof ClientChannelInfo) { + ClientChannelInfo clientChannelInfo = (ClientChannelInfo) args[0]; + channelManager.onClientOffline(clientChannelInfo.getClientId()); + grpcClientManager.removeClientSettings(clientChannelInfo.getClientId()); + } + } + } + + @Override + public void shutdown() { + + } + } + + protected class ProducerChangeListenerImpl implements ProducerChangeListener { + + @Override + public void handle(ProducerGroupEvent event, String group, ClientChannelInfo clientChannelInfo) { + switch (event) { + case GROUP_UNREGISTER: + connectorManager.getTransactionHeartbeatRegisterService().onProducerGroupOffline(group); + break; + case CLIENT_UNREGISTER: + channelManager.onClientOffline(clientChannelInfo.getClientId()); + grpcClientManager.removeClientSettings(clientChannelInfo.getClientId()); + break; + default: + break; + } + } } public CompletableFuture heartbeat(Context ctx, HeartbeatRequest request) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index 8ebea5c10a..55a62de108 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -48,7 +48,6 @@ import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper; -import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;