From f07904a70acdc99fa2d79c458de9609f85356f8e Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Wed, 16 Mar 2022 11:51:31 +0800 Subject: [PATCH] [ISSUE #3949] Simplify Channels --- .../proxy/channel/ChannelManager.java | 11 +++++----- .../proxy/channel/InvocationChannel.java | 17 ++++++++++---- .../apache/rocketmq/proxy/common/Cleaner.java | 22 +++++++++++++++++++ .../channel/ReceiveMessageChannel.java | 21 +----------------- .../adapter/channel/SendMessageChannel.java | 22 +------------------ 5 files changed, 42 insertions(+), 51 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/common/Cleaner.java 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 53bf9e33c8..61910fb248 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 @@ -26,7 +26,7 @@ import java.util.concurrent.ConcurrentMap; import java.util.function.Supplier; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.grpc.adapter.channel.SendMessageChannel; +import org.apache.rocketmq.proxy.common.Cleaner; import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -102,13 +102,12 @@ public class ChannelManager { Iterator> iterator = clientIdChannelMap.entrySet().iterator(); while (iterator.hasNext()) { Map.Entry entry = iterator.next(); - if (!entry.getValue() - .isActive()) { + if (!entry.getValue().isActive()) { iterator.remove(); } else { - if (entry.getValue() instanceof SendMessageChannel) { - SendMessageChannel channel = (SendMessageChannel) entry.getValue(); - channel.cleanExpiredRequests(); + if (entry.getValue() instanceof Cleaner) { + Cleaner cleaner = (Cleaner) entry.getValue(); + cleaner.clean(); } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java index 4c70cc74c8..23caa9e8e5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/InvocationChannel.java @@ -23,21 +23,29 @@ import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import org.apache.rocketmq.proxy.common.Cleaner; import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; +import org.apache.rocketmq.proxy.grpc.adapter.handler.ResponseHandler; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public class InvocationChannel extends SimpleChannel { +public abstract class InvocationChannel extends SimpleChannel implements Cleaner { protected final ConcurrentMap> inFlightRequestMap; + protected final ResponseHandler handler; - public InvocationChannel(SimpleChannel simpleChannel) { - super(simpleChannel); + public InvocationChannel(ResponseHandler handler) { + super(ChannelManager.createSimpleChannelDirectly()); this.inFlightRequestMap = new ConcurrentHashMap<>(); + this.handler = handler; } @Override public ChannelFuture writeAndFlush(Object msg) { if (msg instanceof RemotingCommand) { RemotingCommand responseCommand = (RemotingCommand) msg; + InvocationContext context = inFlightRequestMap.remove(responseCommand.getOpaque()); + if (null != context) { + handler.handle(responseCommand, context); + } inFlightRequestMap.remove(responseCommand.getOpaque()); } return super.writeAndFlush(msg); @@ -64,7 +72,8 @@ public class InvocationChannel extends SimpleChannel { inFlightRequestMap.remove(opaque); } - public void cleanExpiredRequests() { + @Override + public void clean() { Iterator>> iterator = inFlightRequestMap.entrySet().iterator(); int count = 0; while (iterator.hasNext()) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/Cleaner.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/Cleaner.java new file mode 100644 index 0000000000..a02b08a913 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/Cleaner.java @@ -0,0 +1,22 @@ +/* + * 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.proxy.common; + +public interface Cleaner { + void clean(); +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ReceiveMessageChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ReceiveMessageChannel.java index acc280a73f..934b049b78 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ReceiveMessageChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ReceiveMessageChannel.java @@ -19,30 +19,11 @@ package org.apache.rocketmq.proxy.grpc.adapter.channel; import apache.rocketmq.v1.ReceiveMessageRequest; import apache.rocketmq.v1.ReceiveMessageResponse; -import io.netty.channel.ChannelFuture; -import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.InvocationChannel; -import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHandler; -import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ReceiveMessageChannel extends InvocationChannel { - private final ReceiveMessageResponseHandler handler; - public ReceiveMessageChannel(ReceiveMessageResponseHandler handler) { - super(ChannelManager.createSimpleChannelDirectly()); - this.handler = handler; - } - - @Override - public ChannelFuture writeAndFlush(Object msg) { - if (msg instanceof RemotingCommand) { - RemotingCommand responseCommand = (RemotingCommand) msg; - InvocationContext context = inFlightRequestMap.remove(responseCommand.getOpaque()); - if (null != context) { - handler.handle(responseCommand, context); - } - } - return super.writeAndFlush(msg); + super(handler); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java index d0037d43fb..3eeb98f1ec 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/SendMessageChannel.java @@ -19,31 +19,11 @@ package org.apache.rocketmq.proxy.grpc.adapter.channel; import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; -import io.netty.channel.ChannelFuture; -import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.InvocationChannel; -import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler; -import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class SendMessageChannel extends InvocationChannel { - private final SendMessageResponseHandler handler; - public SendMessageChannel(SendMessageResponseHandler handler) { - super(ChannelManager.createSimpleChannelDirectly()); - this.handler = handler; - } - - @Override - public ChannelFuture writeAndFlush(Object msg) { - if (msg instanceof RemotingCommand) { - RemotingCommand responseCommand = (RemotingCommand) msg; - InvocationContext context = inFlightRequestMap.remove(responseCommand.getOpaque()); - if (null != context) { - handler.handle(responseCommand, context); - } - } - - return super.writeAndFlush(msg); + super(handler); } }