1.消费线程为0的时候,不启动消费线程

2. 提交测试工程
This commit is contained in:
ascrutae 2016-06-06 11:13:15 +08:00
parent c6aa2ac941
commit f706bff794
5 changed files with 209 additions and 98 deletions

View File

@ -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;
}
}

View File

@ -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;
}
}

View File

@ -0,0 +1,34 @@
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>com.ai.cloud</groupId>
<artifactId>skywalking</artifactId>
<version>1.0-Final</version>
</parent>
<artifactId>skywalking-test-api</artifactId>
<packaging>jar</packaging>
<name>skywalking-test-api</name>
<url>http://maven.apache.org</url>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.12</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.ai.cloud</groupId>
<artifactId>skywalking-api</artifactId>
<version>${parent.version}</version>
</dependency>
</dependencies>
</project>

View File

@ -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<Span> traceSpanList = TraceTreeDataAcquirer.acquireCurrentTraceSpanData();
validateTraceId(traceSpanList);
validateSpanData(traceSpanList, traceTree);
}
private static void validateSpanData(List<Span> traceSpanList, String[][] traceTree) {
}
private static void validateTraceId(List<Span> traceSpanList) {
}
}

View File

@ -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<Span> acquireCurrentTraceSpanData() {
return null;
}
}