mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
+36
-1
@@ -17,14 +17,44 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service;
|
||||
|
||||
import apache.rocketmq.v2.Publishing;
|
||||
import apache.rocketmq.v2.RetryPolicy;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import apache.rocketmq.v2.Subscription;
|
||||
import com.google.protobuf.util.Durations;
|
||||
import io.grpc.Context;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
|
||||
public class GrpcClientManager {
|
||||
|
||||
|
||||
private static final Settings DEFAULT_PRODUCER_SETTINGS = Settings.newBuilder()
|
||||
.setPublishing(Publishing.newBuilder()
|
||||
.setRetryPolicy(RetryPolicy.newBuilder()
|
||||
.setMaxAttempts(3)
|
||||
.setInitialBackoff(1)
|
||||
.setMaxBackoff(4)
|
||||
.setBackoffMultiplier(2)
|
||||
.build())
|
||||
.setCompressBodyThreshold(4 * 1024)
|
||||
.setMaxBodySize(4 * 1024 * 1024)
|
||||
.build())
|
||||
.build();
|
||||
private static final Settings DEFAULT_CONSUMER_SETTINGS = Settings.newBuilder()
|
||||
.setSubscription(Subscription.newBuilder()
|
||||
.setFifo(false)
|
||||
.setBackoffPolicy(RetryPolicy.newBuilder()
|
||||
.setMaxAttempts(16)
|
||||
.setInitialBackoff(1)
|
||||
.setMaxBackoff(10)
|
||||
.setBackoffMultiplier(2)
|
||||
.build())
|
||||
.setReceiveBatchSize(ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST)
|
||||
.setLongPollingTimeout(Durations.fromSeconds(30))
|
||||
.build())
|
||||
.build();
|
||||
private static final Map<String, Settings> CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>();
|
||||
|
||||
public Settings getClientSettings(Context ctx) {
|
||||
@@ -37,6 +67,11 @@ public class GrpcClientManager {
|
||||
}
|
||||
|
||||
public void updateClientSettings(String clientId, Settings settings) {
|
||||
if (settings.hasPublishing()) {
|
||||
settings = DEFAULT_PRODUCER_SETTINGS.toBuilder().mergeFrom(settings).build();
|
||||
} else if (settings.hasSubscription()) {
|
||||
settings = DEFAULT_CONSUMER_SETTINGS.toBuilder().mergeFrom(settings).build();
|
||||
}
|
||||
CLIENT_SETTINGS_MAP.put(clientId, settings);
|
||||
}
|
||||
|
||||
|
||||
+1
-2
@@ -124,8 +124,7 @@ public class ConsumerService extends BaseService {
|
||||
|
||||
protected PopMessageRequestHeader buildPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) {
|
||||
checkSubscriptionData(request.getMessageQueue().getTopic(), request.getFilterExpression());
|
||||
// TODO: get fifo config from subscriptionGroupManager
|
||||
boolean fifo = false;
|
||||
boolean fifo = grpcClientManager.getClientSettings(ctx).getSubscription().getFifo();
|
||||
return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx), fifo);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user