KLALB V3.6.0 写了一半
This commit is contained in:
@@ -0,0 +1,280 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
import java.lang.ref.Cleaner;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
import java.util.concurrent.locks.LockSupport;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.kne.cloud.network.ThreadTool;
|
||||
import org.kne.concurrent.SpinLock;
|
||||
import org.kne.opencl64.Releaser;
|
||||
|
||||
/**
|
||||
* 网络消息批量发送器,用于减少小包数量,降低网络性能开销
|
||||
* @param <T> 消息类型
|
||||
*/
|
||||
|
||||
public class ArrayListMessageBatcher<T> implements MessageBatcher<T>{
|
||||
private static final Cleaner clr=Cleaner.create();
|
||||
|
||||
private ArrayListMessageBatcher0<T> impl;
|
||||
|
||||
private ArrayListMessageBatcherReleaser<T> releaser;
|
||||
|
||||
public ArrayListMessageBatcher(int batchSize, long maxDelay) {
|
||||
impl=new ArrayListMessageBatcher0<T>(batchSize, maxDelay);
|
||||
releaser=new ArrayListMessageBatcherReleaser<T>(impl);
|
||||
clr.register(this, releaser);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getBatchSize() {
|
||||
return impl.getBatchSize();
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getMaxDelay() {
|
||||
return impl.getMaxDelay();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setConsumer(Consumer<List<T>> consumer) {
|
||||
impl.setConsumer(consumer);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putMessage(T message) {
|
||||
impl.putMessage(message);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putMessages(List<T> messages) {
|
||||
impl.putMessages(messages);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void flush() {
|
||||
impl.flush();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getQueueSize() {
|
||||
return impl.getQueueSize();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isClosed() {
|
||||
return impl.isClosed();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
impl.close();
|
||||
}
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
for(;;){
|
||||
ArrayListMessageBatcher<Integer>mes=new ArrayListMessageBatcher<>(10, 1000000L);
|
||||
System.gc();
|
||||
}
|
||||
/*ArrayListMessageBatcher<Integer>mes=new ArrayListMessageBatcher<>(10, 1000000L);
|
||||
mes.setConsumer((x)->{System.out.println(x);});
|
||||
int n=0;
|
||||
for(;;){
|
||||
mes.putMessage(n++);
|
||||
mes.putMessage(n++);
|
||||
mes.putMessage(n++);
|
||||
mes.putMessage(n++);
|
||||
Thread.sleep(1);
|
||||
}*/
|
||||
}
|
||||
}
|
||||
class ArrayListMessageBatcherReleaser<T> extends Releaser<ArrayListMessageBatcher0<T>>{
|
||||
|
||||
public ArrayListMessageBatcherReleaser(ArrayListMessageBatcher0<T> resource) {
|
||||
super(resource);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void release(ArrayListMessageBatcher0<T> resource) {
|
||||
resource.close();
|
||||
//System.out.println("release!");
|
||||
}
|
||||
|
||||
}
|
||||
class ArrayListMessageBatcher0<T> implements Runnable,MessageBatcher<T>{
|
||||
private ArrayList<T> messageList;
|
||||
private Lock lock=new SpinLock();
|
||||
private final int batchSize;
|
||||
private final long maxDelay;
|
||||
private Consumer<List<T>> consumer;
|
||||
private final Thread batchThread;
|
||||
private final AtomicBoolean closed;
|
||||
private long firstTime;
|
||||
|
||||
/**
|
||||
* 构造函数
|
||||
* @param batchSize 批量大小,当消息达到此数量时触发发送
|
||||
* @param maxDelayMillis 最大延迟时间(毫秒),即使消息数量不足也会触发发送
|
||||
*/
|
||||
public ArrayListMessageBatcher0(int batchSize, long maxDelay) {
|
||||
this.batchSize = batchSize;
|
||||
this.maxDelay = maxDelay;
|
||||
this.closed = new AtomicBoolean(false);
|
||||
|
||||
// 创建并启动批量处理线程
|
||||
this.batchThread = ThreadTool.makeVDaemonThread("ArrayListMessageBatcher Thread", this );
|
||||
this.batchThread.setDaemon(true);
|
||||
this.batchThread.start();
|
||||
}
|
||||
|
||||
/**
|
||||
* 设置消息消费者回调函数
|
||||
* @param consumer 消费者回调
|
||||
*/
|
||||
public void setConsumer(Consumer<List<T>> consumer) {
|
||||
this.consumer = consumer;
|
||||
}
|
||||
|
||||
/**
|
||||
* 添加消息到批量处理器
|
||||
* @param message 消息对象
|
||||
*/
|
||||
public void putMessage(T message) {
|
||||
if (message != null) {
|
||||
lock.lock();
|
||||
try {
|
||||
if(messageList==null) {
|
||||
messageList=new ArrayList<T>(batchSize);
|
||||
firstTime=System.nanoTime();
|
||||
}
|
||||
messageList.add(message);
|
||||
check(false);
|
||||
}finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* 批量添加消息
|
||||
* @param messages 消息列表
|
||||
*/
|
||||
public void putMessages(List<T> messages) {
|
||||
for(T message :messages) {
|
||||
putMessage(message);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private void check(boolean force) {
|
||||
if(messageList==null)
|
||||
return;
|
||||
long currTime=System.nanoTime();
|
||||
if(messageList.size()>=batchSize||(currTime-firstTime)>maxDelay||force) {
|
||||
if(consumer!=null) {
|
||||
try {
|
||||
consumer.accept(messageList);
|
||||
}catch(Exception e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
messageList=null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止批量处理器
|
||||
*/
|
||||
public void close() {
|
||||
closed.set(true);
|
||||
LockSupport.unpark(batchThread);
|
||||
// 确保所有消息都被发送
|
||||
flush();
|
||||
}
|
||||
|
||||
/**
|
||||
* 批量处理消息的核心方法
|
||||
*/
|
||||
public void run() {
|
||||
while(!closed.get()) {
|
||||
lock.lock();
|
||||
try {
|
||||
check(false);
|
||||
}finally {
|
||||
lock.unlock();
|
||||
}
|
||||
// 更精确的等待策略
|
||||
long sleepTimeNanos = calculateSleepTime();
|
||||
if (sleepTimeNanos > 0) {
|
||||
LockSupport.parkNanos(sleepTimeNanos);
|
||||
} else {
|
||||
// 避免忙等待
|
||||
LockSupport.parkNanos(1_000_000L); // 1ms
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 计算需要等待的时间
|
||||
* @return 等待时间(纳秒)
|
||||
*/
|
||||
private long calculateSleepTime() {
|
||||
lock.lock();
|
||||
try {
|
||||
if (messageList == null || messageList.isEmpty()) {
|
||||
return 10_000_000L; // 10ms
|
||||
}
|
||||
|
||||
long elapsed = System.nanoTime() - firstTime;
|
||||
long remaining = maxDelay - elapsed;
|
||||
|
||||
if (remaining <= 0) {
|
||||
return 0; // 立即处理
|
||||
}
|
||||
|
||||
// 返回剩余时间或10ms中的较小值
|
||||
return Math.min(remaining, 10_000_000L);
|
||||
} finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
public int getBatchSize() {
|
||||
return batchSize;
|
||||
}
|
||||
|
||||
public long getMaxDelay() {
|
||||
return maxDelay;
|
||||
}
|
||||
|
||||
public void flush() {
|
||||
lock.lock();
|
||||
try {
|
||||
check(true);
|
||||
}finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
public int getQueueSize() {
|
||||
List<T> cmessageList =messageList;
|
||||
if(cmessageList==null)
|
||||
return 0;
|
||||
return cmessageList.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isClosed() {
|
||||
return closed.get();
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,150 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class BBRCongestionAlgorithm implements CongestionAlgorithm ,DetnetCongestionAlgorithm{
|
||||
|
||||
private static final long BURST_TIME = 10000000L;
|
||||
private static final long BIAS = 500000L;
|
||||
private static long MIN_SPEED = 512 * 1024L;
|
||||
private static long MIN_WINDOW = 16384;
|
||||
|
||||
private Consumer<Long> windowControlConsumer;
|
||||
private BiConsumer<Long, Long> speedControlConsumer;
|
||||
|
||||
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
|
||||
private long MIN_RTTVAR = 5000000L;
|
||||
private volatile long RTTMin = 1000000000L;
|
||||
private volatile long RTTVar = 1000000000L;
|
||||
private volatile long RTTAvg = 1000000000L;
|
||||
private volatile long RTTAvg2 = 1000000000L;
|
||||
private volatile long RTTTotal = 0;
|
||||
private volatile long RTTCount = 0;
|
||||
private volatile long RTO = 1000000000L;
|
||||
|
||||
|
||||
private volatile long window=MIN_WINDOW;
|
||||
private volatile long speed=MIN_SPEED;
|
||||
private volatile long maxwindow=99999999999L;
|
||||
private long congressSpeed=MIN_SPEED;
|
||||
private double MIN_GAIN=1.05;
|
||||
private double MAX_GAIN=1.20;//2.8853900817779
|
||||
private double windowGain=MIN_GAIN;
|
||||
private double speedGain=MAX_GAIN;//1.20 MAX_GAIN
|
||||
@Override
|
||||
public void reset() {
|
||||
window = MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
speed = MIN_SPEED;
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
windowGain=MAX_GAIN;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
this.windowControlConsumer = windowControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
this.speedControlConsumer = speedControlConsumer;
|
||||
}
|
||||
|
||||
private long RTTstartTime = System.nanoTime();
|
||||
private long bandwidth;
|
||||
@Override
|
||||
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
if (latencyns <= RTTMin) {
|
||||
RTTMin = latencyns;
|
||||
} else {
|
||||
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
|
||||
}
|
||||
if (firstUpdate.compareAndSet(true, false)) {
|
||||
RTTAvg = latencyns;
|
||||
RTTVar = latencyns / 2;
|
||||
RTTAvg2 = latencyns;
|
||||
} else {
|
||||
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
|
||||
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
|
||||
}
|
||||
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
|
||||
|
||||
RTTTotal += latencyns;
|
||||
RTTCount++;
|
||||
/*double rate=latencyns/(double)(RTTMin+BIAS);
|
||||
if(rate>2) {
|
||||
windowGain=MIN_GAIN;
|
||||
}else if(rate<1.2){
|
||||
windowGain=MAX_GAIN;
|
||||
|
||||
}*/
|
||||
long curr=System.nanoTime();
|
||||
if(curr-RTTstartTime>RTTAvg2) {
|
||||
if(RTTCount>0) {
|
||||
RTTAvg2=(RTTAvg2+RTTTotal/RTTCount)/2;
|
||||
RTTTotal=0;
|
||||
RTTCount=0;
|
||||
}
|
||||
|
||||
RTTstartTime=curr;
|
||||
|
||||
if (bandwidth >= congressSpeed) {
|
||||
congressSpeed = bandwidth;
|
||||
//congressSpeed = (congressSpeed + bandwidth) / 2;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
|
||||
}
|
||||
speed=(long) (congressSpeed * speedGain)+MIN_SPEED ;
|
||||
window= ((long) (congressSpeed*windowGain )*(RTTMin+BIAS)/1000000000L)+MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return RTO;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long used) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
this.bandwidth=bandwidth;
|
||||
|
||||
//System.out.println(gain);
|
||||
//System.out.println(bandwidth+" "+congressSpeed+" "+window);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setUpperDelayBound(double upper) {
|
||||
MAX_GAIN=upper;
|
||||
speedGain=MAX_GAIN;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setLowerDelayBound(double lower) {
|
||||
MIN_GAIN=lower;
|
||||
windowGain=MIN_GAIN;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,150 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class BBRVegasCongressAlgorithm implements CongestionAlgorithm {
|
||||
|
||||
private static final long BURST_TIME = 2000000L;
|
||||
private static final long BIAS = 2000000L;
|
||||
private static long MIN_SPEED = 128 * 1024L;
|
||||
private static long MIN_WINDOW = 8192;
|
||||
|
||||
private Consumer<Long> windowControlConsumer;
|
||||
private BiConsumer<Long, Long> speedControlConsumer;
|
||||
|
||||
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
|
||||
private long MIN_RTTVAR = 50000000L;
|
||||
private volatile long RTTMin = 1000000000L;
|
||||
private volatile long RTTVar = 1000000000L;
|
||||
private volatile long RTTAvg = 1000000000L;
|
||||
private volatile long RTTAvg2 = 1000000000L;
|
||||
private volatile long RTTTotal = 0;
|
||||
private volatile long RTTCount = 0;
|
||||
private volatile long RTO = 1000000000L;
|
||||
|
||||
|
||||
private volatile long window=MIN_WINDOW;
|
||||
private volatile long speed=MIN_SPEED;
|
||||
private volatile long maxwindow=99999999999L;
|
||||
private long congressSpeed=MIN_SPEED;
|
||||
private double MAX_GAIN=2.8853900817779;//2.8853900817779
|
||||
private double MIN_GAIN=1.1;
|
||||
private double gain=MAX_GAIN;
|
||||
@Override
|
||||
public void reset() {
|
||||
window = MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
speed = MIN_SPEED;
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
gain=MAX_GAIN;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
this.windowControlConsumer = windowControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
this.speedControlConsumer = speedControlConsumer;
|
||||
}
|
||||
|
||||
private long RTTstartTime = System.nanoTime();
|
||||
@Override
|
||||
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
if (latencyns <= RTTMin) {
|
||||
RTTMin = latencyns;
|
||||
} else {
|
||||
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
|
||||
}
|
||||
if (firstUpdate.compareAndSet(true, false)) {
|
||||
RTTAvg = latencyns;
|
||||
RTTVar = latencyns / 2;
|
||||
RTTAvg2 = latencyns;
|
||||
} else {
|
||||
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
|
||||
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
|
||||
}
|
||||
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
|
||||
|
||||
RTTTotal += latencyns;
|
||||
RTTCount++;
|
||||
long curr=System.nanoTime();
|
||||
if(curr-RTTstartTime>RTTMin) {
|
||||
if(RTTCount>0) {
|
||||
RTTAvg2=(RTTAvg2+RTTTotal/RTTCount)/2;
|
||||
RTTTotal=0;
|
||||
RTTCount=0;
|
||||
}
|
||||
long RTTAvg2x=Math.max(RTTAvg2-BIAS,RTTMin);
|
||||
double excepted=window/(double)RTTMin;
|
||||
double actual=window/(double)RTTAvg2x;
|
||||
double diff=(excepted-actual)*RTTMin;
|
||||
if(diff>65536*3+3) {
|
||||
gain-=0.0005;
|
||||
}
|
||||
if(diff<65536*3+1) {
|
||||
gain+=0.0005;
|
||||
}
|
||||
if(gain>=MAX_GAIN){
|
||||
gain=MAX_GAIN;
|
||||
//System.out.println("window:"+window);
|
||||
}
|
||||
if(gain<MIN_GAIN) {
|
||||
gain=MIN_GAIN;
|
||||
//System.out.println("window:"+window);
|
||||
}
|
||||
RTTstartTime=curr;
|
||||
|
||||
}
|
||||
//System.out.println(speed+" "+window);
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return RTO;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long used) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
/*if (bandwidth >= congressSpeed) {
|
||||
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed *2 + bandwidth) / 3;
|
||||
}*/
|
||||
if (bandwidth >= congressSpeed) {
|
||||
//congressSpeed = bandwidth;
|
||||
congressSpeed = (congressSpeed + bandwidth) / 2;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
|
||||
}
|
||||
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
|
||||
window=Math.max( (Math.max(MIN_SPEED,(long) (congressSpeed*gain ))*Math.max(2000000L,RTTMin)/1000000000L),MIN_WINDOW);
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
//System.out.println(gain);
|
||||
//System.out.println(bandwidth+" "+congressSpeed+" "+window);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public interface CongestionAlgorithm {
|
||||
public void reset();
|
||||
public void setWindowControlConsumer(Consumer<Long>windowControlConsumer);
|
||||
public void setSpeedControlConsumer(BiConsumer<Long,Long>speedControlConsumer);
|
||||
public void putAck(long packetSize,long latencyns,boolean ecn);
|
||||
public void putLoss(long packetSize,long lossns,int losscounter);
|
||||
public void setCurrentWindowUsed(long maxwindow);
|
||||
public void setCurrentBandwidth(long bandwidth);
|
||||
public long getRTO();
|
||||
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class DCTCP2CongestionAlgorithm implements CongestionAlgorithm {
|
||||
|
||||
private static final long BURST_TIME = 1000000L;
|
||||
private static final long BIAS = 2000000L;
|
||||
private static long MIN_SPEED = 64 * 1024L;
|
||||
private static long MIN_WINDOW = 32768;
|
||||
|
||||
private Consumer<Long> windowControlConsumer;
|
||||
private BiConsumer<Long, Long> speedControlConsumer;
|
||||
|
||||
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
|
||||
private long MIN_RTTVAR = 400000000L;
|
||||
private volatile long RTTMin = 1000000000L;
|
||||
private volatile long RTTVar = 1000000000L;
|
||||
private volatile long RTTAvg = 1000000000L;
|
||||
private double ECNAvg2 = 0L;
|
||||
private double alpha=0;
|
||||
private volatile AtomicLong ECNSize = new AtomicLong(0);
|
||||
private volatile AtomicLong TotalSize = new AtomicLong(0);
|
||||
private volatile long RTO = 1000000000L;
|
||||
|
||||
|
||||
private volatile long window=MIN_WINDOW;
|
||||
private volatile long windowUsed=MIN_WINDOW;
|
||||
private volatile long congressWindowSize=MIN_WINDOW;
|
||||
private volatile long speed=MIN_SPEED;
|
||||
private long congressSpeed=MIN_SPEED;
|
||||
|
||||
private long RTTstartTime = System.nanoTime();
|
||||
private long bandwidth=MIN_SPEED;
|
||||
private double MAX_GAIN=1.5;//2.8853900817779
|
||||
private double MIN_GAIN=1.2;
|
||||
private double gain=MAX_GAIN;
|
||||
|
||||
@Override
|
||||
public void reset() {
|
||||
window = MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
speed = MIN_SPEED;
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
gain=MAX_GAIN;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
this.windowControlConsumer = windowControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
this.speedControlConsumer = speedControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
if (latencyns <= RTTMin) {
|
||||
RTTMin = latencyns;
|
||||
} else {
|
||||
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
|
||||
}
|
||||
if (firstUpdate.compareAndSet(true, false)) {
|
||||
RTTAvg = latencyns;
|
||||
RTTVar = latencyns / 2;
|
||||
} else {
|
||||
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
|
||||
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
|
||||
}
|
||||
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
|
||||
|
||||
if(ecn) {
|
||||
ECNSize.addAndGet(packetSize);
|
||||
}
|
||||
TotalSize.addAndGet(packetSize);
|
||||
|
||||
|
||||
long curr=System.nanoTime();
|
||||
if(curr-RTTstartTime>(RTTAvg+BIAS)) {
|
||||
RTTstartTime=curr;
|
||||
|
||||
long totalSize=TotalSize.getAndSet(0);
|
||||
long ecnSize=ECNSize.getAndSet(0);
|
||||
if(totalSize!=0) {
|
||||
ECNAvg2=ecnSize/(double)totalSize;
|
||||
alpha=(alpha*15+ECNAvg2)/16;
|
||||
}
|
||||
|
||||
if(ECNAvg2>0.5) {
|
||||
gain=MIN_GAIN;
|
||||
if(congressWindowSize>MIN_WINDOW) {
|
||||
congressWindowSize-=200;
|
||||
updateWindowSize();
|
||||
}
|
||||
}else if(ECNAvg2>0.1){
|
||||
if(windowUsed*3L>=congressWindowSize) {
|
||||
congressWindowSize+=200;
|
||||
updateWindowSize();
|
||||
}
|
||||
}else {
|
||||
if(windowUsed*3L>=congressWindowSize) {
|
||||
congressWindowSize+=8000;
|
||||
updateWindowSize();
|
||||
}
|
||||
}
|
||||
//System.out.println(ECNAvg2);
|
||||
//System.out.println(speed+" "+window);
|
||||
if (bandwidth >= congressSpeed) {
|
||||
congressSpeed = bandwidth;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed * 99 + bandwidth) / 100;
|
||||
}
|
||||
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
|
||||
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
private void updateWindowSize() {
|
||||
window=Math.max( congressWindowSize,MIN_WINDOW);
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return RTO;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long used) {
|
||||
this.windowUsed=used;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
this.bandwidth=bandwidth;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class DCTCPCongestionAlgorithm implements CongestionAlgorithm {
|
||||
|
||||
private static final long BURST_TIME = 1000000L;
|
||||
private static final long BIAS = 2000000L;
|
||||
private static long MIN_SPEED = 64 * 1024L;
|
||||
private static long MIN_WINDOW = 32768;
|
||||
|
||||
private Consumer<Long> windowControlConsumer;
|
||||
private BiConsumer<Long, Long> speedControlConsumer;
|
||||
|
||||
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
|
||||
private long MIN_RTTVAR = 400000000L;
|
||||
private volatile long RTTMin = 1000000000L;
|
||||
private volatile long RTTVar = 1000000000L;
|
||||
private volatile long RTTAvg = 1000000000L;
|
||||
private double ECNAvg2 = 0L;
|
||||
private double alpha=0;
|
||||
private volatile AtomicLong ECNSize = new AtomicLong(0);
|
||||
private volatile AtomicLong TotalSize = new AtomicLong(0);
|
||||
private volatile long RTO = 1000000000L;
|
||||
|
||||
|
||||
private volatile long window=MIN_WINDOW;
|
||||
private volatile long windowUsed=MIN_WINDOW;
|
||||
private volatile long congressWindowSize=MIN_WINDOW;
|
||||
private volatile long speed=MIN_SPEED;
|
||||
private long congressSpeed=MIN_SPEED;
|
||||
|
||||
private long RTTstartTime = System.nanoTime();
|
||||
private long bandwidth=MIN_SPEED;
|
||||
private double MAX_GAIN=2.8853900817779;//2.8853900817779
|
||||
private double MIN_GAIN=1.2;
|
||||
private double gain=MAX_GAIN;
|
||||
|
||||
@Override
|
||||
public void reset() {
|
||||
window = MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
speed = MIN_SPEED;
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
gain=MAX_GAIN;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
this.windowControlConsumer = windowControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
this.speedControlConsumer = speedControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
if (latencyns <= RTTMin) {
|
||||
RTTMin = latencyns;
|
||||
} else {
|
||||
RTTMin = (RTTMin * 99999 + latencyns) / 100000;
|
||||
}
|
||||
if (firstUpdate.compareAndSet(true, false)) {
|
||||
RTTAvg = latencyns;
|
||||
RTTVar = latencyns / 2;
|
||||
} else {
|
||||
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
|
||||
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
|
||||
}
|
||||
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);// RTTVar*4
|
||||
|
||||
if(ecn) {
|
||||
ECNSize.addAndGet(packetSize);
|
||||
}
|
||||
TotalSize.addAndGet(packetSize);
|
||||
|
||||
|
||||
long curr=System.nanoTime();
|
||||
if(curr-RTTstartTime>(RTTAvg+BIAS)) {
|
||||
RTTstartTime=curr;
|
||||
|
||||
long totalSize=TotalSize.getAndSet(0);
|
||||
long ecnSize=ECNSize.getAndSet(0);
|
||||
if(totalSize!=0) {
|
||||
ECNAvg2=ecnSize/(double)totalSize;
|
||||
alpha=(alpha*15+ECNAvg2)/16;
|
||||
}
|
||||
|
||||
if(ECNAvg2>0.6) {
|
||||
gain=MIN_GAIN;
|
||||
if(congressWindowSize>MIN_WINDOW) {
|
||||
congressWindowSize=Math.max(MIN_WINDOW,(long) (congressWindowSize*(1-alpha/10)));
|
||||
updateWindowSize();
|
||||
}
|
||||
}else if(ECNAvg2>0.2){
|
||||
if(windowUsed*3L>=congressWindowSize) {
|
||||
congressWindowSize+=200;
|
||||
updateWindowSize();
|
||||
}
|
||||
}else {
|
||||
if(windowUsed*3L>=congressWindowSize) {
|
||||
congressWindowSize+=8000;
|
||||
updateWindowSize();
|
||||
}
|
||||
}
|
||||
//System.out.println(ECNAvg2);
|
||||
//System.out.println(speed+" "+window);
|
||||
if (bandwidth >= congressSpeed) {
|
||||
congressSpeed = bandwidth;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed * 99 + bandwidth) / 100;
|
||||
}
|
||||
speed=Math.max(MIN_SPEED,(long) (congressSpeed * gain)) ;
|
||||
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
private void updateWindowSize() {
|
||||
window=Math.max( congressWindowSize,MIN_WINDOW);
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return RTO;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long used) {
|
||||
this.windowUsed=used;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
this.bandwidth=bandwidth;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
public interface DetnetCongestionAlgorithm extends CongestionAlgorithm {
|
||||
public void setUpperDelayBound(double upper);
|
||||
public void setLowerDelayBound(double lower);
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class EmptyCongestionAlgorithm implements CongestionAlgorithm{
|
||||
|
||||
private long timeout;
|
||||
|
||||
@Override
|
||||
public void reset() {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
public long getTimeout() {
|
||||
return timeout;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long maxwindow) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
// TODO 自动生成的方法存根
|
||||
|
||||
}
|
||||
|
||||
public EmptyCongestionAlgorithm(long timeout) {
|
||||
super();
|
||||
this.timeout = timeout;
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return timeout;
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
import java.util.List;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* 消息批处理器接口
|
||||
* 提供批量消息处理功能,支持按数量和时间触发批量发送
|
||||
*/
|
||||
public interface MessageBatcher<T> {
|
||||
|
||||
// ==================== 配置相关方法 ====================
|
||||
|
||||
/**
|
||||
* 获取批处理大小
|
||||
* @return 批处理大小
|
||||
*/
|
||||
int getBatchSize();
|
||||
|
||||
/**
|
||||
* 获取最大延迟时间
|
||||
* @return 最大延迟时间(纳秒)
|
||||
*/
|
||||
long getMaxDelay();
|
||||
|
||||
/**
|
||||
* 设置消息消费者回调函数
|
||||
* @param consumer 消费者回调
|
||||
*/
|
||||
void setConsumer(Consumer<List<T>> consumer);
|
||||
|
||||
// ==================== 消息操作方法 ====================
|
||||
|
||||
/**
|
||||
* 添加消息到批处理器
|
||||
* @param message 消息对象
|
||||
*/
|
||||
void putMessage(T message);
|
||||
|
||||
/**
|
||||
* 批量添加消息
|
||||
* @param messages 消息列表
|
||||
*/
|
||||
void putMessages(List<T> messages);
|
||||
|
||||
/**
|
||||
* 立即刷新并发送所有待处理消息
|
||||
*/
|
||||
void flush();
|
||||
|
||||
// ==================== 状态查询方法 ====================
|
||||
|
||||
/**
|
||||
* 获取当前队列中的消息数量
|
||||
* @return 队列大小
|
||||
*/
|
||||
int getQueueSize();
|
||||
|
||||
|
||||
// ==================== 生命周期管理方法 ====================
|
||||
|
||||
/**
|
||||
* 检查批处理器是否停止
|
||||
* @return 是否停止
|
||||
*/
|
||||
boolean isClosed();
|
||||
|
||||
|
||||
/**
|
||||
* 停止批处理器
|
||||
*/
|
||||
void close();
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.AtomicReferenceArray;
|
||||
|
||||
public class ReceiveByteSlidingWindow<E> {
|
||||
private AtomicReferenceArray<E> array;
|
||||
private AtomicInteger windowSize =new AtomicInteger();
|
||||
private AtomicLong windowPosition=new AtomicLong();
|
||||
|
||||
public int getWindowSize() {
|
||||
return windowSize.get();
|
||||
}
|
||||
|
||||
public void setWindowSize(int windowSize) {
|
||||
this.windowSize.set(windowSize);
|
||||
}
|
||||
|
||||
public long getWindowPosition() {
|
||||
return windowPosition.get();
|
||||
}
|
||||
|
||||
public void setWindowPosition(long windowPosition) {
|
||||
this.windowPosition.set(windowPosition);
|
||||
}
|
||||
|
||||
public ReceiveByteSlidingWindow(int capacity){
|
||||
array=new AtomicReferenceArray<E>(capacity);
|
||||
}
|
||||
|
||||
public boolean push(E object) {
|
||||
Objects.requireNonNull(object);
|
||||
if(array.get((int) ((windowPosition.get()-windowSize.get())%array.length()))==null ) {
|
||||
array.set((int) (windowPosition.getAndIncrement()%array.length()), object);
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
public E pull(long number) {
|
||||
return array.getAndSet((int) (number%array.length()), null);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.Iterator;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.AtomicReferenceArray;
|
||||
|
||||
public class SendByteSlidingWindow<E> implements Iterable<E>{
|
||||
private Object[] array;
|
||||
private volatile int windowSize ;
|
||||
private volatile long windowPosition;
|
||||
|
||||
public int getWindowSize() {
|
||||
return windowSize;
|
||||
}
|
||||
|
||||
public void setWindowSize(int windowSize) {
|
||||
this.windowSize=windowSize;
|
||||
}
|
||||
|
||||
public long getWindowPosition() {
|
||||
return windowPosition;
|
||||
}
|
||||
|
||||
public void setWindowPosition(long windowPosition) {
|
||||
this.windowPosition=windowPosition;
|
||||
}
|
||||
|
||||
public SendByteSlidingWindow(int capacity,int windowSize){
|
||||
if(windowSize>capacity) {
|
||||
throw new IllegalArgumentException("windowSize>capacity!");
|
||||
}
|
||||
array= new Object[capacity];
|
||||
this.windowSize=windowSize;
|
||||
}
|
||||
|
||||
public boolean checkPush() {
|
||||
if(array[calcPosition(windowPosition-windowSize)]==null ) {
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
public void push(E object) {
|
||||
|
||||
array[calcPosition(windowPosition)]= object;
|
||||
}
|
||||
|
||||
public E remove(long number) {
|
||||
if(number<windowPosition&&number>=windowPosition-windowSize) {
|
||||
int indp=calcPosition(number);
|
||||
E old=(E) array[indp];
|
||||
array[indp]=null;
|
||||
return old;
|
||||
}else {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isEmpty() {
|
||||
for (long i = windowPosition-windowSize; i < windowPosition; i++) {
|
||||
if(array[calcPosition(i)]!=null) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
public E get(long number) {
|
||||
if(number<windowPosition&&number>=windowPosition-windowSize) {
|
||||
return (E) array[calcPosition(number)];
|
||||
}else {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public Iterator<E> iterator() {
|
||||
return new Iterator<E>() {
|
||||
private long pointer=windowPosition-windowSize;
|
||||
@Override
|
||||
public boolean hasNext() {
|
||||
return pointer<windowPosition;
|
||||
}
|
||||
|
||||
@Override
|
||||
public E next() {
|
||||
return (E) array[calcPosition(pointer++)];
|
||||
}
|
||||
};
|
||||
}
|
||||
private int calcPosition(long number) {
|
||||
return (int) ((number+array.length)%array.length);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,221 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.io.Closeable;
|
||||
import java.io.IOException;
|
||||
import java.net.SocketException;
|
||||
import java.util.Collection;
|
||||
import java.util.Iterator;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
import java.util.Map.Entry;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.LongAdder;
|
||||
import java.util.concurrent.locks.Condition;
|
||||
import java.util.concurrent.locks.Lock;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import org.kne.cloud.network.NetworkPacket;
|
||||
import org.kne.cloud.network.ThreadTool;
|
||||
import org.kne.cloud.network.ipv6.IPv6Packet;
|
||||
import org.kne.cloud.network.klalb.SendItem;
|
||||
|
||||
public class SendPacketSlidingWindow<K, V extends NetworkPacket> implements Closeable, AutoCloseable {
|
||||
private Map<K, SendItem<V>> sendMap = new ConcurrentHashMap<>(1024, 0.2f);
|
||||
private AtomicLong sendmapWindowUsed=new AtomicLong(0);
|
||||
//private LongAdder sendmapWindowUsed = new LongAdder();
|
||||
|
||||
private int headerCalibrate = 0;
|
||||
|
||||
public void setHeaderCalibrate(int headerCalibrate) {
|
||||
this.headerCalibrate = headerCalibrate;
|
||||
}
|
||||
|
||||
private CongestionAlgorithm algorithm;
|
||||
|
||||
private Consumer<V> resendConsumer;
|
||||
private Consumer<V> closingConsumer;
|
||||
private volatile boolean closed = false;
|
||||
private ResendChecker checker = new ResendChecker();
|
||||
private AtomicLong windowSize = new AtomicLong();
|
||||
|
||||
private Lock lock;
|
||||
|
||||
private Condition condition;
|
||||
private boolean removeAtResend;
|
||||
|
||||
public SendPacketSlidingWindow(CongestionAlgorithm algorithm, long windowSize) {
|
||||
this(algorithm, windowSize, 0, true);
|
||||
}
|
||||
|
||||
public SendPacketSlidingWindow(CongestionAlgorithm algorithm, long windowSize, int headerCalibrate) {
|
||||
this(algorithm, windowSize, headerCalibrate, true);
|
||||
}
|
||||
|
||||
public SendPacketSlidingWindow(CongestionAlgorithm algorithm, long windowSize, int headerCalibrate,
|
||||
boolean removeAtResend) {
|
||||
this.headerCalibrate = headerCalibrate;
|
||||
this.removeAtResend = removeAtResend;
|
||||
this.algorithm = algorithm;
|
||||
this.windowSize.set(windowSize);
|
||||
Thread t = ThreadTool.makeVThread("重路由计时器线程", checker);
|
||||
t.start();
|
||||
}
|
||||
|
||||
public long getWindowSize() {
|
||||
return windowSize.get();
|
||||
}
|
||||
|
||||
public void setWindowSize(long windowSize) {
|
||||
this.windowSize.set(windowSize);
|
||||
}
|
||||
|
||||
public Consumer<V> getResendConsumer() {
|
||||
return resendConsumer;
|
||||
}
|
||||
|
||||
public void setResendConsumer(Consumer<V> resendConsumer) {
|
||||
this.resendConsumer = resendConsumer;
|
||||
}
|
||||
|
||||
public Consumer<V> getClosingConsumer() {
|
||||
return closingConsumer;
|
||||
}
|
||||
|
||||
public void setClosingConsumer(Consumer<V> closingConsumer) {
|
||||
this.closingConsumer = closingConsumer;
|
||||
}
|
||||
|
||||
public long getSendmapWindowUsed() {
|
||||
return sendmapWindowUsed.get();
|
||||
}
|
||||
|
||||
public int getHeaderCalibrate() {
|
||||
return headerCalibrate;
|
||||
}
|
||||
|
||||
public CongestionAlgorithm getAlgorithm() {
|
||||
return algorithm;
|
||||
}
|
||||
|
||||
public boolean isCongress(IPv6Packet iPv6Packet, double scale) {
|
||||
boolean congress = sendmapWindowUsed.get() > windowSize.get() * scale;
|
||||
return congress;
|
||||
|
||||
}
|
||||
|
||||
public void put(K sequence, V packet) {
|
||||
SendItem<V> newitem = new SendItem<V>(packet);
|
||||
SendItem<V> old = sendMap.put(sequence, newitem);
|
||||
long dx=0;
|
||||
if (old != null) {
|
||||
dx=-(old.getPacketLength() + headerCalibrate);
|
||||
}
|
||||
long newWindow =sendmapWindowUsed.addAndGet(newitem.getPacketLength() + headerCalibrate+dx);
|
||||
algorithm.setCurrentWindowUsed(newWindow);
|
||||
}
|
||||
public V ack(K sequence) {
|
||||
return ack(sequence,false);
|
||||
}
|
||||
|
||||
public V ack(K sequence,boolean ecn) {
|
||||
SendItem<V> ipv6;
|
||||
if ((ipv6 = sendMap.remove(sequence)) != null) {
|
||||
long newWindow =sendmapWindowUsed.addAndGet(-(ipv6.getPacketLength() + headerCalibrate));
|
||||
algorithm.setCurrentWindowUsed(newWindow);
|
||||
long RTTC = System.nanoTime() - ipv6.getSendtime();
|
||||
algorithm.putAck(ipv6.getPacketLength(), RTTC, ecn);
|
||||
// System.out.println("remove:"+aseq.getSequence()+" "+sendMap.size());
|
||||
Lock lockx = lock;
|
||||
if (lockx != null) {
|
||||
lockx.lock();
|
||||
try {
|
||||
Condition cds = condition;
|
||||
if (cds != null) {
|
||||
cds.signalAll();
|
||||
}
|
||||
} finally {
|
||||
lockx.unlock();
|
||||
}
|
||||
}
|
||||
return ipv6.getPacket();
|
||||
} else {
|
||||
// System.out.println("miss:"+aseq.getSequence()+" "+sendMap.size());
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private class ResendChecker implements Runnable {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
while (!closed) {
|
||||
|
||||
Collection<SendItem<V>> vals = sendMap.values();
|
||||
for (Iterator<SendItem<V>> iterator = vals.iterator(); iterator.hasNext();) {
|
||||
SendItem<V> sitm = (SendItem<V>) iterator.next();
|
||||
long current = System.nanoTime();
|
||||
if (current - sitm.getSendtime() > algorithm.getRTO()) {
|
||||
if (removeAtResend) {
|
||||
iterator.remove();
|
||||
long newWindow = sendmapWindowUsed.addAndGet((int) -(sitm.getPacketLength() + headerCalibrate));
|
||||
algorithm.setCurrentWindowUsed(newWindow);
|
||||
} else {
|
||||
sitm.setSendtime(current);
|
||||
}
|
||||
V pktr = sitm.getPacket();
|
||||
if (resendConsumer != null)
|
||||
resendConsumer.accept(pktr);
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
Thread.sleep(2);
|
||||
} catch (InterruptedException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
// 清理资源
|
||||
if (closingConsumer != null) {
|
||||
for (SendItem<V> item : sendMap.values()) {
|
||||
closingConsumer.accept(item.getPacket());
|
||||
}
|
||||
}
|
||||
sendMap.clear();
|
||||
sendmapWindowUsed.set(0);
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isClosed() {
|
||||
return closed;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
closed = true;
|
||||
}
|
||||
|
||||
public void setCongressCondition(Lock lock2, Condition condition2) {
|
||||
this.lock = lock2;
|
||||
this.condition = condition2;
|
||||
}
|
||||
|
||||
public Map<K, SendItem<V>> getSendmap() {
|
||||
return sendMap;
|
||||
}
|
||||
|
||||
public int getWindowPacketCount() {
|
||||
return sendMap.size();
|
||||
}
|
||||
|
||||
public boolean isEmpty() {
|
||||
return sendMap.isEmpty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "SendPacketSlidingWindow [sendmapWindowUsed=" + sendmapWindowUsed + ", windowSize=" + windowSize + "]";
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,218 @@
|
||||
package org.kne.cloud.network.congestion;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
public class Vegas2CongestionAlgorithm implements CongestionAlgorithm,DetnetCongestionAlgorithm {
|
||||
|
||||
// 可配置的时延膨胀比上下限 - Vegas2.0核心参数
|
||||
private double delayUpperBound = 1.20; // RTT上限为基础RTT的1.2倍
|
||||
private double delayLowerBound = 1.15; // RTT下限为基础RTT的1.05倍
|
||||
|
||||
private static final long BURST_TIME = 2000000L;
|
||||
private static final long BIAS = 500000L;
|
||||
private static long MIN_SPEED = 256 * 1024L;
|
||||
private static long MIN_WINDOW = 16384;
|
||||
|
||||
// 新增:窗口调整参数(可基于当前窗口大小动态调整)
|
||||
private static final double WINDOW_ADJUST_FACTOR = 0.01; // 1%的窗口调整幅度
|
||||
|
||||
private Consumer<Long> windowControlConsumer;
|
||||
private BiConsumer<Long, Long> speedControlConsumer;
|
||||
|
||||
private AtomicBoolean firstUpdate = new AtomicBoolean(true);
|
||||
private long MIN_RTTVAR = 30000000L;
|
||||
private volatile long RTTMin = 1000000000L; // BaseRTT
|
||||
private volatile long RTTVar = 1000000000L;
|
||||
private volatile long RTTAvg = 1000000000L;
|
||||
private volatile long RTTAvg2 = 1000000000L; // 当前平滑RTT
|
||||
private volatile long RTTTotal = 0;
|
||||
private volatile long RTTCount = 0;
|
||||
private volatile long RTO = 1000000000L;
|
||||
|
||||
private volatile long window = MIN_WINDOW;
|
||||
private volatile long window2 = MIN_WINDOW;
|
||||
private volatile long speed = MIN_SPEED;
|
||||
private volatile long maxwindow = 99999999999L;
|
||||
private long congressSpeed=MIN_SPEED;
|
||||
|
||||
public Vegas2CongestionAlgorithm() {
|
||||
reset();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void reset() {
|
||||
window = MIN_WINDOW;
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
speed = MIN_SPEED;
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setWindowControlConsumer(Consumer<Long> windowControlConsumer) {
|
||||
this.windowControlConsumer = windowControlConsumer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setSpeedControlConsumer(BiConsumer<Long, Long> speedControlConsumer) {
|
||||
this.speedControlConsumer = speedControlConsumer;
|
||||
}
|
||||
|
||||
private long RTTstartTime = System.nanoTime();
|
||||
|
||||
@Override
|
||||
public synchronized void putAck(long packetSize, long latencyns, boolean ecn) {
|
||||
// 处理ECN信号 - 将其视为强烈的拥塞信号
|
||||
if (ecn) {
|
||||
// 当收到ECN时,更激进地减少窗口
|
||||
window2 = Math.max(0, window2 * 3 / 4);
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window2);
|
||||
}
|
||||
}
|
||||
|
||||
// 更新最小RTT(BaseRTT)
|
||||
if (latencyns <= RTTMin) {
|
||||
RTTMin = latencyns;
|
||||
} else {
|
||||
// 缓慢适应:当网络路径真正变化时,BaseRTT应能缓慢更新
|
||||
// 这里使用极慢的衰减因子,只有在持续观测到更低RTT时才快速更新
|
||||
RTTMin = (RTTMin * 9999 + latencyns) / 10000;
|
||||
}
|
||||
|
||||
// 更新平滑RTT估计
|
||||
if (firstUpdate.compareAndSet(true, false)) {
|
||||
RTTAvg = latencyns;
|
||||
RTTVar = latencyns / 2;
|
||||
RTTAvg2 = latencyns;
|
||||
} else {
|
||||
RTTVar = (RTTVar * 3 + Math.abs(RTTAvg - latencyns)) / 4;
|
||||
RTTAvg = (RTTAvg * 7 + latencyns) / 8;
|
||||
}
|
||||
|
||||
RTO = RTTAvg + Math.max(MIN_RTTVAR, RTTVar * 4);
|
||||
|
||||
RTTTotal += latencyns;
|
||||
RTTCount++;
|
||||
|
||||
long curr = System.nanoTime();
|
||||
// 每个RTT周期调整一次窗口(使用当前估计的RTT)
|
||||
if (curr - RTTstartTime > Math.max(BIAS, RTTMin)) {
|
||||
if (RTTCount > 0) {
|
||||
long currentRTT = RTTTotal / RTTCount;
|
||||
RTTAvg2 = (RTTAvg2 * 3 + currentRTT) / 4; // 平滑当前RTT
|
||||
RTTTotal = 0;
|
||||
RTTCount = 0;
|
||||
}
|
||||
|
||||
// ================== Vegas2.0核心逻辑 ==================
|
||||
// 1. 计算时延膨胀比
|
||||
double delayRatio = (double) RTTAvg2 / (double) Math.max(BIAS, RTTMin);
|
||||
|
||||
// 2. 基于比值的窗口调整(取代原来的基于差值的调整)
|
||||
long adjustStep = Math.max(512, (long)(window2 * WINDOW_ADJUST_FACTOR));
|
||||
|
||||
if (delayRatio > delayUpperBound) {
|
||||
// 时延过高:减小窗口,减少幅度与超标程度成正比
|
||||
double exceedRatio = delayRatio / delayUpperBound;
|
||||
window2 -= (long)(adjustStep * exceedRatio);
|
||||
} else if (delayRatio < delayLowerBound) {
|
||||
// 时延过低:增大窗口,增加幅度与低于目标程度成正比
|
||||
double belowRatio = delayLowerBound / delayRatio;
|
||||
window2 += (long)(adjustStep * belowRatio);
|
||||
} else {
|
||||
// 在理想区间内:微调以维持稳定
|
||||
// 计算目标RTT = BaseRTT × 目标膨胀比
|
||||
long targetRTT = (long)(Math.max(BIAS, RTTMin) * (delayLowerBound + delayUpperBound) / 2.0);
|
||||
|
||||
if (RTTAvg2 > targetRTT) {
|
||||
// 略高于目标:小幅减少
|
||||
window2 = Math.max(0, window2 - 256);
|
||||
} else if (RTTAvg2 < targetRTT) {
|
||||
// 略低于目标:小幅增加
|
||||
window2 = Math.min(maxwindow, window2 + 256);
|
||||
}
|
||||
// 非常接近目标:保持窗口不变
|
||||
}
|
||||
// ===================================================
|
||||
|
||||
// 窗口边界检查
|
||||
if (window2 >= maxwindow) {
|
||||
window2 = maxwindow;
|
||||
}
|
||||
if (window2 < 0) {
|
||||
window2 = 0;
|
||||
}
|
||||
|
||||
RTTstartTime = curr;
|
||||
|
||||
if (bandwidth >= congressSpeed) {
|
||||
congressSpeed = bandwidth;
|
||||
//congressSpeed = (congressSpeed + bandwidth) / 2;
|
||||
} else {
|
||||
congressSpeed = (congressSpeed * 49 + bandwidth) / 50;
|
||||
}
|
||||
|
||||
// 速度计算:基于当前窗口和RTT,考虑一个安全边界
|
||||
// 使用目标膨胀比而非固定1.2倍,确保与窗口控制逻辑一致
|
||||
window=window2+MIN_WINDOW;
|
||||
speed=(long) (congressSpeed * delayUpperBound)+MIN_SPEED ;
|
||||
|
||||
if (windowControlConsumer != null) {
|
||||
windowControlConsumer.accept(window);
|
||||
}
|
||||
if (speedControlConsumer != null) {
|
||||
speedControlConsumer.accept(speed, BURST_TIME);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void putLoss(long packetSize, long lossns, int losscounter) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getRTO() {
|
||||
return RTO;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCurrentWindowUsed(long used) {
|
||||
this.maxwindow = used * 16 + 65536;
|
||||
}
|
||||
private long bandwidth;
|
||||
|
||||
@Override
|
||||
public void setCurrentBandwidth(long bandwidth) {
|
||||
// 可选:基于已知带宽信息调整参数
|
||||
this.bandwidth=bandwidth;
|
||||
}
|
||||
|
||||
// 新增方法:获取当前状态信息(用于监控和调试)
|
||||
public double getCurrentDelayRatio() {
|
||||
return (double) RTTAvg2 / (double) Math.max(BIAS, RTTMin);
|
||||
}
|
||||
|
||||
public long getBaseRTT() {
|
||||
return RTTMin;
|
||||
}
|
||||
|
||||
public long getCurrentRTT() {
|
||||
return RTTAvg2;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setUpperDelayBound(double upper) {
|
||||
delayUpperBound=upper;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setLowerDelayBound(double lower) {
|
||||
delayLowerBound=lower;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user