diff --git a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/buffer/Channels.java b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/buffer/Channels.java index e4b2e6182..e577ac16f 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/buffer/Channels.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/buffer/Channels.java @@ -12,9 +12,11 @@ import org.skywalking.apm.commons.datacarrier.partition.IDataPartitioner; public class Channels { private final Buffer[] bufferChannels; private IDataPartitioner dataPartitioner; + private BufferStrategy strategy; public Channels(int channelSize, int bufferSize, IDataPartitioner partitioner, BufferStrategy strategy) { this.dataPartitioner = partitioner; + this.strategy = strategy; bufferChannels = new Buffer[channelSize]; for (int i = 0; i < channelSize; i++) { bufferChannels[i] = new Buffer(bufferSize, strategy); @@ -23,7 +25,19 @@ public class Channels { public boolean save(T data) { int index = dataPartitioner.partition(bufferChannels.length, data); - return bufferChannels[index].save(data); + int retryCountDown = 1; + if (BufferStrategy.IF_POSSIBLE.equals(strategy)) { + int maxRetryCount = dataPartitioner.maxRetryCount(); + if (maxRetryCount > 1) { + retryCountDown = maxRetryCount; + } + } + for (; retryCountDown > 0; retryCountDown--) { + if (bufferChannels[index].save(data)) { + return true; + } + } + return false; } public void setPartitioner(IDataPartitioner dataPartitioner) { diff --git a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/IDataPartitioner.java b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/IDataPartitioner.java index 105bcbc4e..4a9089c99 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/IDataPartitioner.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/IDataPartitioner.java @@ -1,8 +1,17 @@ package org.skywalking.apm.commons.datacarrier.partition; +import org.skywalking.apm.commons.datacarrier.buffer.BufferStrategy; + /** * Created by wusheng on 2016/10/25. */ public interface IDataPartitioner { int partition(int total, T data); + + /** + * @return an integer represents how many times should retry when {@link BufferStrategy#IF_POSSIBLE}. + * + * Less or equal 1, means not support retry. + */ + int maxRetryCount(); } diff --git a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/ProducerThreadPartitioner.java b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/ProducerThreadPartitioner.java index 3ebbee592..f3cbd1f1a 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/ProducerThreadPartitioner.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/ProducerThreadPartitioner.java @@ -6,8 +6,22 @@ package org.skywalking.apm.commons.datacarrier.partition; * Created by wusheng on 2016/10/25. */ public class ProducerThreadPartitioner implements IDataPartitioner { + private int retryTime = 3; + + public ProducerThreadPartitioner() { + } + + public ProducerThreadPartitioner(int retryTime) { + this.retryTime = retryTime; + } + @Override public int partition(int total, T data) { return (int)Thread.currentThread().getId() % total; } + + @Override + public int maxRetryCount() { + return 1; + } } diff --git a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/SimpleRollingPartitioner.java b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/SimpleRollingPartitioner.java index e44c89963..c9f425e81 100644 --- a/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/SimpleRollingPartitioner.java +++ b/apm-commons/apm-datacarrier/src/main/java/org/skywalking/apm/commons/datacarrier/partition/SimpleRollingPartitioner.java @@ -14,8 +14,8 @@ public class SimpleRollingPartitioner implements IDataPartitioner { return Math.abs(i++ % total); } - public static void main(String[] args) { - SimpleRollingPartitioner s = new SimpleRollingPartitioner(); - System.out.print(s.i++ % 10); + @Override + public int maxRetryCount() { + return 3; } }