mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-19 10:54:54 +08:00
This commit is contained in:
+12
-2
@@ -235,13 +235,23 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen
|
||||
@Override
|
||||
public CompletableFuture<RemotingCommand> request(ProxyContext ctx, String brokerName, RemotingCommand request,
|
||||
long timeoutMillis) {
|
||||
return this.requestBrokerProcessor.request(ctx, brokerName, request, timeoutMillis);
|
||||
int originalRequestOpaque = request.getOpaque();
|
||||
request.setOpaque(RemotingCommand.createNewRequestId());
|
||||
return this.requestBrokerProcessor.request(ctx, brokerName, request, timeoutMillis).thenApply(r -> {
|
||||
request.setOpaque(originalRequestOpaque);
|
||||
return r;
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<Void> requestOneway(ProxyContext ctx, String brokerName, RemotingCommand request,
|
||||
long timeoutMillis) {
|
||||
return this.requestBrokerProcessor.requestOneway(ctx, brokerName, request, timeoutMillis);
|
||||
int originalRequestOpaque = request.getOpaque();
|
||||
request.setOpaque(RemotingCommand.createNewRequestId());
|
||||
return this.requestBrokerProcessor.requestOneway(ctx, brokerName, request, timeoutMillis).thenApply(r -> {
|
||||
request.setOpaque(originalRequestOpaque);
|
||||
return r;
|
||||
});
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Reference in New Issue
Block a user