[ISSUE #9396] Use fastjson2 in all modules (#9397)

* Use fastjson2 in all modules

* Update test

* Update test

* Update test

* Add serialization compatibility test tool class

* Update RemotingSerializableCompatTest.java

* Update RemotingSerializableCompatTest.java

* Update RemotingSerializableCompatTest.java

* Update BitSet problem

* Update

* Update

* Update test

* Update test

* Update BUILD.bazel

* Update BUILD.bazel

* Update test

* Update BitSet problem

* Add test

* Add compat test

* merge develop

* Update test

* merge develop
This commit is contained in:
yx9o
2025-12-04 19:09:08 +08:00
committed by GitHub
parent 47c6e89589
commit d7e27d6d69
116 changed files with 4209 additions and 595 deletions
-2
View File
@@ -29,7 +29,6 @@ java_library(
"//srvutil",
"@maven//:ch_qos_logback_logback_classic",
"@maven//:ch_qos_logback_logback_core",
"@maven//:com_alibaba_fastjson",
"@maven//:com_alibaba_fastjson2_fastjson2",
"@maven//:com_github_ben_manes_caffeine_caffeine",
"@maven//:com_github_luben_zstd_jni",
@@ -88,7 +87,6 @@ java_library(
"//srvutil",
"//remoting",
"@maven//:ch_qos_logback_logback_core",
"@maven//:com_alibaba_fastjson",
"@maven//:com_alibaba_fastjson2_fastjson2",
"@maven//:com_github_ben_manes_caffeine_caffeine",
"@maven//:com_google_guava_guava",
@@ -17,7 +17,7 @@
package org.apache.rocketmq.proxy.config;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson2.JSON;
import com.google.common.base.Charsets;
import com.google.common.io.CharStreams;
import java.io.File;
@@ -17,8 +17,8 @@
package org.apache.rocketmq.proxy.config;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.serializer.SerializerFeature;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONWriter.Feature;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.auth.config.AuthConfig;
import org.apache.rocketmq.common.MixAll;
@@ -59,6 +59,6 @@ public class ConfigurationManager {
public static String formatProxyConfig() {
return JSON.toJSONString(ConfigurationManager.getProxyConfig(),
SerializerFeature.PrettyFormat, SerializerFeature.WriteMapNullValue, SerializerFeature.WriteDateUseDateFormat, SerializerFeature.WriteNullListAsEmpty);
Feature.PrettyFormat, Feature.WriteMapNullValue, Feature.WriteNullListAsEmpty);
}
}
@@ -17,15 +17,16 @@
package org.apache.rocketmq.proxy.processor.channel;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import java.util.HashMap;
import java.util.Map;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONObject;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.Map;
public class RemoteChannelSerializer {
private static final Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
private static final String REMOTE_PROXY_IP_KEY = "remoteProxyIp";
@@ -17,15 +17,10 @@
package org.apache.rocketmq.proxy.remoting.activity;
import com.alibaba.fastjson.serializer.SerializerFeature;
import com.alibaba.fastjson2.JSONWriter;
import com.google.common.net.HostAndPort;
import io.netty.channel.ChannelHandlerContext;
import java.util.ArrayList;
import java.util.List;
import org.apache.rocketmq.common.MQVersion;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.header.namesrv.GetRouteInfoRequestHeader;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.apache.rocketmq.proxy.common.Address;
import org.apache.rocketmq.proxy.common.ProxyContext;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
@@ -34,6 +29,12 @@ import org.apache.rocketmq.proxy.processor.MessagingProcessor;
import org.apache.rocketmq.proxy.remoting.pipeline.RequestPipeline;
import org.apache.rocketmq.proxy.service.route.ProxyTopicRouteData;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.header.namesrv.GetRouteInfoRequestHeader;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import java.util.ArrayList;
import java.util.List;
public class GetTopicRouteActivity extends AbstractRemotingActivity {
public GetTopicRouteActivity(RequestPipeline requestPipeline,
@@ -57,9 +58,7 @@ public class GetTopicRouteActivity extends AbstractRemotingActivity {
byte[] content;
Boolean standardJsonOnly = requestHeader.getAcceptStandardJsonOnly();
if (request.getVersion() >= MQVersion.Version.V4_9_4.ordinal() || null != standardJsonOnly && standardJsonOnly) {
content = topicRouteData.encode(SerializerFeature.BrowserCompatible,
SerializerFeature.QuoteFieldNames, SerializerFeature.SkipTransientField,
SerializerFeature.MapSortField);
content = topicRouteData.encode(JSONWriter.Feature.BrowserCompatible, JSONWriter.Feature.MapSortField);
} else {
content = topicRouteData.encode();
}
@@ -17,25 +17,22 @@
package org.apache.rocketmq.proxy.remoting.channel;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.TypeReference;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.TypeReference;
import com.google.common.base.MoreObjects;
import io.netty.channel.Channel;
import io.netty.channel.ChannelConfig;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelMetadata;
import java.time.Duration;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.utils.ExceptionUtils;
import org.apache.rocketmq.common.utils.FutureUtils;
import org.apache.rocketmq.common.utils.NetworkUtil;
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
import org.apache.rocketmq.proxy.common.channel.ChannelHelper;
import org.apache.rocketmq.common.utils.ExceptionUtils;
import org.apache.rocketmq.common.utils.FutureUtils;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.processor.channel.ChannelExtendAttributeGetter;
import org.apache.rocketmq.proxy.processor.channel.ChannelProtocolType;
@@ -58,6 +55,10 @@ import org.apache.rocketmq.remoting.protocol.header.ConsumeMessageDirectlyResult
import org.apache.rocketmq.remoting.protocol.header.GetConsumerRunningInfoRequestHeader;
import org.apache.rocketmq.remoting.protocol.heartbeat.SubscriptionData;
import java.time.Duration;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
public class RemotingChannel extends ProxyChannel implements RemoteChannelConverter, ChannelExtendAttributeGetter {
private static final Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
private static final long DEFAULT_MQ_CLIENT_TIMEOUT = Duration.ofSeconds(3).toMillis();
@@ -17,35 +17,36 @@
package org.apache.rocketmq.proxy.service.sysmessage;
import com.alibaba.fastjson.JSON;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import com.alibaba.fastjson2.JSON;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.impl.mqclient.MQClientAPIFactory;
import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.topic.TopicValidator;
import org.apache.rocketmq.common.utils.StartAndShutdown;
import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
import org.apache.rocketmq.proxy.common.ProxyContext;
import org.apache.rocketmq.proxy.common.ProxyException;
import org.apache.rocketmq.proxy.common.ProxyExceptionCode;
import org.apache.rocketmq.common.utils.StartAndShutdown;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.config.ProxyConfig;
import org.apache.rocketmq.proxy.service.admin.AdminService;
import org.apache.rocketmq.client.impl.mqclient.MQClientAPIFactory;
import org.apache.rocketmq.proxy.service.route.AddressableMessageQueue;
import org.apache.rocketmq.proxy.service.route.TopicRouteService;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.header.SendMessageRequestHeader;
import org.apache.rocketmq.remoting.protocol.heartbeat.MessageModel;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
public abstract class AbstractSystemMessageSyncer implements StartAndShutdown, MessageListenerConcurrently {
protected static final Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME);
protected final TopicRouteService topicRouteService;
@@ -17,21 +17,15 @@
package org.apache.rocketmq.proxy.service.sysmessage;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson2.JSON;
import io.netty.channel.Channel;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.broker.client.ClientChannelInfo;
import org.apache.rocketmq.broker.client.ConsumerGroupEvent;
import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener;
import org.apache.rocketmq.broker.client.ConsumerManager;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.impl.mqclient.MQClientAPIFactory;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
@@ -40,13 +34,20 @@ import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.config.ProxyConfig;
import org.apache.rocketmq.proxy.processor.channel.RemoteChannel;
import org.apache.rocketmq.proxy.service.admin.AdminService;
import org.apache.rocketmq.client.impl.mqclient.MQClientAPIFactory;
import org.apache.rocketmq.proxy.service.route.TopicRouteService;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.heartbeat.ConsumeType;
import org.apache.rocketmq.remoting.protocol.heartbeat.MessageModel;
import org.apache.rocketmq.remoting.protocol.heartbeat.SubscriptionData;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class HeartbeatSyncer extends AbstractSystemMessageSyncer {
protected ThreadPoolExecutor threadPoolExecutor;
@@ -21,6 +21,8 @@ import org.apache.rocketmq.proxy.ProxyMode;
import org.junit.Test;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
public class ConfigurationManagerTest extends InitConfigTest {
@@ -44,4 +46,12 @@ public class ConfigurationManagerTest extends InitConfigTest {
assertThat(ConfigurationManager.getProxyConfig()).isNotNull();
}
@Test
public void testFormatProxyConfig() {
String actual = ConfigurationManager.formatProxyConfig();
assertNotNull(actual);
ProxyConfig expected = ConfigurationManager.getProxyConfig();
assertTrue(actual.contains(expected.getProxyMode()));
assertTrue(actual.contains(expected.getProxyName()));
}
}
@@ -0,0 +1,46 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.rocketmq.proxy.config;
import org.apache.rocketmq.auth.config.AuthConfig;
import org.junit.Before;
import org.junit.Test;
import static org.junit.Assert.assertNotNull;
import static org.mockito.Mockito.spy;
public class ConfigurationTest {
private Configuration configuration;
@Before
public void init() {
configuration = spy(new Configuration());
}
@Test
public void testInit() throws Exception {
configuration.init();
ProxyConfig loadedProxyConfig = configuration.getProxyConfig();
assertNotNull(loadedProxyConfig);
AuthConfig loadedAuthConfig = configuration.getAuthConfig();
assertNotNull(loadedAuthConfig);
}
}
@@ -0,0 +1,171 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.rocketmq.proxy.remoting.activity;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONWriter;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import org.apache.rocketmq.common.MQVersion;
import org.apache.rocketmq.proxy.common.ProxyContext;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.processor.MessagingProcessor;
import org.apache.rocketmq.proxy.remoting.pipeline.RequestPipeline;
import org.apache.rocketmq.proxy.service.channel.SimpleChannel;
import org.apache.rocketmq.proxy.service.channel.SimpleChannelHandlerContext;
import org.apache.rocketmq.proxy.service.route.ProxyTopicRouteData;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
import org.apache.rocketmq.remoting.protocol.RequestCode;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.header.namesrv.GetRouteInfoRequestHeader;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.remoting.protocol.route.QueueData;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.junit.MockitoJUnitRunner;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@RunWith(MockitoJUnitRunner.class)
public class GetTopicRouteActivityTest {
@Mock
private RequestPipeline requestPipeline;
@Mock
private MessagingProcessor messagingProcessor;
private GetTopicRouteActivity getTopicRouteActivity;
private ChannelHandlerContext ctx;
private ProxyContext context;
@Before
public void setup() throws Exception {
getTopicRouteActivity = new GetTopicRouteActivity(requestPipeline, messagingProcessor);
ConfigurationManager.initEnv();
ConfigurationManager.intConfig();
Channel channel = new SimpleChannel(null, "0.0.0.0:0", "1.1.1.1:1");
ctx = new SimpleChannelHandlerContext(channel);
context = ProxyContext.create();
}
@Test
public void testProcessRequest0_HighVersion_SerializeWithFeatures() throws Exception {
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_ROUTEINFO_BY_TOPIC, null);
request.setVersion(MQVersion.Version.V4_9_4.ordinal());
GetRouteInfoRequestHeader header = new GetRouteInfoRequestHeader();
header.setTopic("TestTopic");
header.setAcceptStandardJsonOnly(false);
request.writeCustomHeader(header);
TopicRouteData topicRouteData = prepareTopicRouteData();
TopicRouteData spyTopicRouteData = Mockito.spy(topicRouteData);
ProxyTopicRouteData proxyTopicRouteData = mock(ProxyTopicRouteData.class);
when(proxyTopicRouteData.buildTopicRouteData()).thenReturn(spyTopicRouteData);
when(messagingProcessor.getTopicRouteDataForProxy(any(ProxyContext.class), anyList(), any()))
.thenReturn(proxyTopicRouteData);
RemotingCommand response = getTopicRouteActivity.processRequest0(ctx, request, context);
assertNotNull(response);
assertEquals(ResponseCode.SUCCESS, response.getCode());
verify(spyTopicRouteData).encode(
JSONWriter.Feature.BrowserCompatible,
JSONWriter.Feature.MapSortField
);
TopicRouteData deserializedData = JSON.parseObject(response.getBody(), TopicRouteData.class);
assertEquals(topicRouteData.getOrderTopicConf(), deserializedData.getOrderTopicConf());
assertEquals(topicRouteData.getQueueDatas().size(), deserializedData.getQueueDatas().size());
}
@Test
public void testProcessRequest0_LowVersion_StandardJsonOnly_SerializeWithFeatures() throws Exception {
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_ROUTEINFO_BY_TOPIC, null);
request.setVersion(MQVersion.Version.V4_9_3.ordinal());
GetRouteInfoRequestHeader header = new GetRouteInfoRequestHeader();
header.setTopic("TestTopic");
header.setAcceptStandardJsonOnly(true);
request.writeCustomHeader(header);
TopicRouteData topicRouteData = prepareTopicRouteData();
TopicRouteData spyTopicRouteData = Mockito.spy(topicRouteData);
ProxyTopicRouteData proxyTopicRouteData = mock(ProxyTopicRouteData.class);
when(proxyTopicRouteData.buildTopicRouteData()).thenReturn(spyTopicRouteData);
when(messagingProcessor.getTopicRouteDataForProxy(any(ProxyContext.class), anyList(), any()))
.thenReturn(proxyTopicRouteData);
RemotingCommand response = getTopicRouteActivity.processRequest0(ctx, request, context);
assertNotNull(response);
assertEquals(ResponseCode.SUCCESS, response.getCode());
verify(spyTopicRouteData).encode();
}
private TopicRouteData prepareTopicRouteData() {
TopicRouteData result = new TopicRouteData();
result.setOrderTopicConf("orderTopicConf");
List<QueueData> queueDatas = new ArrayList<>();
QueueData queueData = new QueueData();
queueData.setBrokerName("broker-a");
queueData.setPerm(6);
queueData.setReadQueueNums(4);
queueData.setWriteQueueNums(4);
queueData.setTopicSysFlag(0);
queueDatas.add(queueData);
result.setQueueDatas(queueDatas);
List<BrokerData> brokerDatas = new ArrayList<>();
BrokerData brokerData = new BrokerData();
brokerData.setBrokerName("broker-a");
HashMap<Long, String> brokerAddrs = new HashMap<>();
brokerAddrs.put(0L, "127.0.0.1:10911");
brokerData.setBrokerAddrs(brokerAddrs);
brokerDatas.add(brokerData);
result.setBrokerDatas(brokerDatas);
return result;
}
}