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 258499838..f012b66a9 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
@@ -1,13 +1,5 @@
package com.ai.cloud.skywalking.buffer;
-import static com.ai.cloud.skywalking.conf.Config.Buffer.BUFFER_MAX_SIZE;
-import static com.ai.cloud.skywalking.conf.Config.Consumer.CONSUMER_FAIL_RETRY_WAIT_INTERVAL;
-import static com.ai.cloud.skywalking.conf.Config.Consumer.MAX_CONSUMER;
-import static com.ai.cloud.skywalking.conf.Config.Consumer.MAX_WAIT_TIME;
-
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-
import com.ai.cloud.skywalking.conf.Config;
import com.ai.cloud.skywalking.conf.Constants;
import com.ai.cloud.skywalking.protocol.Span;
@@ -15,107 +7,117 @@ import com.ai.cloud.skywalking.selfexamination.HeathReading;
import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector;
import com.ai.cloud.skywalking.sender.DataSenderFactoryWithBalance;
import com.ai.cloud.skywalking.util.AtomicRangeInteger;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import static com.ai.cloud.skywalking.conf.Config.Buffer.BUFFER_MAX_SIZE;
+import static com.ai.cloud.skywalking.conf.Config.Consumer.*;
public class BufferGroup {
- private static Logger logger = LogManager.getLogger(BufferGroup.class);
- private String groupName;
- private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
- AtomicRangeInteger index = new AtomicRangeInteger(0, BUFFER_MAX_SIZE);
+ private static Logger logger = LogManager.getLogger(BufferGroup.class);
+ private String groupName;
+ private Span[] dataBuffer = new Span[BUFFER_MAX_SIZE];
+ AtomicRangeInteger index = new AtomicRangeInteger(0, BUFFER_MAX_SIZE);
- public BufferGroup(String groupName) {
- this.groupName = groupName;
+ public BufferGroup(String groupName) {
+ this.groupName = groupName;
+ startConsumerWorker();
+ }
- int step = (int) Math.ceil(BUFFER_MAX_SIZE * 1.0 / MAX_CONSUMER);
- int start = 0, end = 0;
- while (true) {
- if (end + step >= BUFFER_MAX_SIZE) {
- new ConsumerWorker(start, BUFFER_MAX_SIZE).start();
- break;
- }
- end += step;
- new ConsumerWorker(start, end).start();
- start = end;
- }
- }
+ private void startConsumerWorker() {
+ if (MAX_CONSUMER > 0) {
+ int step = (int) Math.ceil(BUFFER_MAX_SIZE * 1.0 / MAX_CONSUMER);
+ int start = 0, end = 0;
+ while (true) {
+ if (end + step >= BUFFER_MAX_SIZE) {
+ new ConsumerWorker(start, BUFFER_MAX_SIZE).start();
+ break;
+ }
+ end += step;
+ new ConsumerWorker(start, end).start();
+ start = end;
+ }
+ }
+ }
- public void save(Span span) {
- int i = index.getAndIncrement();
- if (dataBuffer[i] != null) {
- logger.warn(
- "Group[{}] index[{}] data collision, discard old data.",
- groupName, i);
- SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.WARNING, "BufferGroup index[" + i + "] data collision, data been coverd.");
- }
- dataBuffer[i] = span;
- SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.INFO, "save span");
- }
+ public void save(Span span) {
+ int i = index.getAndIncrement();
+ if (dataBuffer[i] != null) {
+ logger.warn(
+ "Group[{}] index[{}] data collision, discard old data.",
+ groupName, i);
+ SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.WARNING, "BufferGroup index[" + i + "] data collision, data been coverd.");
+ }
+ dataBuffer[i] = span;
+ SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.INFO, "save span");
+ }
- class ConsumerWorker extends Thread {
- private int start = 0;
- private int end = BUFFER_MAX_SIZE;
+ class ConsumerWorker extends Thread {
+ private int start = 0;
+ private int end = BUFFER_MAX_SIZE;
- private ConsumerWorker(int start, int end) {
- super("ConsumerWorker");
- this.start = start;
- this.end = end;
- }
+ private ConsumerWorker(int start, int end) {
+ super("ConsumerWorker");
+ this.start = start;
+ this.end = end;
+ }
- @Override
- public void run() {
- StringBuilder data = new StringBuilder();
- while (true) {
- boolean bool = false;
- try {
- for (int i = start; i < end; i++) {
- if (dataBuffer[i] == null) {
- continue;
- }
- bool = true;
- if (data.length() + dataBuffer[i].toString().length() >= Config.Sender.MAX_SEND_LENGTH) {
- while (!DataSenderFactoryWithBalance.getSender()
- .send(data.toString())) {
- try {
- Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
- } catch (InterruptedException e) {
- logger.error("Sleep Failure");
- }
- }
- logger.debug("send buried-point data, size:{}", data.length());
- data = new StringBuilder();
- }
+ @Override
+ public void run() {
+ StringBuilder data = new StringBuilder();
+ while (true) {
+ boolean bool = false;
+ try {
+ for (int i = start; i < end; i++) {
+ if (dataBuffer[i] == null) {
+ continue;
+ }
+ bool = true;
+ if (data.length() + dataBuffer[i].toString().length() >= Config.Sender.MAX_SEND_LENGTH) {
+ while (!DataSenderFactoryWithBalance.getSender()
+ .send(data.toString())) {
+ try {
+ Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
+ } catch (InterruptedException e) {
+ logger.error("Sleep Failure");
+ }
+ }
+ logger.debug("send buried-point data, size:{}", data.length());
+ data = new StringBuilder();
+ }
- data.append(dataBuffer[i] + Constants.DATA_SPILT);
- dataBuffer[i] = null;
- }
+ data.append(dataBuffer[i] + Constants.DATA_SPILT);
+ dataBuffer[i] = null;
+ }
- if (data != null && data.length() > 0) {
- while (!DataSenderFactoryWithBalance.getSender().send(
- data.toString())) {
- try {
- Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
- } catch (InterruptedException e) {
- logger.error("Sleep Failure");
- }
- }
- data = new StringBuilder();
- }
- } catch (Throwable e) {
- logger.error("buffer group running failed", e);
- }
+ if (data != null && data.length() > 0) {
+ while (!DataSenderFactoryWithBalance.getSender().send(
+ data.toString())) {
+ try {
+ Thread.sleep(CONSUMER_FAIL_RETRY_WAIT_INTERVAL);
+ } catch (InterruptedException e) {
+ logger.error("Sleep Failure");
+ }
+ }
+ data = new StringBuilder();
+ }
+ } catch (Throwable e) {
+ logger.error("buffer group running failed", e);
+ }
- if (!bool) {
- try {
- Thread.sleep(MAX_WAIT_TIME);
- } catch (InterruptedException e) {
- logger.error("Sleep Failure");
- }
- }
- }
- }
- }
+ if (!bool) {
+ try {
+ Thread.sleep(MAX_WAIT_TIME);
+ } catch (InterruptedException e) {
+ logger.error("Sleep Failure");
+ }
+ }
+ }
+ }
+ }
- public String getGroupName() {
- return groupName;
- }
+ public String getGroupName() {
+ return groupName;
+ }
}
diff --git a/skywalking-api/src/test/java/test/ai/cloud/bufferGroup.java b/skywalking-api/src/test/java/test/ai/cloud/bufferGroup.java
new file mode 100644
index 000000000..968786fc9
--- /dev/null
+++ b/skywalking-api/src/test/java/test/ai/cloud/bufferGroup.java
@@ -0,0 +1,38 @@
+package test.ai.cloud;
+
+import com.ai.cloud.skywalking.buffer.BufferGroup;
+import com.ai.cloud.skywalking.conf.Config;
+import org.junit.Test;
+
+import static org.junit.Assert.assertEquals;
+
+public class bufferGroup {
+
+ @Test
+ public void checkConsumerWorkerIsStartIfConsumerSizeIsZero() {
+ Config.Consumer.MAX_CONSUMER = 0;
+ BufferGroup bufferGroup = new BufferGroup("testBufferGroup");
+ int count = getConsumerWorkerThreadCount();
+ assertEquals(Config.Consumer.MAX_CONSUMER, count);
+ }
+
+ private int getConsumerWorkerThreadCount() {
+ ThreadGroup group = Thread.currentThread().getThreadGroup();
+ ThreadGroup topGroup = group;
+ while (group != null) {
+ topGroup = group;
+ group = group.getParent();
+ }
+
+ int activeCount = topGroup.activeCount();
+ Thread[] threads = new Thread[activeCount];
+ topGroup.enumerate(threads);
+ int count = 0;
+ for (Thread thread : threads) {
+ if (thread != null && "ConsumerWorker".equals(thread.getName())) {
+ count++;
+ }
+ }
+ return count;
+ }
+}
diff --git a/test/skywalking-test-api/pom.xml b/test/skywalking-test-api/pom.xml
new file mode 100644
index 000000000..bbec2da81
--- /dev/null
+++ b/test/skywalking-test-api/pom.xml
@@ -0,0 +1,34 @@
+
+ 4.0.0
+
+
+ com.ai.cloud
+ skywalking
+ 1.0-Final
+
+
+ skywalking-test-api
+ jar
+
+ skywalking-test-api
+ http://maven.apache.org
+
+
+ UTF-8
+
+
+
+
+ junit
+ junit
+ 4.12
+ test
+
+
+ com.ai.cloud
+ skywalking-api
+ ${parent.version}
+
+
+
diff --git a/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeAssert.java b/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeAssert.java
new file mode 100644
index 000000000..d3609813a
--- /dev/null
+++ b/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeAssert.java
@@ -0,0 +1,22 @@
+package test.com.ai.skywalking.api;
+
+import com.ai.cloud.skywalking.protocol.Span;
+
+import java.util.List;
+
+public class TraceTreeAssert {
+
+ public static void assertEquals(String[][] traceTree) {
+ List traceSpanList = TraceTreeDataAcquirer.acquireCurrentTraceSpanData();
+
+ validateTraceId(traceSpanList);
+ validateSpanData(traceSpanList, traceTree);
+ }
+
+ private static void validateSpanData(List traceSpanList, String[][] traceTree) {
+
+ }
+
+ private static void validateTraceId(List traceSpanList) {
+ }
+}
diff --git a/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeDataAcquirer.java b/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeDataAcquirer.java
new file mode 100644
index 000000000..002e33dee
--- /dev/null
+++ b/test/skywalking-test-api/src/main/java/test/com/ai/skywalking/api/TraceTreeDataAcquirer.java
@@ -0,0 +1,15 @@
+package test.com.ai.skywalking.api;
+
+import com.ai.cloud.skywalking.protocol.Span;
+
+import java.util.List;
+
+/**
+ * Created by xin on 16-6-6.
+ */
+public class TraceTreeDataAcquirer {
+ public static List acquireCurrentTraceSpanData() {
+
+ return null;
+ }
+}