[ISSUE #3949] ack msg when tag not match; refactor client factory

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:12 +08:00
committed by zhouxiang
parent 1dc5b7c415
commit 3f07761f60
11 changed files with 186 additions and 66 deletions
@@ -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<PopResult> 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);
@@ -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 +
']';
}
}
@@ -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;
}
@@ -86,14 +86,13 @@ public class ForwardProducer extends AbstractForwardClient {
public CompletableFuture<SendResult> sendMessage(String address, String brokerName, Message msg,
SendMessageRequestHeader requestHeader, long timeoutMillis) {
CompletableFuture<SendResult> 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<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
@@ -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<T> {
private static final Logger LOGGER = LoggerFactory.getLogger(AbstractClientFactory.class);
protected final ScheduledExecutorService scheduledExecutorService;
protected Map<String, T> 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<T> {
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<T> {
return object;
}
// @Override
public void shutdownAll() {
this.cacheTable.forEach((k, v) -> {
try {
@@ -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<MQClientAPIExtImpl> {
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();
}
}
@@ -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() {
@@ -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<MQClientAPIExtImpl> {
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);
}
}
@@ -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<MQClientAPIExtImpl> {
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);
}
}
@@ -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();
@@ -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<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook = null;
private volatile ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook = null;
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook = null;
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> 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<Message> 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<AckResult> 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<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request) {
CompletableFuture<AckMessageResponse> future = new CompletableFuture<>();
future.whenComplete((response, throwable) -> {
@@ -242,6 +285,11 @@ public class ConsumerService extends BaseService {
this.receiveMessageHook = receiveMessageHook;
}
public void setAckNoMatchedMessageHook(
ResponseHook<AckMessageRequestHeader, AckResult> ackNoMatchedMessageHook) {
this.ackNoMatchedMessageHook = ackNoMatchedMessageHook;
}
public void setAckMessageHook(
ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook) {
this.ackMessageHook = ackMessageHook;