From 053c6437fbc4e6646a5507f14cd57ae831ebbdae Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 18 Mar 2022 14:15:34 +0800 Subject: [PATCH] [ISSUE #3949] Add DelayPolicy --- .../rocketmq/proxy/config/ProxyConfig.java | 19 +++++ .../rocketmq/proxy/grpc/common/Converter.java | 6 +- .../proxy/grpc/common/DelayPolicy.java | 82 +++++++++++++++++++ .../proxy/grpc/service/LocalGrpcService.java | 5 +- .../grpc/service/cluster/ConsumerService.java | 6 +- 5 files changed, 114 insertions(+), 4 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/DelayPolicy.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 9c40facf87..cb23d66c6b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -79,6 +79,9 @@ public class ProxyConfig { private int longPollingReserveTimeInMillis = 10000; + private int retryDelayLevelDelta = 3; + private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"; + public Integer getHealthCheckPort() { return healthCheckPort; } @@ -366,4 +369,20 @@ public class ProxyConfig { public void setLongPollingReserveTimeInMillis(int longPollingReserveTimeInMillis) { this.longPollingReserveTimeInMillis = longPollingReserveTimeInMillis; } + + public int getRetryDelayLevelDelta() { + return retryDelayLevelDelta; + } + + public void setRetryDelayLevelDelta(int retryDelayLevelDelta) { + this.retryDelayLevelDelta = retryDelayLevelDelta; + } + + public String getMessageDelayLevel() { + return messageDelayLevel; + } + + public void setMessageDelayLevel(String messageDelayLevel) { + this.messageDelayLevel = messageDelayLevel; + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index 69b0a27f4d..03f4877313 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -88,6 +88,7 @@ import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.common.sysflag.MessageSysFlag; import org.apache.rocketmq.common.sysflag.PullSysFlag; import org.apache.rocketmq.common.utils.BinaryUtil; +import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -179,7 +180,8 @@ public class Converter { return ackMessageRequestHeader; } - public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(NackMessageRequest request) { + public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(NackMessageRequest request, + DelayPolicy delayPolicy) { String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); @@ -191,7 +193,7 @@ public class Converter { changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId()); changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle()); changeInvisibleTimeRequestHeader.setOffset(handle.getOffset()); - changeInvisibleTimeRequestHeader.setInvisibleTime(0L); + changeInvisibleTimeRequestHeader.setInvisibleTime(delayPolicy.getDelayInterval(ConfigurationManager.getProxyConfig().getRetryDelayLevelDelta() + request.getDeliveryAttempt())); return changeInvisibleTimeRequestHeader; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/DelayPolicy.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/DelayPolicy.java new file mode 100644 index 0000000000..18d442e7a7 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/DelayPolicy.java @@ -0,0 +1,82 @@ +/* + * 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.grpc.common; + +import com.google.common.base.Splitter; +import com.google.common.collect.Lists; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +public class DelayPolicy { + private List delayIntervalList; + + private DelayPolicy(List delayIntervalList) { + this.delayIntervalList = delayIntervalList; + } + + public long getDelayInterval(int index) { + int size = delayIntervalList.size(); + if (index >= size) { + throw new IllegalArgumentException("Out of index, size: " + size); + } + return delayIntervalList.get(index); + } + + public void refresh(String messageDelayLevel) { + delayIntervalList = buildList(messageDelayLevel); + } + + public static DelayPolicy build(String messageDelayLevel) { + return new DelayPolicy(buildList(messageDelayLevel)); + } + + private static List buildList(String messageDelayLevel) { + List delayLevelList = Lists.newArrayList(Splitter.on(" ").split(messageDelayLevel)); + List delayIntervalList = new ArrayList<>(); + for (String delayLevel : delayLevelList) { + final Pattern p = Pattern.compile("(\\d+)([smhd])"); + final Matcher m = p.matcher(delayLevel); + while (m.find()) + { + final int duration = Integer.parseInt(m.group(1)); + final String timeUnitString = m.group(2); + final long interval = toInterval(duration, timeUnitString); + delayIntervalList.add(interval); + } + } + return delayIntervalList; + } + + private static long toInterval(int duration, final String timeUnitString) { + switch (timeUnitString) { + case "s": + return TimeUnit.SECONDS.toMillis(duration); + case "m": + return TimeUnit.MINUTES.toMillis(duration); + case "h": + return TimeUnit.HOURS.toMillis(duration); + case "d": + return TimeUnit.DAYS.toMillis(duration); + default: + throw new IllegalArgumentException(String.format("%s is not a valid code [smhd]", timeUnitString)); + } + } +} \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index a13e7c517c..2072a21eae 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -97,6 +97,7 @@ import org.apache.rocketmq.proxy.grpc.adapter.handler.PullMessageResponseHandler import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.common.Converter; +import org.apache.rocketmq.proxy.grpc.common.DelayPolicy; import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseFuture; import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; @@ -119,6 +120,7 @@ public class LocalGrpcService implements GrpcForwardService { private final ChannelManager channelManager; private final PollCommandResponseManager pollCommandResponseManager; private final RouteService routeService; + private final DelayPolicy delayPolicy; public LocalGrpcService(BrokerController brokerController) { this.brokerController = brokerController; @@ -127,6 +129,7 @@ public class LocalGrpcService implements GrpcForwardService { ConnectorManager connectorManager = new ConnectorManager(null); this.pollCommandResponseManager = new PollCommandResponseManager(); this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager); + this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel()); } @Override public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { @@ -279,7 +282,7 @@ public class LocalGrpcService implements GrpcForwardService { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - ChangeInvisibleTimeRequestHeader requestHeader = Converter.buildChangeInvisibleTimeRequestHeader(request); + ChangeInvisibleTimeRequestHeader requestHeader = Converter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); command.makeCustomHeaderToNet(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 505fb46e18..fb4d10bb7b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -41,6 +41,7 @@ import org.apache.rocketmq.proxy.connector.ForwardReadConsumer; import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.grpc.common.Converter; +import org.apache.rocketmq.proxy.grpc.common.DelayPolicy; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; @@ -59,12 +60,15 @@ public class ConsumerService extends BaseService { private volatile ResponseHook ackMessageHook = null; private volatile ResponseHook nackMessageHook = null; + private final DelayPolicy delayPolicy; + public ConsumerService(ConnectorManager connectorManager) { super(connectorManager); this.readConsumer = connectorManager.getForwardReadConsumer(); this.writeConsumer = connectorManager.getForwardWriteConsumer(); this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache()); + this.delayPolicy = DelayPolicy.build(ConfigurationManager.getProxyConfig().getMessageDelayLevel()); } public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { @@ -214,7 +218,7 @@ public class ConsumerService extends BaseService { } protected ChangeInvisibleTimeRequestHeader convertToChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) { - return Converter.buildChangeInvisibleTimeRequestHeader(request); + return Converter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); } protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, AckResult ackResult) {