From a081ef4b296188e1c1f5cf17862b305c841e8095 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Sun, 24 Apr 2022 10:28:30 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../grpc/v2/service/GrpcClientManager.java | 37 ++++++++++++++++++- .../v2/service/cluster/ConsumerService.java | 3 +- 2 files changed, 37 insertions(+), 3 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java index 2a93676e7b..f784c7fe29 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java @@ -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 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); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index 819349ede5..e0ebe9658f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -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); }