From 3f07761f60efd00ae2f809f049b9802912b5f781 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Fri, 18 Mar 2022 17:22:37 +0800 Subject: [PATCH] [ISSUE #3949] ack msg when tag not match; refactor client factory --- .../client/impl/MQClientAPIExtImpl.java | 19 +++++- .../ChangeInvisibleTimeRequestHeader.java | 11 ++++ .../proxy/common/utils/ProxyUtils.java | 2 + .../proxy/connector/ForwardProducer.java | 3 +- .../factory/AbstractClientFactory.java | 7 +- .../factory/AbstractMQClientFactory.java | 66 +++++++++++++++++++ .../factory/ForwardClientFactory.java | 10 ++- .../connector/factory/MQClientFactory.java | 34 +++------- .../factory/TransactionProducerFactory.java | 34 +++------- .../rocketmq/proxy/grpc/common/Converter.java | 6 ++ .../grpc/service/cluster/ConsumerService.java | 60 +++++++++++++++-- 11 files changed, 186 insertions(+), 66 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java index d77159d87d..8b0463bca6 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java @@ -54,9 +54,13 @@ import org.apache.rocketmq.remoting.exception.RemotingException; import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.apache.rocketmq.remoting.netty.ResponseFuture; import org.apache.rocketmq.remoting.protocol.RemotingCommand; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; public class MQClientAPIExtImpl { + private static final Logger log = LoggerFactory.getLogger(MQClientAPIExtImpl.class); + private final MQClientAPIImpl mqClientAPI; private final ClientConfig clientConfig; @@ -75,6 +79,19 @@ public class MQClientAPIExtImpl { this.mqClientAPI.shutdown(); } + public void fetchNameServerAddr() { + this.mqClientAPI.fetchNameServerAddr(); + } + + public boolean updateNameServerAddressList() { + if (this.clientConfig.getNamesrvAddr() != null) { + this.mqClientAPI.updateNameServerAddressList(this.clientConfig.getNamesrvAddr()); + log.info("user specified name server address: {}", this.clientConfig.getNamesrvAddr()); + return true; + } + return false; + } + protected static MQClientException processNullResponseErr(ResponseFuture responseFuture) { MQClientException ex; if (!responseFuture.isSendRequestOK()) { @@ -174,7 +191,7 @@ public class MQClientAPIExtImpl { long timeoutMillis) { CompletableFuture future = new CompletableFuture<>(); try { - this.mqClientAPI.popMessageAsync(brokerAddr, brokerName, requestHeader, timeoutMillis, new PopCallback() { + this.mqClientAPI.popMessageAsync(brokerName, brokerAddr, requestHeader, timeoutMillis, new PopCallback() { @Override public void onSuccess(PopResult popResult) { future.complete(popResult); diff --git a/common/src/main/java/org/apache/rocketmq/common/protocol/header/ChangeInvisibleTimeRequestHeader.java b/common/src/main/java/org/apache/rocketmq/common/protocol/header/ChangeInvisibleTimeRequestHeader.java index a586e490cf..f01e89c725 100644 --- a/common/src/main/java/org/apache/rocketmq/common/protocol/header/ChangeInvisibleTimeRequestHeader.java +++ b/common/src/main/java/org/apache/rocketmq/common/protocol/header/ChangeInvisibleTimeRequestHeader.java @@ -94,4 +94,15 @@ public class ChangeInvisibleTimeRequestHeader implements CommandCustomHeader { this.queueId = queueId; } + @Override + public String toString() { + return "ChangeInvisibleTimeRequestHeader [" + + "consumerGroup='" + consumerGroup + '\'' + + ", topic='" + topic + '\'' + + ", queueId=" + queueId + + ", extraInfo='" + extraInfo + '\'' + + ", offset=" + offset + + ", invisibleTime=" + invisibleTime + + ']'; + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java index f335cdfa0a..f7c05450b6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java @@ -20,4 +20,6 @@ import java.time.Duration; public class ProxyUtils { public static final long DEFAULT_MQ_CLIENT_TIMEOUT = Duration.ofSeconds(3).toMillis(); + + public static final int MAX_MSG_NUMS_FOR_POP_REQUEST = 32; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java index d35b65fd24..30bc4a939b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java @@ -86,14 +86,13 @@ public class ForwardProducer extends AbstractForwardClient { public CompletableFuture sendMessage(String address, String brokerName, Message msg, SendMessageRequestHeader requestHeader, long timeoutMillis) { CompletableFuture future = this.getClient().sendMessage(address, brokerName, msg, requestHeader, timeoutMillis); - future.thenApply(sendResult -> { + return future.thenApply(sendResult -> { if (SendStatus.SEND_OK.equals(sendResult.getSendStatus()) && !StringUtils.isEmpty(sendResult.getTransactionId())) { TransactionId transactionId = TransactionId.genFromBrokerTransactionId(address, sendResult); sendResult.setTransactionId(transactionId.getProxyTransactionId()); } return sendResult; }); - return future; } public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java index 7e0519097b..daae5bc2bd 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java @@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.connector.factory; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledExecutorService; import org.apache.rocketmq.remoting.RPCHook; import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.slf4j.Logger; @@ -26,10 +27,12 @@ import org.slf4j.LoggerFactory; public abstract class AbstractClientFactory { private static final Logger LOGGER = LoggerFactory.getLogger(AbstractClientFactory.class); + protected final ScheduledExecutorService scheduledExecutorService; protected Map cacheTable = new ConcurrentHashMap<>(); protected RPCHook rpcHook; - public AbstractClientFactory(RPCHook rpcHook) { + public AbstractClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { + this.scheduledExecutorService = scheduledExecutorService; this.rpcHook = rpcHook; } @@ -47,7 +50,6 @@ public abstract class AbstractClientFactory { return nettyClientConfig; } -// @Override public T getOne(String instanceName, int bootstrapWorkerThreads) { if (cacheTable.containsKey(instanceName)) { return cacheTable.get(instanceName); @@ -71,7 +73,6 @@ public abstract class AbstractClientFactory { return object; } -// @Override public void shutdownAll() { this.cacheTable.forEach((k, v) -> { try { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java new file mode 100644 index 0000000000..c7417d4d46 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java @@ -0,0 +1,66 @@ +/* + * 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.connector.factory; + +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.client.ClientConfig; +import org.apache.rocketmq.client.impl.ClientRemotingProcessor; +import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.remoting.RPCHook; + +public abstract class AbstractMQClientFactory extends AbstractClientFactory { + + public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, + RPCHook rpcHook) { + super(scheduledExecutorService, rpcHook); + } + + protected abstract ClientRemotingProcessor createClientRemotingProcessor(); + + @Override + protected MQClientAPIExtImpl newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) { + ClientConfig clientConfig = new ClientConfig(); + clientConfig.setInstanceName(instanceName); + + return new MQClientAPIExtImpl( + createNettyClientConfig(bootstrapWorkerThreads), + createClientRemotingProcessor(), + rpcHook, + clientConfig + ); + } + + @Override + protected boolean tryStart(MQClientAPIExtImpl client) { + if (!client.updateNameServerAddressList()) { + this.scheduledExecutorService.scheduleAtFixedRate( + client::fetchNameServerAddr, + 1000 * 10, + 1000 * 60 * 2, + TimeUnit.MILLISECONDS + ); + } + client.start(); + return true; + } + + @Override + protected void shutdown(MQClientAPIExtImpl client) { + client.shutdown(); + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java index 08b43acf2c..350111011a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java @@ -16,6 +16,9 @@ */ package org.apache.rocketmq.proxy.connector.factory; +import com.google.common.util.concurrent.ThreadFactoryBuilder; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.ClientConfig; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; @@ -35,8 +38,11 @@ public class ForwardClientFactory implements StartAndShutdown { public ForwardClientFactory(TransactionStateChecker transactionStateChecker) { this.init(); - this.mqClientFactory = new MQClientFactory(this.rpcHook); - this.transactionalProducerFactory = new TransactionProducerFactory(this.rpcHook, transactionStateChecker); + ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor( + new ThreadFactoryBuilder().setNameFormat("ForwardClientFactoryScheduledThread" + "-%d").build() + ); + this.mqClientFactory = new MQClientFactory(scheduledExecutorService, this.rpcHook); + this.transactionalProducerFactory = new TransactionProducerFactory(scheduledExecutorService, this.rpcHook, transactionStateChecker); } private void init() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java index 3b9a75f33b..07484bd54b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java @@ -16,38 +16,20 @@ */ package org.apache.rocketmq.proxy.connector.factory; -import org.apache.rocketmq.client.ClientConfig; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import java.util.concurrent.ScheduledExecutorService; +import org.apache.rocketmq.client.impl.ClientRemotingProcessor; import org.apache.rocketmq.proxy.connector.processor.DoNothingClientRemotingProcessor; import org.apache.rocketmq.remoting.RPCHook; -public class MQClientFactory extends AbstractClientFactory { +public class MQClientFactory extends AbstractMQClientFactory { - public MQClientFactory(RPCHook rpcHook) { - super(rpcHook); + public MQClientFactory(ScheduledExecutorService scheduledExecutorService, + RPCHook rpcHook) { + super(scheduledExecutorService, rpcHook); } @Override - protected MQClientAPIExtImpl newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) { - ClientConfig clientConfig = new ClientConfig(); - clientConfig.setInstanceName(instanceName); - - return new MQClientAPIExtImpl( - createNettyClientConfig(bootstrapWorkerThreads), - new DoNothingClientRemotingProcessor(null), - rpcHook, - clientConfig - ); - } - - @Override - protected boolean tryStart(MQClientAPIExtImpl client) { - client.start(); - return true; - } - - @Override - protected void shutdown(MQClientAPIExtImpl client) { - client.shutdown(); + protected ClientRemotingProcessor createClientRemotingProcessor() { + return new DoNothingClientRemotingProcessor(null); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/TransactionProducerFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/TransactionProducerFactory.java index a97727c11d..92a4006479 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/TransactionProducerFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/TransactionProducerFactory.java @@ -16,41 +16,23 @@ */ package org.apache.rocketmq.proxy.connector.factory; -import org.apache.rocketmq.client.ClientConfig; -import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import java.util.concurrent.ScheduledExecutorService; +import org.apache.rocketmq.client.impl.ClientRemotingProcessor; import org.apache.rocketmq.proxy.connector.processor.ProxyClientRemotingProcessor; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; import org.apache.rocketmq.remoting.RPCHook; -public class TransactionProducerFactory extends AbstractClientFactory { +public class TransactionProducerFactory extends AbstractMQClientFactory { private final TransactionStateChecker transactionStateChecker; - public TransactionProducerFactory(RPCHook rpcHook, TransactionStateChecker transactionStateChecker) { - super(rpcHook); + public TransactionProducerFactory(ScheduledExecutorService scheduledExecutorService, + RPCHook rpcHook, TransactionStateChecker transactionStateChecker) { + super(scheduledExecutorService, rpcHook); this.transactionStateChecker = transactionStateChecker; } @Override - public MQClientAPIExtImpl newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) { - ClientConfig clientConfig = new ClientConfig(); - clientConfig.setInstanceName(instanceName); - - return new MQClientAPIExtImpl( - createNettyClientConfig(bootstrapWorkerThreads), - new ProxyClientRemotingProcessor(this.transactionStateChecker), - rpcHook, - clientConfig - ); - } - - @Override - protected boolean tryStart(MQClientAPIExtImpl client) { - client.start(); - return true; - } - - @Override - protected void shutdown(MQClientAPIExtImpl client) { - client.shutdown(); + protected ClientRemotingProcessor createClientRemotingProcessor() { + return new ProxyClientRemotingProcessor(this.transactionStateChecker); } } 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 f60fac135f..e1bd1f35d6 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.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.slf4j.Logger; @@ -140,6 +141,11 @@ public class Converter { String topicName = Converter.getResourceNameWithNamespace(topic); int queueId = partition.getId(); int maxMessageNumbers = request.getBatchSize(); + if (maxMessageNumbers > ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST) { + LOGGER.warn("change maxNums from {} to {} for pop request, with info: topic:{}, group:{}", + maxMessageNumbers, ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST, topicName, groupName); + maxMessageNumbers = ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST; + } long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); long bornTime = Timestamps.toMillis(request.getInitializationTimestamp()); ConsumePolicy policy = request.getConsumePolicy(); 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 fb4d10bb7b..21b8edc980 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 @@ -34,6 +34,8 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; +import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.proxy.common.utils.FilterUtil; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; @@ -57,6 +59,7 @@ public class ConsumerService extends BaseService { private volatile ReadQueueSelector readQueueSelector; private volatile ResponseHook receiveMessageHook = null; + private volatile ResponseHook ackNoMatchedMessageHook = null; private volatile ResponseHook ackMessageHook = null; private volatile ResponseHook nackMessageHook = null; @@ -87,13 +90,18 @@ public class ConsumerService extends BaseService { messageQueue.getBrokerName(), requestHeader, requestHeader.getPollTime()); - popResultFuture.thenAccept(result -> { - try { - future.complete(convertToReceiveMessageResponse(ctx, request, result)); - } catch (Throwable throwable) { + popResultFuture + .thenAccept(result -> { + try { + future.complete(convertToReceiveMessageResponse(ctx, request, result)); + } catch (Throwable throwable) { + future.completeExceptionally(throwable); + } + }) + .exceptionally(throwable -> { future.completeExceptionally(throwable); - } - }); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } @@ -101,6 +109,9 @@ public class ConsumerService extends BaseService { } protected PopMessageRequestHeader convertToPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) { + // check filterExpression is correct or not + Converter.buildSubscriptionData(Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); + long timeRemaining = ctx.getDeadline() .timeRemaining(TimeUnit.MILLISECONDS); long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); @@ -112,6 +123,8 @@ public class ConsumerService extends BaseService { } protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) { + SubscriptionData subscriptionData = Converter.buildSubscriptionData( + Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); PopStatus status = result.getPopStatus(); switch (status) { case FOUND: @@ -130,6 +143,10 @@ public class ConsumerService extends BaseService { List messages = new ArrayList<>(); for (MessageExt messageExt : result.getMsgFoundList()) { + if (FilterUtil.isTagNotMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { + this.ackNoMatchedMessage(ctx, request, messageExt); + continue; + } messages.add(Converter.buildMessage(messageExt)); } @@ -139,6 +156,32 @@ public class ConsumerService extends BaseService { .build(); } + protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { + CompletableFuture future = new CompletableFuture<>(); + AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader(); + try { + ReceiptHandle handle = ReceiptHandle.create(messageExt); + if (handle == null) { + return; + } + String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName()); + ackMessageRequestHeader.setConsumerGroup(Converter.getResourceNameWithNamespace(request.getGroup())); + ackMessageRequestHeader.setTopic(messageExt.getTopic()); + ackMessageRequestHeader.setQueueId(handle.getQueueId()); + ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); + ackMessageRequestHeader.setOffset(handle.getOffset()); + + future = this.writeConsumer.ackMessage(brokerAddr, ackMessageRequestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + } catch (Throwable t) { + future.completeExceptionally(t); + } + future.whenComplete((ackResult, throwable) -> { + if (ackNoMatchedMessageHook != null) { + ackNoMatchedMessageHook.beforeResponse(ackMessageRequestHeader, ackResult, throwable); + } + }); + } + public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> { @@ -242,6 +285,11 @@ public class ConsumerService extends BaseService { this.receiveMessageHook = receiveMessageHook; } + public void setAckNoMatchedMessageHook( + ResponseHook ackNoMatchedMessageHook) { + this.ackNoMatchedMessageHook = ackNoMatchedMessageHook; + } + public void setAckMessageHook( ResponseHook ackMessageHook) { this.ackMessageHook = ackMessageHook;