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:
@@ -32,7 +32,7 @@ import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.common.Cleaner;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ package org.apache.rocketmq.proxy.common.utils;
|
||||
import java.util.Set;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
|
||||
public class FilterUtil {
|
||||
public class FilterUtils {
|
||||
/**
|
||||
* Whether the message's tag matches consumerGroup's SubscriptionData
|
||||
*
|
||||
@@ -64,12 +64,22 @@ public class DefaultForwardClient extends AbstractForwardClient {
|
||||
return getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> getMaxOffset(String brokerAddr, String topic, int queueId, long timeoutMillis) {
|
||||
public CompletableFuture<Long> getMaxOffset(
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
int queueId,
|
||||
long timeoutMillis
|
||||
) {
|
||||
return getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<Long> searchOffset(String brokerAddr, String topic, int queueId, long timestamp,
|
||||
long timeoutMillis) {
|
||||
public CompletableFuture<Long> searchOffset(
|
||||
String brokerAddr,
|
||||
String topic,
|
||||
int queueId,
|
||||
long timestamp,
|
||||
long timeoutMillis
|
||||
) {
|
||||
return getClient().searchOffset(brokerAddr, topic, queueId, timestamp, timeoutMillis);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -193,7 +193,8 @@ public class Converter {
|
||||
changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId());
|
||||
changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle());
|
||||
changeInvisibleTimeRequestHeader.setOffset(handle.getOffset());
|
||||
changeInvisibleTimeRequestHeader.setInvisibleTime(delayPolicy.getDelayInterval(ConfigurationManager.getProxyConfig().getRetryDelayLevelDelta() + request.getDeliveryAttempt()));
|
||||
changeInvisibleTimeRequestHeader.setInvisibleTime(
|
||||
delayPolicy.getDelayInterval(ConfigurationManager.getProxyConfig().getRetryDelayLevelDelta() + request.getDeliveryAttempt()));
|
||||
return changeInvisibleTimeRequestHeader;
|
||||
}
|
||||
|
||||
|
||||
@@ -23,7 +23,6 @@ import io.grpc.Metadata;
|
||||
import io.grpc.ServerCall;
|
||||
import io.grpc.ServerCallHandler;
|
||||
import io.grpc.ServerInterceptor;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
|
||||
public class ContextInterceptor implements ServerInterceptor {
|
||||
|
||||
|
||||
@@ -25,7 +25,6 @@ import io.grpc.ServerCallHandler;
|
||||
import io.grpc.ServerInterceptor;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.net.SocketAddress;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
|
||||
public class HeaderInterceptor implements ServerInterceptor {
|
||||
@Override
|
||||
|
||||
+1
-1
@@ -15,7 +15,7 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.common;
|
||||
package org.apache.rocketmq.proxy.grpc.interceptor;
|
||||
|
||||
import io.grpc.Context;
|
||||
import io.grpc.Metadata;
|
||||
@@ -98,7 +98,7 @@ import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHand
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseFuture;
|
||||
import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ProxyMode;
|
||||
|
||||
+1
-1
@@ -32,7 +32,7 @@ import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.slf4j.Logger;
|
||||
|
||||
+2
-2
@@ -34,7 +34,7 @@ import org.apache.rocketmq.client.consumer.PullResult;
|
||||
import org.apache.rocketmq.client.consumer.PullStatus;
|
||||
import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
|
||||
import org.apache.rocketmq.proxy.common.utils.FilterUtil;
|
||||
import org.apache.rocketmq.proxy.common.utils.FilterUtils;
|
||||
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
@@ -156,7 +156,7 @@ public class PullMessageService extends BaseService {
|
||||
PullStatus status = result.getPullStatus();
|
||||
if (status.equals(PullStatus.FOUND)) {
|
||||
List<Message> messageList = result.getMsgFoundList().stream()
|
||||
.filter(msg -> FilterUtil.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages.
|
||||
.filter(msg -> FilterUtils.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages.
|
||||
.map(Converter::buildMessage)
|
||||
.collect(Collectors.toList());
|
||||
|
||||
|
||||
@@ -27,25 +27,25 @@ public class FilterUtilTest {
|
||||
@Test
|
||||
public void testIsTagMatched() throws Exception {
|
||||
SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA");
|
||||
assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue();
|
||||
assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIsTagNotMatched() throws Exception {
|
||||
SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA");
|
||||
assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagB")).isFalse();
|
||||
assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagB")).isFalse();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIsTagMatchedStar() throws Exception {
|
||||
SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "*");
|
||||
assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue();
|
||||
assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testIsTagNotMatchedNull() throws Exception {
|
||||
SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA");
|
||||
assertThat(FilterUtil.isTagMatched(subscriptionData.getTagsSet(), null)).isFalse();
|
||||
assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), null)).isFalse();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
+1
-1
@@ -78,7 +78,7 @@ import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader;
|
||||
import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.MessageStore;
|
||||
|
||||
Reference in New Issue
Block a user