diff --git a/acl/pom.xml b/acl/pom.xml
index c4eac5aac1..9df280713d 100644
--- a/acl/pom.xml
+++ b/acl/pom.xml
@@ -18,6 +18,10 @@
rocketmq-acl
rocketmq-acl ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/acl/src/main/java/org/apache/rocketmq/acl/common/AclUtils.java b/acl/src/main/java/org/apache/rocketmq/acl/common/AclUtils.java
index 69e03523c8..f2c1b40824 100644
--- a/acl/src/main/java/org/apache/rocketmq/acl/common/AclUtils.java
+++ b/acl/src/main/java/org/apache/rocketmq/acl/common/AclUtils.java
@@ -19,7 +19,6 @@ package org.apache.rocketmq.acl.common;
import com.alibaba.fastjson.JSONObject;
import java.io.FileInputStream;
import java.io.FileNotFoundException;
-import java.io.FileWriter;
import java.io.InputStream;
import java.io.PrintWriter;
import java.util.Map;
@@ -256,7 +255,7 @@ public class AclUtils {
public static boolean writeDataObject(String path, Map dataMap) {
Yaml yaml = new Yaml();
- try (PrintWriter pw = new PrintWriter(new FileWriter(path))) {
+ try (PrintWriter pw = new PrintWriter(path, "UTF-8")) {
String dumpAsMap = yaml.dumpAsMap(dataMap);
pw.print(dumpAsMap);
pw.flush();
diff --git a/broker/pom.xml b/broker/pom.xml
index 077ca82f42..b551266cc5 100644
--- a/broker/pom.xml
+++ b/broker/pom.xml
@@ -21,6 +21,10 @@
rocketmq-broker
rocketmq-broker ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/BrokerStartup.java b/broker/src/main/java/org/apache/rocketmq/broker/BrokerStartup.java
index f5fc038d1e..92ace559aa 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/BrokerStartup.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/BrokerStartup.java
@@ -53,7 +53,7 @@ public class BrokerStartup {
public static CommandLine commandLine = null;
public static String configFile = null;
public static InternalLogger log;
- public static SystemConfigFileHelper configFileHelper = new SystemConfigFileHelper();
+ public static final SystemConfigFileHelper configFileHelper = new SystemConfigFileHelper();
public static void main(String[] args) {
start(createBrokerController(args));
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
index ab9463a969..3af82641a2 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java
@@ -22,6 +22,7 @@ import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import java.io.UnsupportedEncodingException;
import java.net.UnknownHostException;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
@@ -374,7 +375,7 @@ public class AdminBrokerProcessor implements NettyRequestProcessor {
groupForbidden.setGroup(group);
groupForbidden.setTopic(topic);
groupForbidden.setReadable(!groupManager.getForbidden(group, topic, PermName.INDEX_PERM_READ));
- response.setBody(groupForbidden.toJson().getBytes());
+ response.setBody(groupForbidden.toJson().getBytes(StandardCharsets.UTF_8));
return response;
}
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
index bb432a8515..a8f93c22bb 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopBufferMergeService.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.broker.processor;
import com.alibaba.fastjson.JSON;
+import java.nio.charset.StandardCharsets;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -605,7 +606,7 @@ public class PopBufferMergeService extends ServiceThread {
PopCheckPoint point = pointWrapper.getCk();
MessageExtBrokerInner msgInner = new MessageExtBrokerInner();
msgInner.setTopic(popMessageProcessor.reviveTopic);
- msgInner.setBody((pointWrapper.getReviveQueueId() + "-" + pointWrapper.getReviveQueueOffset()).getBytes());
+ msgInner.setBody((pointWrapper.getReviveQueueId() + "-" + pointWrapper.getReviveQueueOffset()).getBytes(StandardCharsets.UTF_8));
msgInner.setQueueId(pointWrapper.getReviveQueueId());
msgInner.setTags(PopAckConstants.CK_TAG);
msgInner.setBornTimestamp(System.currentTimeMillis());
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/ReplyMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/ReplyMessageProcessor.java
index 183f64f380..da4d8db1f6 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/processor/ReplyMessageProcessor.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/ReplyMessageProcessor.java
@@ -301,7 +301,7 @@ public class ReplyMessageProcessor extends AbstractSendMessageProcessor {
int commercialBaseCount = brokerController.getBrokerConfig().getCommercialBaseCount();
int wroteSize = putMessageResult.getAppendMessageResult().getWroteBytes();
- int incValue = (int) Math.ceil(wroteSize / commercialSizePerMsg) * commercialBaseCount;
+ int incValue = (int) Math.ceil(wroteSize * 1.0 / commercialSizePerMsg) * commercialBaseCount;
sendMessageContext.setCommercialSendStats(BrokerStatsManager.StatsType.SEND_SUCCESS);
sendMessageContext.setCommercialSendTimes(incValue);
@@ -311,7 +311,7 @@ public class ReplyMessageProcessor extends AbstractSendMessageProcessor {
} else {
if (hasSendMessageHook()) {
int wroteSize = request.getBody().length;
- int incValue = (int) Math.ceil(wroteSize / commercialSizePerMsg);
+ int incValue = (int) Math.ceil(wroteSize * 1.0 / commercialSizePerMsg);
sendMessageContext.setCommercialSendStats(BrokerStatsManager.StatsType.SEND_FAILURE);
sendMessageContext.setCommercialSendTimes(incValue);
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageUtil.java b/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageUtil.java
index 2d6774432e..cdb010482f 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageUtil.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/transaction/queue/TransactionalMessageUtil.java
@@ -24,7 +24,7 @@ import java.nio.charset.StandardCharsets;
public class TransactionalMessageUtil {
public static final String REMOVETAG = "d";
- public static Charset charset = StandardCharsets.UTF_8;
+ public static final Charset charset = StandardCharsets.UTF_8;
public static String buildOpTopic() {
return TopicValidator.RMQ_SYS_TRANS_OP_HALF_TOPIC;
diff --git a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java
index 298442c454..6a09475f6b 100644
--- a/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java
+++ b/broker/src/main/java/org/apache/rocketmq/broker/util/HookUtils.java
@@ -127,7 +127,7 @@ public class HookUtils {
|| tranType == MessageSysFlag.TRANSACTION_COMMIT_TYPE) {
if (!isRolledTimerMessage(msg)) {
if (checkIfTimerMessage(msg)) {
- if (!MessageStoreConfig.isTimerWheelEnable()) {
+ if (!brokerController.getMessageStoreConfig().isTimerWheelEnable()) {
//wheel timer is not enabled, reject the message
return new PutMessageResult(PutMessageStatus.WHEEL_TIMER_NOT_ENABLE, null);
}
diff --git a/client/pom.xml b/client/pom.xml
index 4954db03fa..bf57e6275b 100644
--- a/client/pom.xml
+++ b/client/pom.xml
@@ -27,6 +27,10 @@
rocketmq-client
rocketmq-client ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/client/src/main/java/org/apache/rocketmq/client/common/ThreadLocalIndex.java b/client/src/main/java/org/apache/rocketmq/client/common/ThreadLocalIndex.java
index 41056fac63..b74efd6eba 100644
--- a/client/src/main/java/org/apache/rocketmq/client/common/ThreadLocalIndex.java
+++ b/client/src/main/java/org/apache/rocketmq/client/common/ThreadLocalIndex.java
@@ -27,10 +27,9 @@ public class ThreadLocalIndex {
public int incrementAndGet() {
Integer index = this.threadLocalIndex.get();
if (null == index) {
- index = Math.abs(random.nextInt());
+ index = random.nextInt();
this.threadLocalIndex.set(index);
}
-
this.threadLocalIndex.set(++index);
return Math.abs(index & POSITIVE_MASK);
}
diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java
index dff7efd637..ef763bc99b 100644
--- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java
+++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java
@@ -479,7 +479,6 @@ public abstract class RebalanceImpl {
final boolean isOrder) {
boolean changed = false;
- Map upgradeMqTable = new HashMap();
// drop process queues no longer belong me
HashMap removeQueueMap = new HashMap(this.processQueueTable.size());
Iterator> it = this.processQueueTable.entrySet().iterator();
diff --git a/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java b/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java
index 139d7232f7..9239527345 100644
--- a/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java
+++ b/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.client.trace;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
@@ -394,7 +395,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
* @param traceTopic the topic which message trace data will send to
*/
private void sendTraceDataByMQ(Set keySet, final String data, String traceTopic) {
- final Message message = new Message(traceTopic, data.getBytes());
+ final Message message = new Message(traceTopic, data.getBytes(StandardCharsets.UTF_8));
// Keyset of message trace includes msgId of or original message
message.setKeys(keySet);
try {
diff --git a/common/pom.xml b/common/pom.xml
index fc810f26cd..aa7f9c3372 100644
--- a/common/pom.xml
+++ b/common/pom.xml
@@ -27,6 +27,10 @@
rocketmq-common
rocketmq-common ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/common/src/main/java/org/apache/rocketmq/common/MixAll.java b/common/src/main/java/org/apache/rocketmq/common/MixAll.java
index f8eb4a33af..6e99565898 100644
--- a/common/src/main/java/org/apache/rocketmq/common/MixAll.java
+++ b/common/src/main/java/org/apache/rocketmq/common/MixAll.java
@@ -19,13 +19,13 @@ package org.apache.rocketmq.common;
import org.apache.rocketmq.common.annotation.ImportantField;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.common.help.FAQUrl;
+import org.apache.rocketmq.common.utils.IOTinyUtils;
import org.apache.rocketmq.logging.InternalLogger;
import org.apache.rocketmq.logging.InternalLoggerFactory;
import java.io.ByteArrayInputStream;
import java.io.File;
import java.io.FileInputStream;
-import java.io.FileWriter;
import java.io.IOException;
import java.io.InputStream;
import java.lang.annotation.Annotation;
@@ -184,18 +184,7 @@ public class MixAll {
if (fileParent != null) {
fileParent.mkdirs();
}
- FileWriter fileWriter = null;
-
- try {
- fileWriter = new FileWriter(file);
- fileWriter.write(str);
- } catch (IOException e) {
- throw e;
- } finally {
- if (fileWriter != null) {
- fileWriter.close();
- }
- }
+ IOTinyUtils.writeStringToFile(file, str, "UTF-8");
}
public static String file2String(final String fileName) throws IOException {
diff --git a/common/src/main/java/org/apache/rocketmq/common/consistenthash/ConsistentHashRouter.java b/common/src/main/java/org/apache/rocketmq/common/consistenthash/ConsistentHashRouter.java
index fca1d877d4..33fae11cd9 100644
--- a/common/src/main/java/org/apache/rocketmq/common/consistenthash/ConsistentHashRouter.java
+++ b/common/src/main/java/org/apache/rocketmq/common/consistenthash/ConsistentHashRouter.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.common.consistenthash;
+import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.Collection;
@@ -122,7 +123,7 @@ public class ConsistentHashRouter {
@Override
public long hash(String key) {
instance.reset();
- instance.update(key.getBytes());
+ instance.update(key.getBytes(StandardCharsets.UTF_8));
byte[] digest = instance.digest();
long h = 0;
diff --git a/common/src/main/java/org/apache/rocketmq/common/protocol/header/ExtraInfoUtil.java b/common/src/main/java/org/apache/rocketmq/common/protocol/header/ExtraInfoUtil.java
index 19f37f6cd1..39cbe8b2aa 100644
--- a/common/src/main/java/org/apache/rocketmq/common/protocol/header/ExtraInfoUtil.java
+++ b/common/src/main/java/org/apache/rocketmq/common/protocol/header/ExtraInfoUtil.java
@@ -60,7 +60,7 @@ public class ExtraInfoUtil {
if (extraInfoStrs == null || extraInfoStrs.length < 4) {
throw new IllegalArgumentException("getReviveQid fail, extraInfoStrs length " + (extraInfoStrs == null ? 0 : extraInfoStrs.length));
}
- return Integer.valueOf(extraInfoStrs[3]);
+ return Integer.parseInt(extraInfoStrs[3]);
}
public static String getRealTopic(String[] extraInfoStrs, String topic, String cid) {
@@ -85,14 +85,14 @@ public class ExtraInfoUtil {
if (extraInfoStrs == null || extraInfoStrs.length < 7) {
throw new IllegalArgumentException("getQueueId fail, extraInfoStrs length " + (extraInfoStrs == null ? 0 : extraInfoStrs.length));
}
- return Integer.valueOf(extraInfoStrs[6]);
+ return Integer.parseInt(extraInfoStrs[6]);
}
public static long getQueueOffset(String[] extraInfoStrs) {
if (extraInfoStrs == null || extraInfoStrs.length < 8) {
throw new IllegalArgumentException("getQueueOffset fail, extraInfoStrs length " + (extraInfoStrs == null ? 0 : extraInfoStrs.length));
}
- return Long.valueOf(extraInfoStrs[7]);
+ return Long.parseLong(extraInfoStrs[7]);
}
public static String buildExtraInfo(long ckQueueOffset, long popTime, long invisibleTime, int reviveQid, String topic, String brokerName, int queueId) {
diff --git a/common/src/main/java/org/apache/rocketmq/common/utils/MessageUtils.java b/common/src/main/java/org/apache/rocketmq/common/utils/MessageUtils.java
index a2affc7e90..0e5ac7add9 100644
--- a/common/src/main/java/org/apache/rocketmq/common/utils/MessageUtils.java
+++ b/common/src/main/java/org/apache/rocketmq/common/utils/MessageUtils.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.common.utils;
+import java.nio.charset.StandardCharsets;
import java.util.Collection;
import java.util.HashSet;
import java.util.Set;
@@ -27,7 +28,7 @@ import org.apache.rocketmq.common.message.MessageExt;
public class MessageUtils {
public static int getShardingKeyIndex(String shardingKey, int indexSize) {
- return Math.abs(Hashing.murmur3_32().hashBytes(shardingKey.getBytes()).asInt() % indexSize);
+ return Math.abs(Hashing.murmur3_32().hashBytes(shardingKey.getBytes(StandardCharsets.UTF_8)).asInt() % indexSize);
}
public static int getShardingKeyIndexByMsg(MessageExt msg, int indexSize) {
diff --git a/container/pom.xml b/container/pom.xml
index 105862e17f..d62b957acb 100644
--- a/container/pom.xml
+++ b/container/pom.xml
@@ -27,6 +27,10 @@
rocketmq-container
rocketmq-container ${project.version}
+
+ ${basedir}/..
+
+
org.apache.rocketmq
diff --git a/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java b/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java
index d4e94a6984..3d78275a10 100644
--- a/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java
+++ b/container/src/main/java/org/apache/rocketmq/container/BrokerContainerStartup.java
@@ -59,9 +59,9 @@ public class BrokerContainerStartup {
public static CommandLine commandLine = null;
public static String configFile = null;
public static InternalLogger log;
- public static SystemConfigFileHelper configFileHelper = new SystemConfigFileHelper();
+ public static final SystemConfigFileHelper configFileHelper = new SystemConfigFileHelper();
public static String rocketmqHome = null;
- public static JoranConfigurator configurator = new JoranConfigurator();
+ public static final JoranConfigurator configurator = new JoranConfigurator();
public static void main(String[] args) {
final BrokerContainer brokerContainer = startBrokerContainer(createBrokerContainer(args));
diff --git a/controller/pom.xml b/controller/pom.xml
index 04dabb0ff9..dc0e7a313a 100644
--- a/controller/pom.xml
+++ b/controller/pom.xml
@@ -26,6 +26,10 @@
rocketmq-controller
rocketmq-controller ${project.version}
+
+ ${basedir}/..
+
+
io.openmessaging.storage
diff --git a/distribution/pom.xml b/distribution/pom.xml
index ce27e43dea..4ce8f520ab 100644
--- a/distribution/pom.xml
+++ b/distribution/pom.xml
@@ -26,6 +26,10 @@
rocketmq-distribution ${project.version}
pom
+
+ ${basedir}/..
+
+
release-all
diff --git a/example/pom.xml b/example/pom.xml
index 5f9f4ffb7a..3849fa00af 100644
--- a/example/pom.xml
+++ b/example/pom.xml
@@ -27,6 +27,10 @@
rocketmq-example
rocketmq-example ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java b/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java
index 30863032ef..cf82c2a874 100644
--- a/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/batch/SimpleBatchProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.example.batch;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
@@ -39,9 +40,9 @@ public class SimpleBatchProducer {
//If you just send messages of no more than 1MiB at a time, it is easy to use batch
//Messages of the same batch should have: same topic, same waitStoreMsgOK and no schedule support
List messages = new ArrayList<>();
- messages.add(new Message(TOPIC, TAG, "OrderID001", "Hello world 0".getBytes()));
- messages.add(new Message(TOPIC, TAG, "OrderID002", "Hello world 1".getBytes()));
- messages.add(new Message(TOPIC, TAG, "OrderID003", "Hello world 2".getBytes()));
+ messages.add(new Message(TOPIC, TAG, "OrderID001", "Hello world 0".getBytes(StandardCharsets.UTF_8)));
+ messages.add(new Message(TOPIC, TAG, "OrderID002", "Hello world 1".getBytes(StandardCharsets.UTF_8)));
+ messages.add(new Message(TOPIC, TAG, "OrderID003", "Hello world 2".getBytes(StandardCharsets.UTF_8)));
SendResult sendResult = producer.send(messages);
System.out.printf("%s", sendResult);
diff --git a/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java b/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java
index aca4f16837..d33a5a5bba 100644
--- a/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/batch/SplitBatchProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.example.batch;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
@@ -44,7 +45,7 @@ public class SplitBatchProducer {
//large batch
List messages = new ArrayList<>(MESSAGE_COUNT);
for (int i = 0; i < MESSAGE_COUNT; i++) {
- messages.add(new Message(TOPIC, TAG, "OrderID" + i, ("Hello world " + i).getBytes()));
+ messages.add(new Message(TOPIC, TAG, "OrderID" + i, ("Hello world " + i).getBytes(StandardCharsets.UTF_8)));
}
//split the large batch into small ones:
diff --git a/example/src/main/java/org/apache/rocketmq/example/namespace/ProducerWithNamespace.java b/example/src/main/java/org/apache/rocketmq/example/namespace/ProducerWithNamespace.java
index 66c7ec0964..ccbd6d0de9 100644
--- a/example/src/main/java/org/apache/rocketmq/example/namespace/ProducerWithNamespace.java
+++ b/example/src/main/java/org/apache/rocketmq/example/namespace/ProducerWithNamespace.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.example.namespace;
+import java.nio.charset.StandardCharsets;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
@@ -36,7 +37,7 @@ public class ProducerWithNamespace {
producer.setNamesrvAddr(DEFAULT_NAMESRVADDR);
producer.start();
for (int i = 0; i < MESSAGE_COUNT; i++) {
- Message message = new Message(TOPIC, TAG, "Hello world".getBytes());
+ Message message = new Message(TOPIC, TAG, "Hello world".getBytes(StandardCharsets.UTF_8));
try {
SendResult result = producer.send(message);
System.out.printf("Topic:%s send success, misId is:%s%n", message.getTopic(), result.getMsgId());
diff --git a/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimpleProducer.java b/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimpleProducer.java
index 803faaa236..ec72aa0dd7 100644
--- a/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimpleProducer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimpleProducer.java
@@ -22,6 +22,7 @@ import io.openmessaging.MessagingAccessPoint;
import io.openmessaging.OMS;
import io.openmessaging.producer.Producer;
import io.openmessaging.producer.SendResult;
+import java.nio.charset.StandardCharsets;
import java.util.concurrent.CountDownLatch;
public class SimpleProducer {
@@ -43,7 +44,7 @@ public class SimpleProducer {
System.out.printf("Producer startup OK%n");
{
- Message message = producer.createBytesMessage(QUEUE, "OMS_HELLO_BODY".getBytes());
+ Message message = producer.createBytesMessage(QUEUE, "OMS_HELLO_BODY".getBytes(StandardCharsets.UTF_8));
SendResult sendResult = producer.send(message);
//final Void aVoid = result.get(3000L);
@@ -52,7 +53,8 @@ public class SimpleProducer {
final CountDownLatch countDownLatch = new CountDownLatch(1);
{
- final Future result = producer.sendAsync(producer.createBytesMessage(QUEUE, "OMS_HELLO_BODY".getBytes()));
+ final Future result = producer.sendAsync(producer.createBytesMessage(QUEUE,
+ "OMS_HELLO_BODY".getBytes(StandardCharsets.UTF_8)));
result.addListener(future -> {
if (future.getThrowable() != null) {
System.out.printf("Send async message Failed, error: %s%n", future.getThrowable().getMessage());
@@ -64,7 +66,7 @@ public class SimpleProducer {
}
{
- producer.sendOneway(producer.createBytesMessage("OMS_HELLO_TOPIC", "OMS_HELLO_BODY".getBytes()));
+ producer.sendOneway(producer.createBytesMessage("OMS_HELLO_TOPIC", "OMS_HELLO_BODY".getBytes(StandardCharsets.UTF_8)));
System.out.printf("Send oneway message OK%n");
}
diff --git a/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimplePullConsumer.java b/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimplePullConsumer.java
index 2c0059ab8d..9ad69b31b3 100644
--- a/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimplePullConsumer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/openmessaging/SimplePullConsumer.java
@@ -23,6 +23,7 @@ import io.openmessaging.OMSBuiltinKeys;
import io.openmessaging.consumer.PullConsumer;
import io.openmessaging.producer.Producer;
import io.openmessaging.producer.SendResult;
+import java.nio.charset.StandardCharsets;
public class SimplePullConsumer {
@@ -48,7 +49,7 @@ public class SimplePullConsumer {
final String queueName = "TopicTest";
producer.startup();
- Message msg = producer.createBytesMessage(queueName, "Hello Open Messaging".getBytes());
+ Message msg = producer.createBytesMessage(queueName, "Hello Open Messaging".getBytes(StandardCharsets.UTF_8));
SendResult sendResult = producer.send(msg);
System.out.printf("Send Message OK. MsgId: %s%n", sendResult.messageId());
producer.shutdown();
diff --git a/example/src/main/java/org/apache/rocketmq/example/rpc/ResponseConsumer.java b/example/src/main/java/org/apache/rocketmq/example/rpc/ResponseConsumer.java
index 421297cc3f..a1c18ae698 100644
--- a/example/src/main/java/org/apache/rocketmq/example/rpc/ResponseConsumer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/rpc/ResponseConsumer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.example.rpc;
+import java.nio.charset.StandardCharsets;
import java.util.List;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
@@ -59,7 +60,7 @@ public class ResponseConsumer {
try {
System.out.printf("handle message: %s %n", msg.toString());
String replyTo = MessageUtil.getReplyToClient(msg);
- byte[] replyContent = "reply message contents.".getBytes();
+ byte[] replyContent = "reply message contents.".getBytes(StandardCharsets.UTF_8);
// create reply message with given util, do not create reply message by yourself
Message replyMessage = MessageUtil.createReplyMessage(msg, replyContent);
diff --git a/example/src/main/java/org/apache/rocketmq/example/schedule/ScheduledMessageProducer.java b/example/src/main/java/org/apache/rocketmq/example/schedule/ScheduledMessageProducer.java
index 994c81e64f..aeae492dd9 100644
--- a/example/src/main/java/org/apache/rocketmq/example/schedule/ScheduledMessageProducer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/schedule/ScheduledMessageProducer.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.example.schedule;
+import java.nio.charset.StandardCharsets;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
@@ -37,7 +38,7 @@ public class ScheduledMessageProducer {
producer.start();
int totalMessagesToSend = 100;
for (int i = 0; i < totalMessagesToSend; i++) {
- Message message = new Message(TOPIC, ("Hello scheduled message " + i).getBytes());
+ Message message = new Message(TOPIC, ("Hello scheduled message " + i).getBytes(StandardCharsets.UTF_8));
// This message will be delivered to consumer 10 seconds later.
message.setDelayTimeLevel(3);
// Send the message
diff --git a/example/src/main/java/org/apache/rocketmq/example/schedule/TimerMessageProducer.java b/example/src/main/java/org/apache/rocketmq/example/schedule/TimerMessageProducer.java
index 18a1b61f5b..baa8da7c82 100644
--- a/example/src/main/java/org/apache/rocketmq/example/schedule/TimerMessageProducer.java
+++ b/example/src/main/java/org/apache/rocketmq/example/schedule/TimerMessageProducer.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.example.schedule;
+import java.nio.charset.StandardCharsets;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
@@ -37,7 +38,7 @@ public class TimerMessageProducer {
producer.start();
int totalMessagesToSend = 10;
for (int i = 0; i < totalMessagesToSend; i++) {
- Message message = new Message(TOPIC, ("Hello scheduled message " + i).getBytes());
+ Message message = new Message(TOPIC, ("Hello scheduled message " + i).getBytes(StandardCharsets.UTF_8));
// This message will be delivered to consumer 10 seconds later.
//message.setDelayTimeSec(10);
// The effect is the same as the above
diff --git a/example/src/main/java/org/apache/rocketmq/example/simple/AclClient.java b/example/src/main/java/org/apache/rocketmq/example/simple/AclClient.java
index 0c97cd3321..c1f56ec9c2 100644
--- a/example/src/main/java/org/apache/rocketmq/example/simple/AclClient.java
+++ b/example/src/main/java/org/apache/rocketmq/example/simple/AclClient.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.example.simple;
+import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -145,7 +146,7 @@ public class AclClient {
return;
for (MessageExt m : msg) {
if (m != null) {
- System.out.printf("msgId : %s body : %s \n\r", m.getMsgId(), new String(m.getBody()));
+ System.out.printf("msgId : %s body : %s \n\r", m.getMsgId(), new String(m.getBody(), StandardCharsets.UTF_8));
}
}
}
diff --git a/filter/pom.xml b/filter/pom.xml
index 5c5080f398..6b6791f602 100644
--- a/filter/pom.xml
+++ b/filter/pom.xml
@@ -28,6 +28,10 @@
rocketmq-filter
rocketmq-filter ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/filter/src/main/java/org/apache/rocketmq/filter/parser/SimpleCharStream.java b/filter/src/main/java/org/apache/rocketmq/filter/parser/SimpleCharStream.java
index c10250e3bb..42626f0f23 100644
--- a/filter/src/main/java/org/apache/rocketmq/filter/parser/SimpleCharStream.java
+++ b/filter/src/main/java/org/apache/rocketmq/filter/parser/SimpleCharStream.java
@@ -19,6 +19,8 @@
/* JavaCCOptions:STATIC=false,SUPPORT_CLASS_VISIBILITY_PUBLIC=true */
package org.apache.rocketmq.filter.parser;
+import java.nio.charset.StandardCharsets;
+
/**
* An implementation of interface CharStream, where the stream is assumed to
* contain only ASCII characters (without unicode processing).
@@ -331,7 +333,7 @@ public class SimpleCharStream {
public SimpleCharStream(java.io.InputStream dstream, String encoding, int startline,
int startcolumn, int buffersize) throws java.io.UnsupportedEncodingException {
this(encoding == null ?
- new java.io.InputStreamReader(dstream) :
+ new java.io.InputStreamReader(dstream, StandardCharsets.UTF_8) :
new java.io.InputStreamReader(dstream, encoding), startline, startcolumn, buffersize);
}
@@ -340,7 +342,7 @@ public class SimpleCharStream {
*/
public SimpleCharStream(java.io.InputStream dstream, int startline,
int startcolumn, int buffersize) {
- this(new java.io.InputStreamReader(dstream), startline, startcolumn, buffersize);
+ this(new java.io.InputStreamReader(dstream, StandardCharsets.UTF_8), startline, startcolumn, buffersize);
}
/**
@@ -379,7 +381,7 @@ public class SimpleCharStream {
public void ReInit(java.io.InputStream dstream, String encoding, int startline,
int startcolumn, int buffersize) throws java.io.UnsupportedEncodingException {
ReInit(encoding == null ?
- new java.io.InputStreamReader(dstream) :
+ new java.io.InputStreamReader(dstream, StandardCharsets.UTF_8) :
new java.io.InputStreamReader(dstream, encoding), startline, startcolumn, buffersize);
}
@@ -388,7 +390,7 @@ public class SimpleCharStream {
*/
public void ReInit(java.io.InputStream dstream, int startline,
int startcolumn, int buffersize) {
- ReInit(new java.io.InputStreamReader(dstream), startline, startcolumn, buffersize);
+ ReInit(new java.io.InputStreamReader(dstream, StandardCharsets.UTF_8), startline, startcolumn, buffersize);
}
/**
diff --git a/logging/pom.xml b/logging/pom.xml
index 6a1572e720..4df7df5fbc 100644
--- a/logging/pom.xml
+++ b/logging/pom.xml
@@ -27,6 +27,10 @@
rocketmq-logging
rocketmq-logging ${project.version}
+
+ ${basedir}/..
+
+
org.slf4j
diff --git a/logging/src/main/java/org/apache/rocketmq/logging/InternalLoggerFactory.java b/logging/src/main/java/org/apache/rocketmq/logging/InternalLoggerFactory.java
index e4b57328cc..2370a911ed 100644
--- a/logging/src/main/java/org/apache/rocketmq/logging/InternalLoggerFactory.java
+++ b/logging/src/main/java/org/apache/rocketmq/logging/InternalLoggerFactory.java
@@ -38,7 +38,7 @@ public abstract class InternalLoggerFactory {
private static String loggerType = null;
- public static ThreadLocal brokerIdentity = new ThreadLocal();
+ public static final ThreadLocal brokerIdentity = new ThreadLocal();
private static ConcurrentHashMap loggerFactoryCache = new ConcurrentHashMap();
diff --git a/logging/src/main/java/org/apache/rocketmq/logging/inner/LoggingBuilder.java b/logging/src/main/java/org/apache/rocketmq/logging/inner/LoggingBuilder.java
index 7468cd498f..996551e77d 100644
--- a/logging/src/main/java/org/apache/rocketmq/logging/inner/LoggingBuilder.java
+++ b/logging/src/main/java/org/apache/rocketmq/logging/inner/LoggingBuilder.java
@@ -27,6 +27,7 @@ import java.io.InterruptedIOException;
import java.io.OutputStream;
import java.io.OutputStreamWriter;
import java.io.Writer;
+import java.nio.charset.StandardCharsets;
import java.text.MessageFormat;
import java.text.SimpleDateFormat;
import java.util.ArrayList;
@@ -534,7 +535,7 @@ public class LoggingBuilder {
}
}
if (retval == null) {
- retval = new OutputStreamWriter(os);
+ retval = new OutputStreamWriter(os, StandardCharsets.UTF_8);
}
return retval;
}
diff --git a/namesrv/pom.xml b/namesrv/pom.xml
index e6c96a862c..5334204f76 100644
--- a/namesrv/pom.xml
+++ b/namesrv/pom.xml
@@ -27,6 +27,10 @@
rocketmq-namesrv
rocketmq-namesrv ${project.version}
+
+ ${basedir}/..
+
+
${project.groupId}
diff --git a/openmessaging/pom.xml b/openmessaging/pom.xml
index 071fede40e..2481a6ac15 100644
--- a/openmessaging/pom.xml
+++ b/openmessaging/pom.xml
@@ -27,6 +27,10 @@
rocketmq-openmessaging
rocketmq-openmessaging ${project.version}
+
+ ${basedir}/..
+
+
io.openmessaging
diff --git a/pom.xml b/pom.xml
index 6f53a33cb4..25d4d4d1a8 100644
--- a/pom.xml
+++ b/pom.xml
@@ -91,6 +91,7 @@
UTF-8
UTF-8
+ ${basedir}
false
@@ -152,8 +153,8 @@
4.3.0
0.8.5
2.19.1
- 3.0.4
3.0.2
+ 4.2.2
3.0.0
2.10.4
2.19.1
@@ -380,16 +381,33 @@
-
- org.codehaus.mojo
- findbugs-maven-plugin
- ${findbugs-maven-plugin.version}
-
org.sonarsource.scanner.maven
sonar-maven-plugin
${sonar-maven-plugin.version}
+
+ com.github.spotbugs
+ spotbugs-maven-plugin
+ ${spotbugs-plugin.version}
+
+
+ check
+ compile
+
+ check
+
+
+
+
+ true
+ false
+ true
+ ${project.root}/style/spotbugs-suppressions.xml
+ High
+ Max
+
+
diff --git a/proxy/pom.xml b/proxy/pom.xml
index a72f30c820..f8978cc057 100644
--- a/proxy/pom.xml
+++ b/proxy/pom.xml
@@ -33,6 +33,7 @@
8
8
+ ${basedir}/..
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java
index aaba6ef9cb..59078c7129 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java
@@ -23,6 +23,7 @@ import com.google.common.io.CharStreams;
import java.io.File;
import java.io.InputStream;
import java.io.InputStreamReader;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.rocketmq.common.constant.LoggerName;
@@ -65,7 +66,7 @@ public class Configuration {
return null;
}
- return new String(Files.readAllBytes(file.toPath()));
+ return new String(Files.readAllBytes(file.toPath()), StandardCharsets.UTF_8);
}
public ProxyConfig getProxyConfig() {
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java
index 3943b33927..2e6d4f67e4 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java
@@ -194,9 +194,6 @@ public class GrpcClientSettingsManager {
settings = mergeSubscriptionData(ctx, settings,
GrpcConverter.getInstance().wrapResourceWithNamespace(settings.getSubscription().getGroup()));
}
- if (settings == null) {
- return null;
- }
return mergeMetric(settings);
}
}
diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java
index cf5314c79f..a05eedd508 100644
--- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java
+++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.proxy.service.route;
import com.google.common.base.MoreObjects;
+import com.google.common.math.IntMath;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
@@ -51,9 +52,9 @@ public class MessageQueueSelector {
this.queues.addAll(buildWrite(topicRouteWrapper));
}
buildBrokerActingQueues(topicRouteWrapper.getTopicName(), this.queues);
-
- this.queueIndex = new AtomicInteger(Math.abs(new Random().nextInt()));
- this.brokerIndex = new AtomicInteger(Math.abs(new Random().nextInt()));
+ Random random = new Random();
+ this.queueIndex = new AtomicInteger(random.nextInt());
+ this.brokerIndex = new AtomicInteger(random.nextInt());
}
private static List buildRead(TopicRouteWrapper topicRoute) {
@@ -172,13 +173,13 @@ public class MessageQueueSelector {
if (brokerActingQueues.isEmpty()) {
return null;
}
- return brokerActingQueues.get(Math.abs(index) % brokerActingQueues.size());
+ return brokerActingQueues.get(IntMath.mod(index, brokerActingQueues.size()));
}
if (queues.isEmpty()) {
return null;
}
- return queues.get(Math.abs(index) % queues.size());
+ return queues.get(IntMath.mod(index, queues.size()));
}
public List getQueues() {
diff --git a/remoting/pom.xml b/remoting/pom.xml
index 484fa95f61..f567d84eab 100644
--- a/remoting/pom.xml
+++ b/remoting/pom.xml
@@ -27,6 +27,10 @@
rocketmq-remoting
rocketmq-remoting ${project.version}
+
+ ${basedir}/..
+
+
com.alibaba
diff --git a/srvutil/pom.xml b/srvutil/pom.xml
index a59e00894c..99ee96bce6 100644
--- a/srvutil/pom.xml
+++ b/srvutil/pom.xml
@@ -27,6 +27,9 @@
rocketmq-srvutil
rocketmq-srvutil ${project.version}
+
+ ${basedir}/..
+
diff --git a/store/pom.xml b/store/pom.xml
index 332a5f51fe..03f6aa93d6 100644
--- a/store/pom.xml
+++ b/store/pom.xml
@@ -27,6 +27,10 @@
rocketmq-store
rocketmq-store ${project.version}
+
+ ${basedir}/..
+
+
io.openmessaging.storage
diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
index c48d066543..d3e5ef06c6 100644
--- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
+++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java
@@ -26,6 +26,7 @@ import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.FileLock;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
@@ -327,7 +328,7 @@ public class DefaultMessageStore implements MessageStore {
throw new RuntimeException("Lock failed,MQ already started");
}
- lockFile.getChannel().write(ByteBuffer.wrap("lock".getBytes()));
+ lockFile.getChannel().write(ByteBuffer.wrap("lock".getBytes(StandardCharsets.UTF_8)));
lockFile.getChannel().force(true);
if (this.getMessageStoreConfig().isDuplicationEnable()) {
diff --git a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
index 9b69e75d91..1d5b733364 100644
--- a/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
+++ b/store/src/main/java/org/apache/rocketmq/store/config/MessageStoreConfig.java
@@ -62,7 +62,7 @@ public class MessageStoreConfig {
private boolean timerEnableCheckMetrics = true;
private boolean timerInterceptDelayLevel = false;
private int timerMaxDelaySec = 3600 * 24 * 3;
- private static boolean timerWheelEnable = true;
+ private boolean timerWheelEnable = true;
/**
* 1. Register to broker after (startTime + disappearTimeAfterStart)
@@ -1441,7 +1441,7 @@ public class MessageStoreConfig {
return timerWarmEnable;
}
- public static boolean isTimerWheelEnable() {
+ public boolean isTimerWheelEnable() {
return timerWheelEnable;
}
public void setTimerWheelEnable(boolean timerWheelEnable) {
diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
index f41c0a81b4..53903a1df7 100644
--- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
+++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java
@@ -23,6 +23,7 @@ import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicReference;
@@ -268,7 +269,7 @@ public class AutoSwitchHAClient extends ServiceThread implements HAClient {
// Address length
this.handshakeHeaderBuffer.putInt(this.localAddress == null ? 0 : this.localAddress.length());
// Slave address
- this.handshakeHeaderBuffer.put(this.localAddress == null ? new byte[0] : this.localAddress.getBytes());
+ this.handshakeHeaderBuffer.put(this.localAddress == null ? new byte[0] : this.localAddress.getBytes(StandardCharsets.UTF_8));
this.handshakeHeaderBuffer.flip();
return this.haWriter.write(this.socketChannel, this.handshakeHeaderBuffer);
diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
index 98fd778916..63e4b6b9e5 100644
--- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
+++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java
@@ -22,6 +22,7 @@ import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
+import java.nio.charset.StandardCharsets;
import java.util.List;
import org.apache.rocketmq.common.EpochEntry;
import org.apache.rocketmq.common.ServiceThread;
@@ -313,7 +314,7 @@ public class AutoSwitchHAConnection implements HAConnection {
final byte[] addressData = new byte[addressLength];
byteBufferRead.position(readPosition + AutoSwitchHAClient.HANDSHAKE_HEADER_SIZE);
byteBufferRead.get(addressData);
- AutoSwitchHAConnection.this.slaveAddress = new String(addressData);
+ AutoSwitchHAConnection.this.slaveAddress = new String(addressData, StandardCharsets.UTF_8);
isSlaveSendHandshake = true;
byteBufferRead.position(readSocketPos);
diff --git a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMetrics.java b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMetrics.java
index 1d963aea92..b82569b899 100644
--- a/store/src/main/java/org/apache/rocketmq/store/timer/TimerMetrics.java
+++ b/store/src/main/java/org/apache/rocketmq/store/timer/TimerMetrics.java
@@ -19,6 +19,10 @@ package org.apache.rocketmq.store.timer;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.serializer.SerializerFeature;
import com.google.common.io.Files;
+import java.io.FileOutputStream;
+import java.io.OutputStreamWriter;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Paths;
import org.apache.rocketmq.common.ConfigManager;
import org.apache.rocketmq.common.DataVersion;
import org.apache.rocketmq.common.constant.LoggerName;
@@ -86,18 +90,26 @@ public class TimerMetrics extends ConfigManager {
public Metric getDistPair(Integer period) {
Metric pair = timingDistribution.get(period);
- if (null == pair) {
- pair = new Metric();
- timingDistribution.putIfAbsent(period, pair);
+ if (null != pair) {
+ return pair;
+ }
+ pair = new Metric();
+ final Metric previous = timingDistribution.putIfAbsent(period, pair);
+ if (null != previous) {
+ return previous;
}
return pair;
}
public Metric getTopicPair(String topic) {
Metric pair = timingCount.get(topic);
- if (null == pair) {
- pair = new Metric();
- timingCount.putIfAbsent(topic, pair);
+ if (null != pair) {
+ return pair;
+ }
+ pair = new Metric();
+ final Metric previous = timingCount.putIfAbsent(topic, pair);
+ if (null != previous) {
+ return previous;
}
return pair;
}
@@ -110,7 +122,6 @@ public class TimerMetrics extends ConfigManager {
this.timerDist = timerDist;
}
-
public long getTimingCount(String topic) {
Metric pair = timingCount.get(topic);
if (null == pair) {
@@ -226,7 +237,8 @@ public class TimerMetrics extends ConfigManager {
return;
}
}
- bufferedWriter = new BufferedWriter(new FileWriter(tmpFile, false));
+ bufferedWriter = new BufferedWriter(new OutputStreamWriter(new FileOutputStream(tmpFile, false),
+ StandardCharsets.UTF_8));
write0(bufferedWriter);
bufferedWriter.flush();
log.debug("Finished writing tmp file: {}", temp);
diff --git a/style/spotbugs-suppressions.xml b/style/spotbugs-suppressions.xml
new file mode 100644
index 0000000000..607080cfbd
--- /dev/null
+++ b/style/spotbugs-suppressions.xml
@@ -0,0 +1,43 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/test/pom.xml b/test/pom.xml
index 2e5863775e..bb6d28c06e 100644
--- a/test/pom.xml
+++ b/test/pom.xml
@@ -27,6 +27,10 @@
rocketmq-test
rocketmq-test ${project.version}
+
+ ${basedir}/..
+
+
log4j
diff --git a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQAsyncSendProducer.java b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQAsyncSendProducer.java
index 9907cac800..d28a5fd3aa 100644
--- a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQAsyncSendProducer.java
+++ b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQAsyncSendProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.test.client.rmq;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
@@ -108,7 +109,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.send(metaqMsg, sendCallback);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -119,7 +120,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
this.msgSize = msgSize;
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.asyncSend(msg);
}
}
@@ -128,7 +129,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.send(metaqMsg, selector, arg, sendCallback);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -138,7 +139,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
public void asyncSend(int msgSize, MessageQueueSelector selector) {
this.msgSize = msgSize;
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.asyncSend(msg, selector, i);
}
}
@@ -147,7 +148,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.send(metaqMsg, mq, sendCallback);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -157,7 +158,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
public void asyncSend(int msgSize, MessageQueue mq) {
this.msgSize = msgSize;
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.asyncSend(msg, mq);
}
}
@@ -178,7 +179,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.sendOneway(metaqMsg);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -187,7 +188,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
public void sendOneWay(int msgSize) {
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.sendOneWay(msg);
}
}
@@ -196,7 +197,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.sendOneway(metaqMsg, mq);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -205,7 +206,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
public void sendOneWay(int msgSize, MessageQueue mq) {
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.sendOneWay(msg, mq);
}
}
@@ -214,7 +215,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
Message metaqMsg = (Message) msg;
try {
producer.sendOneway(metaqMsg, selector, arg);
- msgBodys.addData(new String(metaqMsg.getBody()));
+ msgBodys.addData(new String(metaqMsg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
} catch (Exception e) {
e.printStackTrace();
@@ -223,7 +224,7 @@ public class RMQAsyncSendProducer extends AbstractMQProducer {
public void sendOneWay(int msgSize, MessageQueueSelector selector) {
for (int i = 0; i < msgSize; i++) {
- Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes());
+ Message msg = new Message(topic, RandomUtil.getStringByUUID().getBytes(StandardCharsets.UTF_8));
this.sendOneWay(msg, selector, i);
}
}
diff --git a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQNormalProducer.java b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQNormalProducer.java
index 001db9584e..eb8cf44be9 100644
--- a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQNormalProducer.java
+++ b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQNormalProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.test.client.rmq;
+import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Map;
import org.apache.log4j.Logger;
@@ -105,9 +106,9 @@ public class RMQNormalProducer extends AbstractMQProducer {
sendResult.setMsgId(metaqResult.getMsgId());
sendResult.setSendResult(metaqResult.getSendStatus().equals(SendStatus.SEND_OK));
sendResult.setBrokerIp(metaqResult.getMessageQueue().getBrokerName());
- msgBodys.addData(new String(message.getBody()));
+ msgBodys.addData(new String(message.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
- originMsgIndex.put(new String(message.getBody()), metaqResult);
+ originMsgIndex.put(new String(message.getBody(), StandardCharsets.UTF_8), metaqResult);
} catch (Exception e) {
if (isDebug) {
e.printStackTrace();
@@ -151,9 +152,9 @@ public class RMQNormalProducer extends AbstractMQProducer {
sendResult.setMsgId(metaqResult.getMsgId());
sendResult.setSendResult(metaqResult.getSendStatus().equals(SendStatus.SEND_OK));
sendResult.setBrokerIp(metaqResult.getMessageQueue().getBrokerName());
- msgBodys.addData(new String(msg.getBody()));
+ msgBodys.addData(new String(msg.getBody(), StandardCharsets.UTF_8));
originMsgs.addData(msg);
- originMsgIndex.put(new String(msg.getBody()), metaqResult);
+ originMsgIndex.put(new String(msg.getBody(), StandardCharsets.UTF_8), metaqResult);
} catch (Exception e) {
if (isDebug) {
e.printStackTrace();
diff --git a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQTransactionalProducer.java b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQTransactionalProducer.java
index dcc76b2d84..69563e0e10 100644
--- a/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQTransactionalProducer.java
+++ b/test/src/main/java/org/apache/rocketmq/test/client/rmq/RMQTransactionalProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.test.client.rmq;
+import java.nio.charset.StandardCharsets;
import org.apache.log4j.Logger;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.LocalTransactionState;
@@ -82,10 +83,10 @@ public class RMQTransactionalProducer extends AbstractMQProducer {
sendResult.setSendResult(true);
sendResult.setBrokerIp(metaqResult.getMessageQueue().getBrokerName());
if (commitMsg) {
- msgBodys.addData(new String(message.getBody()));
+ msgBodys.addData(new String(message.getBody(), StandardCharsets.UTF_8));
}
originMsgs.addData(msg);
- originMsgIndex.put(new String(message.getBody()), metaqResult);
+ originMsgIndex.put(new String(message.getBody(), StandardCharsets.UTF_8), metaqResult);
} catch (MQClientException e) {
if (isDebug) {
e.printStackTrace();
diff --git a/test/src/main/java/org/apache/rocketmq/test/clientinterface/AbstractMQProducer.java b/test/src/main/java/org/apache/rocketmq/test/clientinterface/AbstractMQProducer.java
index 1e0a19a2b3..258000b090 100644
--- a/test/src/main/java/org/apache/rocketmq/test/clientinterface/AbstractMQProducer.java
+++ b/test/src/main/java/org/apache/rocketmq/test/clientinterface/AbstractMQProducer.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.test.clientinterface;
+import java.nio.charset.StandardCharsets;
import java.util.Date;
import java.util.List;
import org.apache.rocketmq.common.message.MessageQueue;
@@ -98,7 +99,7 @@ public abstract class AbstractMQProducer extends MQCollector implements MQProduc
Object objMsg = null;
if (this instanceof RMQNormalProducer) {
org.apache.rocketmq.common.message.Message msg = new org.apache.rocketmq.common.message.Message(
- topic, (RandomUtil.getStringByUUID() + "." + new Date()).getBytes());
+ topic, (RandomUtil.getStringByUUID() + "." + new Date()).getBytes(StandardCharsets.UTF_8));
objMsg = msg;
if (tag != null) {
msg.setTags(tag);
diff --git a/test/src/main/java/org/apache/rocketmq/test/factory/MQMessageFactory.java b/test/src/main/java/org/apache/rocketmq/test/factory/MQMessageFactory.java
index f998fcb13e..ca472f0393 100644
--- a/test/src/main/java/org/apache/rocketmq/test/factory/MQMessageFactory.java
+++ b/test/src/main/java/org/apache/rocketmq/test/factory/MQMessageFactory.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.test.factory;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -33,7 +34,7 @@ public class MQMessageFactory {
public static List