mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Improve readability.
This commit is contained in:
+17
-19
@@ -39,10 +39,8 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor {
|
||||
}
|
||||
|
||||
@Override
|
||||
public RemotingCommand processRequest(
|
||||
ChannelHandlerContext ctx,
|
||||
RemotingCommand request
|
||||
) throws RemotingCommandException {
|
||||
public RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand request)
|
||||
throws RemotingCommandException {
|
||||
if (request.getCode() == RequestCode.CHECK_TRANSACTION_STATE) {
|
||||
return this.checkTransactionState(ctx, request);
|
||||
}
|
||||
@@ -50,10 +48,8 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor {
|
||||
}
|
||||
|
||||
@Override
|
||||
public RemotingCommand checkTransactionState(
|
||||
ChannelHandlerContext ctx,
|
||||
RemotingCommand request
|
||||
) throws RemotingCommandException {
|
||||
public RemotingCommand checkTransactionState(ChannelHandlerContext ctx, RemotingCommand request)
|
||||
throws RemotingCommandException {
|
||||
final CheckTransactionStateRequestHeader requestHeader =
|
||||
(CheckTransactionStateRequestHeader) request.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class);
|
||||
final ByteBuffer byteBuffer = ByteBuffer.wrap(request.getBody());
|
||||
@@ -61,18 +57,20 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor {
|
||||
if (messageExt != null) {
|
||||
final String group = messageExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP);
|
||||
if (group != null) {
|
||||
transactionStateChecker.checkTransactionState(new TransactionStateCheckRequest(
|
||||
group,
|
||||
requestHeader.getTranStateTableOffset(),
|
||||
requestHeader.getCommitLogOffset(),
|
||||
requestHeader.getMsgId(),
|
||||
TransactionId.genFromBrokerTransactionId(
|
||||
ctx.channel().remoteAddress(),
|
||||
requestHeader.getTransactionId(),
|
||||
transactionStateChecker.checkTransactionState(
|
||||
new TransactionStateCheckRequest(
|
||||
group,
|
||||
requestHeader.getTranStateTableOffset(),
|
||||
requestHeader.getCommitLogOffset(),
|
||||
requestHeader.getTranStateTableOffset()),
|
||||
messageExt
|
||||
));
|
||||
requestHeader.getMsgId(),
|
||||
TransactionId.genFromBrokerTransactionId(
|
||||
ctx.channel().remoteAddress(),
|
||||
requestHeader.getTransactionId(),
|
||||
requestHeader.getCommitLogOffset(),
|
||||
requestHeader.getTranStateTableOffset()),
|
||||
messageExt
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
|
||||
@@ -259,17 +259,15 @@ public class Converter {
|
||||
return endTransactionRequestHeader;
|
||||
}
|
||||
|
||||
public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request, long pollTime) {
|
||||
public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request, long pollTimeoutInMillis) {
|
||||
Partition partition = request.getPartition();
|
||||
String groupName = Converter.getResourceNameWithNamespace(request.getGroup());
|
||||
String topicName = Converter.getResourceNameWithNamespace(partition.getTopic());
|
||||
|
||||
int queueId = partition.getId();
|
||||
int sysFlag = PullSysFlag.buildSysFlag(false, true, true, false, false);
|
||||
String expression = request.getFilterExpression()
|
||||
.getExpression();
|
||||
String expressionType = Converter.buildExpressionType(request.getFilterExpression()
|
||||
.getType());
|
||||
String expression = request.getFilterExpression().getExpression();
|
||||
String expressionType = Converter.buildExpressionType(request.getFilterExpression().getType());
|
||||
|
||||
PullMessageRequestHeader requestHeader = new PullMessageRequestHeader();
|
||||
requestHeader.setConsumerGroup(groupName);
|
||||
@@ -279,7 +277,7 @@ public class Converter {
|
||||
requestHeader.setMaxMsgNums(request.getBatchSize());
|
||||
requestHeader.setSysFlag(sysFlag);
|
||||
requestHeader.setCommitOffset(0L);
|
||||
requestHeader.setSuspendTimeoutMillis(pollTime);
|
||||
requestHeader.setSuspendTimeoutMillis(pollTimeoutInMillis);
|
||||
requestHeader.setSubscription(expression);
|
||||
requestHeader.setSubVersion(0L);
|
||||
requestHeader.setExpressionType(expressionType);
|
||||
@@ -296,22 +294,26 @@ public class Converter {
|
||||
}
|
||||
}
|
||||
MessageAccessor.setProperties(messageWithHeader, Maps.newHashMap(userProperties));
|
||||
|
||||
// set tag
|
||||
String tag = message.getSystemAttribute().getTag();
|
||||
if (!"".equals(tag)) {
|
||||
messageWithHeader.setTags(tag);
|
||||
}
|
||||
|
||||
// set keys
|
||||
List<String> keysList = message.getSystemAttribute().getKeysList();
|
||||
if (keysList.size() > 0) {
|
||||
messageWithHeader.setKeys(keysList);
|
||||
}
|
||||
|
||||
// set message id
|
||||
String messageId = message.getSystemAttribute().getMessageId();
|
||||
if ("".equals(messageId)) {
|
||||
throw new IllegalArgumentException("message id is empty");
|
||||
}
|
||||
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX, messageId);
|
||||
|
||||
// set transaction property
|
||||
MessageType messageType = message.getSystemAttribute().getMessageType();
|
||||
if (messageType.equals(MessageType.TRANSACTION)) {
|
||||
@@ -328,8 +330,7 @@ public class Converter {
|
||||
case DELAY_LEVEL:
|
||||
int delayLevel = message.getSystemAttribute().getDelayLevel();
|
||||
if (delayLevel > 0) {
|
||||
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_DELAY_TIME_LEVEL,
|
||||
String.valueOf(delayLevel));
|
||||
MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_DELAY_TIME_LEVEL, String.valueOf(delayLevel));
|
||||
}
|
||||
break;
|
||||
case DELIVERY_TIMESTAMP:
|
||||
|
||||
@@ -68,9 +68,9 @@ import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService;
|
||||
import org.slf4j.Logger;
|
||||
@@ -115,9 +115,11 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
@Override
|
||||
public CompletableFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
this.clientService.heartbeat(ctx, request, channelManager);
|
||||
return CompletableFuture.completedFuture(HeartbeatResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build());
|
||||
return CompletableFuture.completedFuture(
|
||||
HeartbeatResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build()
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -194,13 +196,16 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
public CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
|
||||
NotifyClientTerminationRequest request) {
|
||||
this.clientService.unregister(ctx, request, channelManager);
|
||||
return CompletableFuture.completedFuture(NotifyClientTerminationResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build());
|
||||
return CompletableFuture.completedFuture(
|
||||
NotifyClientTerminationResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build()
|
||||
);
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
@Override
|
||||
public CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
+8
-8
@@ -40,7 +40,7 @@ import org.apache.rocketmq.proxy.grpc.common.ResponseHook;
|
||||
|
||||
public class TransactionService extends BaseService implements TransactionStateChecker {
|
||||
|
||||
private ChannelManager channelManager;
|
||||
private final ChannelManager channelManager;
|
||||
private final ForwardProducer forwardProducer;
|
||||
|
||||
private volatile ResponseHook<TransactionStateCheckRequest, PollCommandResponse> checkTransactionStateHook = null;
|
||||
@@ -56,23 +56,23 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
public void checkTransactionState(TransactionStateCheckRequest checkData) {
|
||||
try {
|
||||
List<String> clientIdList = this.channelManager.getClientIdList(checkData.getGroupId());
|
||||
// if clientIdList's size is 0, here will throw: java.lang.IllegalArgumentException: bound must be positive
|
||||
String clientId = clientIdList.get(ThreadLocalRandom.current().nextInt(clientIdList.size()));
|
||||
|
||||
GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId);
|
||||
String transactionId = checkData.getTransactionId().getProxyTransactionId();
|
||||
|
||||
Message message = Converter.buildMessage(checkData.getMessageExt());
|
||||
|
||||
PollCommandResponse commandResponse = PollCommandResponse.newBuilder()
|
||||
PollCommandResponse response = PollCommandResponse.newBuilder()
|
||||
.setRecoverOrphanedTransactionCommand(
|
||||
RecoverOrphanedTransactionCommand.newBuilder()
|
||||
.setOrphanedTransactionalMessage(message)
|
||||
.setTransactionId(transactionId)
|
||||
.build()
|
||||
).build();
|
||||
channel.writeAndFlush(commandResponse);
|
||||
channel.writeAndFlush(response);
|
||||
if (this.checkTransactionStateHook != null) {
|
||||
this.checkTransactionStateHook.beforeResponse(checkData, commandResponse, null);
|
||||
this.checkTransactionStateHook.beforeResponse(checkData, response, null);
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
if (this.checkTransactionStateHook != null) {
|
||||
@@ -88,8 +88,9 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
endTransactionHook.beforeResponse(request, response, throwable);
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
EndTransactionRequestHeader requestHeader = this.convertToEndTransactionRequestHeader(ctx, request);
|
||||
EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request);
|
||||
this.forwardProducer.endTransaction(requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
future.complete(EndTransactionResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
@@ -100,8 +101,7 @@ public class TransactionService extends BaseService implements TransactionStateC
|
||||
return future;
|
||||
}
|
||||
|
||||
protected EndTransactionRequestHeader convertToEndTransactionRequestHeader(Context ctx,
|
||||
EndTransactionRequest request) {
|
||||
protected EndTransactionRequestHeader toEndTransactionRequestHeader(Context ctx, EndTransactionRequest request) {
|
||||
return Converter.buildEndTransactionRequestHeader(request);
|
||||
}
|
||||
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.common.utils;
|
||||
|
||||
import java.util.concurrent.ThreadLocalRandom;
|
||||
import org.apache.rocketmq.common.filter.FilterAPI;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.junit.Test;
|
||||
|
||||
Reference in New Issue
Block a user