From 6278973886b65e2fa043ee17f997996f4539b5b7 Mon Sep 17 00:00:00 2001 From: wusheng Date: Tue, 23 Feb 2016 10:50:35 +0800 Subject: [PATCH] =?UTF-8?q?1.=E6=8F=90=E4=BA=A4=E6=96=B0=E7=9A=84AtomicRan?= =?UTF-8?q?geInteger=EF=BC=8C=E9=81=BF=E5=85=8D=E6=95=B4=E6=95=B0=E8=B6=8A?= =?UTF-8?q?=E7=95=8C=E3=80=81=E6=95=B4=E6=95=B0=E5=BE=AA=E7=8E=AF=EF=BC=8C?= =?UTF-8?q?=E9=80=A0=E6=88=90=E9=97=AE=E9=A2=98=E3=80=82=E4=BF=AE=E5=A4=8D?= =?UTF-8?q?#31=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../cloud/skywalking/buffer/BufferGroup.java | 7 +- .../skywalking/util/AtomicRangeInteger.java | 67 +++++++++++++++++++ .../util/AtomicRangeIntegerTest.java | 57 ++++++++++++++++ .../reciever/buffer/DataBufferThread.java | 24 ++++--- 4 files changed, 142 insertions(+), 13 deletions(-) create mode 100644 skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/AtomicRangeInteger.java create mode 100644 skywalking-protocol/src/test/java/test/ai/cloud/skywalking/util/AtomicRangeIntegerTest.java diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java index cf6cc4182..b74b3f352 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java @@ -5,8 +5,6 @@ import static com.ai.cloud.skywalking.conf.Config.Consumer.CONSUMER_FAIL_RETRY_W import static com.ai.cloud.skywalking.conf.Config.Consumer.MAX_CONSUMER; import static com.ai.cloud.skywalking.conf.Config.Consumer.MAX_WAIT_TIME; -import java.util.concurrent.atomic.AtomicInteger; - import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -14,12 +12,13 @@ import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.conf.Constants; import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.sender.DataSenderFactoryWithBalance; +import com.ai.cloud.skywalking.util.AtomicRangeInteger; public class BufferGroup { private static Logger logger = LogManager.getLogger(BufferGroup.class); private String groupName; private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE]; - AtomicInteger index = new AtomicInteger(0); + AtomicRangeInteger index = new AtomicRangeInteger(0, BUFFER_MAX_SIZE); public BufferGroup(String groupName) { this.groupName = groupName; @@ -38,7 +37,7 @@ public class BufferGroup { } public void save(Span span) { - int i = Math.abs(index.getAndIncrement() % BUFFER_MAX_SIZE); + int i = index.getAndIncrement(); if (dataBuffer[i] != null) { logger.warn( "Group[{}] index[{}] data collision, discard old data.", diff --git a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/AtomicRangeInteger.java b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/AtomicRangeInteger.java new file mode 100644 index 000000000..bf1fe360d --- /dev/null +++ b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/AtomicRangeInteger.java @@ -0,0 +1,67 @@ +package com.ai.cloud.skywalking.util; + +import java.util.concurrent.atomic.AtomicInteger; + +/** + * 原子性,带范围性自增的整数 + * + * @author wusheng + * + */ +public class AtomicRangeInteger extends Number implements java.io.Serializable { + private static final long serialVersionUID = -4099792402691141643L; + + private AtomicInteger value; + + private int startValue; + private int endValue; + + /** + * Creates a new AtomicInteger with the given initial value and max value + * + * @param initialValue + * the initial value + */ + public AtomicRangeInteger(int startValue, int endValue) { + value = new AtomicInteger(startValue); + this.startValue = startValue; + this.endValue = endValue; + } + + /** + * Atomically increments by one the current value. + * + * @return the previous value + */ + public final int getAndIncrement() { + for (;;) { + int current = get(); + int next = current + 1; + if (next >= this.endValue) { + next = this.startValue; + } + if (value.compareAndSet(current, next)) + return current; + } + } + + public final int get() { + return value.get(); + } + + public int intValue() { + return value.intValue(); + } + + public long longValue() { + return value.longValue(); + } + + public float floatValue() { + return value.floatValue(); + } + + public double doubleValue() { + return value.doubleValue(); + } +} diff --git a/skywalking-protocol/src/test/java/test/ai/cloud/skywalking/util/AtomicRangeIntegerTest.java b/skywalking-protocol/src/test/java/test/ai/cloud/skywalking/util/AtomicRangeIntegerTest.java new file mode 100644 index 000000000..c3e9330d2 --- /dev/null +++ b/skywalking-protocol/src/test/java/test/ai/cloud/skywalking/util/AtomicRangeIntegerTest.java @@ -0,0 +1,57 @@ +package test.ai.cloud.skywalking.util; + +import java.util.ArrayList; +import java.util.List; + +import junit.framework.Assert; +import junit.framework.TestCase; + +import com.ai.cloud.skywalking.util.AtomicRangeInteger; + +public class AtomicRangeIntegerTest extends TestCase{ + public void testGet(){ + AtomicRangeInteger ari = new AtomicRangeInteger(0, 12); + for(int i = 0; i < 51; i++){ + System.out.print(ari.getAndIncrement() + ";"); + } + } + + public void testMultiThreads(){ + List tlist = new ArrayList(); + for(int i = 0; i < 20; i++){ + RangeIntegerThread t = new RangeIntegerThread(); + tlist.add(t); + t.start(); + } + + for(int i = 0; i < tlist.size(); i++){ + try { + tlist.get(i).join(); + } catch (InterruptedException e) { + e.printStackTrace(); + } + } + } + +} + +class RangeIntegerThread extends Thread{ + private static String[] buffer = new String[500000000]; + + private static AtomicRangeInteger ari = new AtomicRangeInteger(0, buffer.length); + + @Override + public void run(){ + while(true){ + int i = ari.getAndIncrement(); + if(i % 10000000 == 0){ + System.out.println(ari.get()); + } + if(i >= buffer.length - 100000){ + break; + } + Assert.assertNull(buffer[i]); + buffer[i] = "string"; + } + } +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java index 5071ac07c..3c36f98fd 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java @@ -1,19 +1,25 @@ package com.ai.cloud.skywalking.reciever.buffer; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - -import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; -import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.DATA_BUFFER_FILE_PARENT_DIRECTORY; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.DATA_CONFLICT_WAIT_TIME; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.DATA_FILE_MAX_LENGTH; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.FLUSH_NUMBER_OF_CACHE; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.MAX_WAIT_TIME; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.PER_THREAD_MAX_BUFFER_NUMBER; +import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.WRITE_DATA_FAILURE_RETRY_INTERVAL; import java.io.File; import java.io.FileNotFoundException; import java.io.FileOutputStream; import java.io.IOException; import java.util.UUID; -import java.util.concurrent.atomic.AtomicInteger; -import static com.ai.cloud.skywalking.reciever.conf.Config.Buffer.*; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; +import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; +import com.ai.cloud.skywalking.util.AtomicRangeInteger; public class DataBufferThread extends Thread { @@ -21,7 +27,7 @@ public class DataBufferThread extends Thread { private byte[][] data = new byte[PER_THREAD_MAX_BUFFER_NUMBER][]; private File file; private FileOutputStream outputStream; - private AtomicInteger index = new AtomicInteger(); + private AtomicRangeInteger index = new AtomicRangeInteger(0, PER_THREAD_MAX_BUFFER_NUMBER); public DataBufferThread(int threadIdx) { super("DataBufferThread_" + threadIdx); @@ -142,7 +148,7 @@ public class DataBufferThread extends Thread { } public void saveTemporarily(byte[] s) { - int i = Math.abs(index.getAndIncrement() % data.length); + int i = index.getAndIncrement(); while (data[i] != null) { try { ServerHealthCollector.getCurrentHeathReading(null).updateData(ServerHeathReading.WARNING, "DataBuffer index[" + i + "] data collision, service pausing. ");