[ISSUE #3949] v2 client manager

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:18 +08:00
committed by zhouxiang
parent c7b81c13c9
commit 35f06a4fd2
14 changed files with 366 additions and 86 deletions
@@ -29,5 +29,9 @@ public enum ConsumerGroupEvent {
/**
* The group of consumer is registered.
*/
REGISTER
REGISTER,
/**
* The client of this consumer is unregistered.
*/
CLIENT_UNREGISTER
}
@@ -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;
}
/**
@@ -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<String, ConsumerGroupInfo> consumerTable =
new ConcurrentHashMap<String, ConsumerGroupInfo>(1024);
private final ConsumerIdsChangeListener consumerIdsChangeListener;
private final List<ConsumerIdsChangeListener> 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<String, ConsumerGroupInfo> 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);
}
}
}
}
@@ -90,6 +90,8 @@ public class DefaultConsumerIdsChangeListener implements ConsumerIdsChangeListen
Collection<SubscriptionData> subscriptionDataList = (Collection<SubscriptionData>) args[0];
this.brokerController.getConsumerFilterManager().register(group, subscriptionDataList);
break;
case CLIENT_UNREGISTER:
break;
default:
throw new RuntimeException("Unknown event " + event);
}
@@ -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);
}
@@ -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
}
@@ -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<String, Channel> clientChannelTable = new ConcurrentHashMap<>();
protected final BrokerStatsManager brokerStatsManager;
private PositiveAtomicCounter positiveAtomicCounter = new PositiveAtomicCounter();
private volatile ProducerGroupOfflineListener producerGroupOfflineListener;
private final List<ProducerChangeListener> 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<Channel, ClientChannelInfo> 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);
}
}
@@ -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<ConsumerGroupEvent, List<ConsumerIdsChangeListenerData>> 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);
}
}
@@ -52,7 +52,19 @@ public class ProducerManagerTest {
public void scanNotActiveChannel() throws Exception {
producerManager.registerProducer(group, clientInfo);
AtomicReference<String> groupRef = new AtomicReference<>();
producerManager.setProducerOfflineListener(groupRef::set);
AtomicReference<ClientChannelInfo> 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<String> groupRef = new AtomicReference<>();
producerManager.setProducerOfflineListener(groupRef::set);
AtomicReference<ClientChannelInfo> 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<String> groupRef = new AtomicReference<>();
producerManager.setProducerOfflineListener(groupRef::set);
AtomicReference<ClientChannelInfo> 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<Channel, ClientChannelInfo> 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();
@@ -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<String, SimpleChannel> clientIdChannelMap = new ConcurrentHashMap<>();
private final ConcurrentMap<String /* clientId */, SimpleChannel> clientIdChannelMap = new ConcurrentHashMap<>();
private final ConcurrentMap<String /* group */, Set<String>/* 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<Map.Entry<String, SimpleChannel>> iterator = clientIdChannelMap.entrySet().iterator();
while (iterator.hasNext()) {
Map.Entry<String, SimpleChannel> 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;
});
}
}
}
@@ -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);
}
}
@@ -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());
}
}
}
}
@@ -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<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
@@ -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;