diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java index ca215e6777..d12d51953a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java @@ -40,7 +40,6 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.producer.SendResult; -import org.apache.rocketmq.common.attribute.TopicMessageType; import org.apache.rocketmq.common.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageExt; @@ -53,18 +52,14 @@ import org.apache.rocketmq.proxy.grpc.v2.common.GrpcProxyException; import org.apache.rocketmq.proxy.grpc.v2.common.ResponseBuilder; import org.apache.rocketmq.proxy.processor.MessagingProcessor; import org.apache.rocketmq.proxy.processor.QueueSelector; -import org.apache.rocketmq.proxy.processor.validator.DefaultTopicMessageTypeValidator; -import org.apache.rocketmq.proxy.processor.validator.TopicMessageTypeValidator; import org.apache.rocketmq.proxy.service.route.MessageQueueView; import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue; public class SendMessageActivity extends AbstractMessingActivity { - private final TopicMessageTypeValidator validator; public SendMessageActivity(MessagingProcessor messagingProcessor, GrpcClientSettingsManager grpcClientSettingsManager) { super(messagingProcessor, grpcClientSettingsManager); - this.validator = new DefaultTopicMessageTypeValidator(); } public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { @@ -76,9 +71,6 @@ public class SendMessageActivity extends AbstractMessingActivity { throw new GrpcProxyException(Code.MESSAGE_CORRUPTED, "no message to send"); } - MessageType messageType = parseMessageType(request.getMessagesList()); - TopicMessageType topicMessageType = GrpcConverter.buildTopicMessageType(messageType); - List messageList = request.getMessagesList(); Resource topic = messageList.get(0).getTopic(); future = this.messagingProcessor.sendMessage(