mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Add DelayPolicy
This commit is contained in:
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<Long> delayIntervalList;
|
||||
|
||||
private DelayPolicy(List<Long> 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<Long> buildList(String messageDelayLevel) {
|
||||
List<String> delayLevelList = Lists.newArrayList(Splitter.on(" ").split(messageDelayLevel));
|
||||
List<Long> 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));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<QueryRouteResponse> 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();
|
||||
|
||||
|
||||
+5
-1
@@ -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<AckMessageRequest, AckMessageResponse> ackMessageHook = null;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> 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<ReceiveMessageResponse> 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) {
|
||||
|
||||
Reference in New Issue
Block a user