mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] Do some refactoring work.
This commit is contained in:
@@ -46,7 +46,8 @@ public abstract class AbstractForwardClient implements StartAndShutdown {
|
||||
if (clients.length == 1) {
|
||||
return this.clients[0];
|
||||
}
|
||||
return this.clients[ThreadLocalRandom.current().nextInt(this.clients.length)];
|
||||
int index = ThreadLocalRandom.current().nextInt(this.clients.length);
|
||||
return this.clients[index];
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -87,7 +87,7 @@ public class ForwardProducer extends AbstractForwardClient {
|
||||
return future.thenApply(sendResult -> {
|
||||
int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag());
|
||||
if (SendStatus.SEND_OK.equals(sendResult.getSendStatus()) && tranType == MessageSysFlag.TRANSACTION_PREPARED_TYPE) {
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(address, sendResult);
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(address, sendResult);
|
||||
sendResult.setTransactionId(transactionId.getProxyTransactionId());
|
||||
}
|
||||
return sendResult;
|
||||
|
||||
+1
-1
@@ -30,7 +30,7 @@ import org.apache.rocketmq.remoting.RPCHook;
|
||||
|
||||
public class ForwardClientManager implements StartAndShutdown {
|
||||
|
||||
private RPCHook rpcHook = null;
|
||||
private RPCHook rpcHook;
|
||||
|
||||
private final MQClientFactory mqClientFactory;
|
||||
private final TransactionProducerFactory transactionalProducerFactory;
|
||||
|
||||
+1
-2
@@ -23,8 +23,7 @@ import org.apache.rocketmq.remoting.RPCHook;
|
||||
|
||||
public class MQClientFactory extends AbstractMQClientFactory {
|
||||
|
||||
public MQClientFactory(ScheduledExecutorService scheduledExecutorService,
|
||||
RPCHook rpcHook) {
|
||||
public MQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) {
|
||||
super(scheduledExecutorService, rpcHook);
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -63,7 +63,7 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor {
|
||||
requestHeader.getTranStateTableOffset(),
|
||||
requestHeader.getCommitLogOffset(),
|
||||
requestHeader.getMsgId(),
|
||||
TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId.genByBrokerTransactionId(
|
||||
ctx.channel().remoteAddress(),
|
||||
requestHeader.getTransactionId(),
|
||||
requestHeader.getCommitLogOffset(),
|
||||
|
||||
@@ -81,7 +81,8 @@ public class TopicRouteCache {
|
||||
}
|
||||
|
||||
public SelectableMessageQueue selectOneWriteQueue(String topic, String brokerName, int queueId) throws Exception {
|
||||
return getMessageQueue(topic).getWriteSelector().selectOne(brokerName, queueId);
|
||||
return getMessageQueue(topic).getWriteSelector()
|
||||
.selectOne(brokerName, queueId);
|
||||
}
|
||||
|
||||
public SelectableMessageQueue selectOneWriteQueueByKey(String topic, String shardingKey) throws Exception {
|
||||
|
||||
-1
@@ -40,7 +40,6 @@ import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class TransactionHeartbeatRegisterService implements StartAndShutdown {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(TransactionHeartbeatRegisterService.class);
|
||||
|
||||
private static final String TRANS_HEARTBEAT_CLIENT_ID = "rmq-proxy-producer-client";
|
||||
|
||||
+3
-3
@@ -54,7 +54,7 @@ public class TransactionId {
|
||||
public TransactionId() {
|
||||
}
|
||||
|
||||
public static TransactionId genFromBrokerTransactionId(String brokerAddr, SendResult sendResult) {
|
||||
public static TransactionId genByBrokerTransactionId(String brokerAddr, SendResult sendResult) {
|
||||
MessageId id = new MessageId(null, 0);
|
||||
try {
|
||||
if (sendResult.getOffsetMsgId() != null) {
|
||||
@@ -65,11 +65,11 @@ public class TransactionId {
|
||||
} catch (Exception e) {
|
||||
log.warn("genFromBrokerTransactionId failed. brokerAddr: {}, sendResult: {}", brokerAddr, sendResult, e);
|
||||
}
|
||||
return genFromBrokerTransactionId(RemotingUtil.string2SocketAddress(brokerAddr), sendResult.getTransactionId(),
|
||||
return genByBrokerTransactionId(RemotingUtil.string2SocketAddress(brokerAddr), sendResult.getTransactionId(),
|
||||
id.getOffset(), sendResult.getQueueOffset());
|
||||
}
|
||||
|
||||
public static TransactionId genFromBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId,
|
||||
public static TransactionId genByBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId,
|
||||
long commitLogOffset, long tranStateTableOffset) {
|
||||
byte[] orgTransactionIdByte = new byte[0];
|
||||
if (StringUtils.isNotBlank(orgTransactionId)) {
|
||||
|
||||
@@ -189,7 +189,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
}
|
||||
|
||||
@Override
|
||||
public void forwardMessageToDeadLetterQueue(ForwardMessageToDeadLetterQueueRequest request, StreamObserver<ForwardMessageToDeadLetterQueueResponse> responseObserver) {
|
||||
public void forwardMessageToDeadLetterQueue(ForwardMessageToDeadLetterQueueRequest request,
|
||||
StreamObserver<ForwardMessageToDeadLetterQueueResponse> responseObserver) {
|
||||
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> future = grpcForwardService.forwardMessageToDeadLetterQueue(Context.current(), request);
|
||||
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
|
||||
.exceptionally(e -> {
|
||||
@@ -251,7 +252,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
}
|
||||
|
||||
@Override
|
||||
public void reportThreadStackTrace(ReportThreadStackTraceRequest request, StreamObserver<ReportThreadStackTraceResponse> responseObserver) {
|
||||
public void reportThreadStackTrace(ReportThreadStackTraceRequest request,
|
||||
StreamObserver<ReportThreadStackTraceResponse> responseObserver) {
|
||||
CompletableFuture<ReportThreadStackTraceResponse> future = grpcForwardService.reportThreadStackTrace(Context.current(), request);
|
||||
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
|
||||
.exceptionally(e -> {
|
||||
@@ -264,7 +266,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
}
|
||||
|
||||
@Override
|
||||
public void reportMessageConsumptionResult(ReportMessageConsumptionResultRequest request, StreamObserver<ReportMessageConsumptionResultResponse> responseObserver) {
|
||||
public void reportMessageConsumptionResult(ReportMessageConsumptionResultRequest request,
|
||||
StreamObserver<ReportMessageConsumptionResultResponse> responseObserver) {
|
||||
CompletableFuture<ReportMessageConsumptionResultResponse> future = grpcForwardService.reportMessageConsumptionResult(Context.current(), request);
|
||||
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
|
||||
.exceptionally(e -> {
|
||||
@@ -277,7 +280,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
}
|
||||
|
||||
@Override
|
||||
public void notifyClientTermination(NotifyClientTerminationRequest request, StreamObserver<NotifyClientTerminationResponse> responseObserver) {
|
||||
public void notifyClientTermination(NotifyClientTerminationRequest request,
|
||||
StreamObserver<NotifyClientTerminationResponse> responseObserver) {
|
||||
CompletableFuture<NotifyClientTerminationResponse> future = grpcForwardService.notifyClientTermination(Context.current(), request);
|
||||
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
|
||||
.exceptionally(e -> {
|
||||
@@ -290,7 +294,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
}
|
||||
|
||||
@Override
|
||||
public void changeInvisibleDuration(ChangeInvisibleDurationRequest request, StreamObserver<ChangeInvisibleDurationResponse> responseObserver) {
|
||||
public void changeInvisibleDuration(ChangeInvisibleDurationRequest request,
|
||||
StreamObserver<ChangeInvisibleDurationResponse> responseObserver) {
|
||||
CompletableFuture<ChangeInvisibleDurationResponse> future = grpcForwardService.changeInvisibleDuration(Context.current(), request);
|
||||
future.thenAccept(response -> ResponseWriter.write(responseObserver, response))
|
||||
.exceptionally(e -> {
|
||||
|
||||
@@ -97,6 +97,7 @@ public class GrpcServer implements StartAndShutdown {
|
||||
.addService(messagingProcessor)
|
||||
.executor(this.executor);
|
||||
|
||||
// grpc interceptors, including acl, logging etc.
|
||||
if (ConfigurationManager.getProxyConfig().isEnableACL()) {
|
||||
List<AccessValidator> accessValidators = ServiceProvider.load(ServiceProvider.ACL_VALIDATOR_ID, AccessValidator.class);
|
||||
if (accessValidators.isEmpty()) {
|
||||
@@ -122,7 +123,7 @@ public class GrpcServer implements StartAndShutdown {
|
||||
this.grpcForwardService.start();
|
||||
|
||||
this.server.start();
|
||||
log.info("grpc server has started");
|
||||
log.info("grpc server start successfully.");
|
||||
}
|
||||
|
||||
public void shutdown() {
|
||||
@@ -132,7 +133,7 @@ public class GrpcServer implements StartAndShutdown {
|
||||
|
||||
this.grpcForwardService.shutdown();
|
||||
|
||||
log.info("grpc server has stopped");
|
||||
log.info("grpc server shutdown successfully.");
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
|
||||
@@ -175,8 +175,7 @@ public class GrpcConverter {
|
||||
public static AckMessageRequestHeader buildAckMessageRequestHeader(AckMessageRequest request) {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
String receiptHandleStr = request.getReceiptHandle();
|
||||
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
|
||||
ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
|
||||
|
||||
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
|
||||
ackMessageRequestHeader.setConsumerGroup(groupName);
|
||||
@@ -191,8 +190,7 @@ public class GrpcConverter {
|
||||
DelayPolicy delayPolicy) {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
String receiptHandleStr = request.getReceiptHandle();
|
||||
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
|
||||
ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
|
||||
|
||||
ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader();
|
||||
changeInvisibleTimeRequestHeader.setConsumerGroup(groupName);
|
||||
@@ -209,8 +207,7 @@ public class GrpcConverter {
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
String receiptHandleStr = request.getReceiptHandle();
|
||||
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
|
||||
ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
|
||||
|
||||
ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader();
|
||||
changeInvisibleTimeRequestHeader.setConsumerGroup(groupName);
|
||||
@@ -226,8 +223,7 @@ public class GrpcConverter {
|
||||
ForwardMessageToDeadLetterQueueRequest request) {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
String receiptHandleStr = request.getReceiptHandle();
|
||||
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
|
||||
ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
|
||||
|
||||
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
|
||||
consumerSendMsgBackRequestHeader.setOffset(handle.getCommitLogOffset());
|
||||
@@ -243,8 +239,7 @@ public class GrpcConverter {
|
||||
NackMessageRequest request) {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic());
|
||||
String receiptHandleStr = request.getReceiptHandle();
|
||||
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
|
||||
ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle());
|
||||
|
||||
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
|
||||
consumerSendMsgBackRequestHeader.setOffset(handle.getCommitLogOffset());
|
||||
|
||||
+4
-2
@@ -68,10 +68,12 @@ public class ForwardClientService extends BaseService {
|
||||
this.pollCommandResponseManager = pollCommandResponseManager;
|
||||
|
||||
this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() {
|
||||
@Override public void handle(ConsumerGroupEvent event, String group, Object... args) {
|
||||
@Override
|
||||
public void handle(ConsumerGroupEvent event, String group, Object... args) {
|
||||
}
|
||||
|
||||
@Override public void shutdown() {
|
||||
@Override
|
||||
public void shutdown() {
|
||||
}
|
||||
});
|
||||
this.producerManager = new ProducerManager();
|
||||
|
||||
+3
-3
@@ -10,7 +10,7 @@ public class TransactionIdTest {
|
||||
|
||||
@Test
|
||||
public void test() throws UnknownHostException {
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingHelper.string2SocketAddress("127.0.0.1:8080"),
|
||||
"71F99B78B6E261357FA259CCA6456118", 1234, 5678);
|
||||
|
||||
@@ -24,7 +24,7 @@ public class TransactionIdTest {
|
||||
|
||||
@Test
|
||||
public void testEmptyTransactionId() throws UnknownHostException {
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingHelper.string2SocketAddress("127.0.0.1:8080"),
|
||||
"", 1234, 5678);
|
||||
|
||||
@@ -38,7 +38,7 @@ public class TransactionIdTest {
|
||||
|
||||
@Test
|
||||
public void testNullTransactionId() throws UnknownHostException {
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingHelper.string2SocketAddress("127.0.0.1:8080"),
|
||||
null, 1234, 5678);
|
||||
|
||||
|
||||
+1
-1
@@ -381,7 +381,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
.thenReturn(response);
|
||||
EndTransactionRequest request = EndTransactionRequest.newBuilder()
|
||||
.setMessageId("123")
|
||||
.setTransactionId(TransactionId.genFromBrokerTransactionId(
|
||||
.setTransactionId(TransactionId.genByBrokerTransactionId(
|
||||
new InetSocketAddress("0.0.0.0", 80), "123", 123, 123
|
||||
).getProxyTransactionId()
|
||||
)
|
||||
|
||||
+2
-2
@@ -48,7 +48,7 @@ public class TransactionServiceTest extends BaseServiceTest {
|
||||
return null;
|
||||
}).when(channel).writeAndFlush(any());
|
||||
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingHelper.string2SocketAddress("127.0.0.1:8080"),
|
||||
"71F99B78B6E261357FA259CCA6456118", 1234, 5678);
|
||||
transactionService.checkTransactionState(new TransactionStateCheckRequest(
|
||||
@@ -69,7 +69,7 @@ public class TransactionServiceTest extends BaseServiceTest {
|
||||
public void testEndTransaction() throws Exception {
|
||||
AtomicReference<EndTransactionRequestHeader> headerRef = new AtomicReference<>();
|
||||
AtomicReference<String> brokerAddrRef = new AtomicReference<>();
|
||||
TransactionId transactionId = TransactionId.genFromBrokerTransactionId(
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingHelper.string2SocketAddress("127.0.0.1:8080"),
|
||||
"71F99B78B6E261357FA259CCA6456118", 1234, 5678);
|
||||
doAnswer(mock -> {
|
||||
|
||||
Reference in New Issue
Block a user