mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #2652] change the method name to incrementAndGet
Co-authored-by: 张玻 <zhangbo@mydataway.com>
This commit is contained in:
@@ -23,7 +23,7 @@ public class ThreadLocalIndex {
|
||||
private final ThreadLocal<Integer> threadLocalIndex = new ThreadLocal<Integer>();
|
||||
private final Random random = new Random();
|
||||
|
||||
public int getAndIncrement() {
|
||||
public int incrementAndGet() {
|
||||
Integer index = this.threadLocalIndex.get();
|
||||
if (null == index) {
|
||||
index = Math.abs(random.nextInt());
|
||||
|
||||
@@ -71,7 +71,7 @@ public class TopicPublishInfo {
|
||||
return selectOneMessageQueue();
|
||||
} else {
|
||||
for (int i = 0; i < this.messageQueueList.size(); i++) {
|
||||
int index = this.sendWhichQueue.getAndIncrement();
|
||||
int index = this.sendWhichQueue.incrementAndGet();
|
||||
int pos = Math.abs(index) % this.messageQueueList.size();
|
||||
if (pos < 0)
|
||||
pos = 0;
|
||||
@@ -85,7 +85,7 @@ public class TopicPublishInfo {
|
||||
}
|
||||
|
||||
public MessageQueue selectOneMessageQueue() {
|
||||
int index = this.sendWhichQueue.getAndIncrement();
|
||||
int index = this.sendWhichQueue.incrementAndGet();
|
||||
int pos = Math.abs(index) % this.messageQueueList.size();
|
||||
if (pos < 0)
|
||||
pos = 0;
|
||||
|
||||
+1
-1
@@ -80,7 +80,7 @@ public class LatencyFaultToleranceImpl implements LatencyFaultTolerance<String>
|
||||
if (half <= 0) {
|
||||
return tmpList.get(0).getName();
|
||||
} else {
|
||||
final int i = this.whichItemWorst.getAndIncrement() % half;
|
||||
final int i = this.whichItemWorst.incrementAndGet() % half;
|
||||
return tmpList.get(i).getName();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -58,7 +58,7 @@ public class MQFaultStrategy {
|
||||
public MessageQueue selectOneMessageQueue(final TopicPublishInfo tpInfo, final String lastBrokerName) {
|
||||
if (this.sendLatencyFaultEnable) {
|
||||
try {
|
||||
int index = tpInfo.getSendWhichQueue().getAndIncrement();
|
||||
int index = tpInfo.getSendWhichQueue().incrementAndGet();
|
||||
for (int i = 0; i < tpInfo.getMessageQueueList().size(); i++) {
|
||||
int pos = Math.abs(index++) % tpInfo.getMessageQueueList().size();
|
||||
if (pos < 0)
|
||||
@@ -74,7 +74,7 @@ public class MQFaultStrategy {
|
||||
final MessageQueue mq = tpInfo.selectOneMessageQueue();
|
||||
if (notBestBroker != null) {
|
||||
mq.setBrokerName(notBestBroker);
|
||||
mq.setQueueId(tpInfo.getSendWhichQueue().getAndIncrement() % writeQueueNums);
|
||||
mq.setQueueId(tpInfo.getSendWhichQueue().incrementAndGet() % writeQueueNums);
|
||||
}
|
||||
return mq;
|
||||
} else {
|
||||
|
||||
@@ -387,7 +387,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
|
||||
filterMqs.add(queue);
|
||||
}
|
||||
}
|
||||
int index = sendWhichQueue.getAndIncrement();
|
||||
int index = sendWhichQueue.incrementAndGet();
|
||||
int pos = Math.abs(index) % filterMqs.size();
|
||||
if (pos < 0) {
|
||||
pos = 0;
|
||||
|
||||
@@ -22,17 +22,17 @@ import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
public class ThreadLocalIndexTest {
|
||||
@Test
|
||||
public void testGetAndIncrement() throws Exception {
|
||||
public void testIncrementAndGet() throws Exception {
|
||||
ThreadLocalIndex localIndex = new ThreadLocalIndex();
|
||||
int initialVal = localIndex.getAndIncrement();
|
||||
int initialVal = localIndex.incrementAndGet();
|
||||
|
||||
assertThat(localIndex.getAndIncrement()).isEqualTo(initialVal + 1);
|
||||
assertThat(localIndex.incrementAndGet()).isEqualTo(initialVal + 1);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGetAndIncrement2() throws Exception {
|
||||
public void testIncrementAndGet2() throws Exception {
|
||||
ThreadLocalIndex localIndex = new ThreadLocalIndex();
|
||||
int initialVal = localIndex.getAndIncrement();
|
||||
int initialVal = localIndex.incrementAndGet();
|
||||
assertThat(initialVal >= 0);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user