This commit is contained in:
2026-07-07 00:26:38 +08:00
parent 74e4360d71
commit ba61d3b5ae
28 changed files with 2411 additions and 0 deletions
@@ -0,0 +1,20 @@
package org.kne.cloud.network.congestion;
public class BBRCongestionAlgorithmFactory implements CongestionAlgorithmFactory {
@Override
public CongestionAlgorithm create() {
return new BBRCongestionAlgorithm();
}
@Override
public String getName() {
return "BBR";
}
@Override
public String toString() {
return getName();
}
}
@@ -0,0 +1,7 @@
package org.kne.cloud.network.congestion;
public interface CongestionAlgorithmFactory {
public CongestionAlgorithm create();
public String getName();
}
@@ -0,0 +1,18 @@
package org.kne.cloud.network.congestion;
import java.util.HashMap;
import java.util.Map;
public class CongestionAlgorithms {
private static Map<String,CongestionAlgorithmFactory>regs=new HashMap();
static {
register(new BBRCongestionAlgorithmFactory());
register(new Vegas2CongestionAlgorithmFactory());
}
public static void register(CongestionAlgorithmFactory algorithm) {
regs.put(algorithm.getName(), algorithm);
}
public static CongestionAlgorithm get(String name) {
return regs.get(name).create();
}
}
@@ -0,0 +1,207 @@
package org.kne.cloud.network.congestion;
import org.jctools.queues.MpscArrayQueue;
import org.kne.cloud.network.ThreadTool;
import org.kne.opencl64.Releaser;
import java.lang.ref.Cleaner;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.LockSupport;
import java.util.function.Consumer;
/**
* 基于 MPSC 队列的高性能消息批量发送器
*/
public class MpscMessageBatcher<T> implements MessageBatcher<T> {
private static final Cleaner clr = Cleaner.create();
private final MpscMessageBatcher0<T> impl;
private final MpscMessageBatcherReleaser<T> releaser;
public MpscMessageBatcher(int batchSize, long maxDelayNanos) {
this.impl = new MpscMessageBatcher0<>(batchSize, maxDelayNanos, 65536);
this.releaser = new MpscMessageBatcherReleaser<>(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 {
MpscMessageBatcher<Integer> batcher =
new MpscMessageBatcher<>(10, 1_000_000L);
batcher.setConsumer(System.out::println);
int n = 0;
for (;;) {
batcher.putMessage(n++);
batcher.putMessage(n++);
batcher.putMessage(n++);
batcher.putMessage(n++);
Thread.sleep(1);
}
}
}
class MpscMessageBatcherReleaser<T> extends Releaser<MpscMessageBatcher0<T>> {
public MpscMessageBatcherReleaser(MpscMessageBatcher0<T> resource) {
super(resource);
}
@Override
protected void release(MpscMessageBatcher0<T> resource) {
resource.close();
}
}
class MpscMessageBatcher0<T> implements Runnable, MessageBatcher<T> {
private final MpscArrayQueue<T> queue;
private final int batchSize;
private final long maxDelay;
private volatile Consumer<List<T>> consumer;
private final Thread batchThread;
private final AtomicBoolean closed = new AtomicBoolean(false);
private long firstTime;
public MpscMessageBatcher0(int batchSize, long maxDelay, int queueCapacity) {
this.batchSize = batchSize;
this.maxDelay = maxDelay;
this.queue = new MpscArrayQueue<>(queueCapacity);
this.batchThread = ThreadTool.makeVDaemonThread("MpscMessageBatcher", this);
this.batchThread.start();
}
@Override
public void setConsumer(Consumer<List<T>> consumer) {
this.consumer = consumer;
}
@Override
public void putMessage(T message) {
if (message == null || closed.get()) return;
// 无锁入队
while (!queue.offer(message)) {
// 队列满时,主动尝试 check 一次(自旋等待)
Thread.yield();
}
check(false);
}
@Override
public void putMessages(List<T> messages) {
for (T msg : messages) {
putMessage(msg);
}
}
@Override
public void run() {
while(!closed.get()) {
check(false);
// 更精确的等待策略
long sleepTimeNanos = calculateSleepTime();
if (sleepTimeNanos > 0) {
LockSupport.parkNanos(sleepTimeNanos);
} else {
// 避免忙等待
LockSupport.parkNanos(1_000_000L); // 1ms
}
}
}
private void check(boolean force) {
long currTime=System.nanoTime();
if(queue.size()>=batchSize||(currTime-firstTime)>maxDelay||force) {
if(consumer!=null) {
ArrayList<T>batch=new ArrayList<>(batchSize);
queue.drain(batch::add, batchSize);
if(!batch.isEmpty())
try {
consumer.accept(batch);
}catch(Exception e) {
e.printStackTrace();
}
firstTime=currTime;
}
}
}
/**
* 计算需要等待的时间
* @return 等待时间(纳秒)
*/
private long calculateSleepTime() {
long elapsed = System.nanoTime() - firstTime;
long remaining = maxDelay - elapsed;
if (remaining <= 0) {
return 0; // 立即处理
}
// 返回剩余时间或10ms中的较小值
return Math.min(remaining, 10_000_000L);
}
@Override
public void flush() {
check(true);
}
@Override
public int getQueueSize() {
return queue.size();
}
@Override
public boolean isClosed() {
return closed.get();
}
@Override
public void close() {
closed.set(true);
LockSupport.unpark(batchThread);
flush(); // 确保剩余消息被处理
}
@Override
public int getBatchSize() {
return batchSize;
}
@Override
public long getMaxDelay() {
return maxDelay;
}
}
@@ -0,0 +1,146 @@
package org.kne.cloud.network.congestion;
import java.io.Closeable;
import java.io.IOException;
import java.net.SocketException;
import java.net.SocketTimeoutException;
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.CountDownLatch;
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.DATATPacket;
import org.kne.cloud.network.klalb.SendItem;
import org.kne.cloud.network.kltp.KLTPPacket;
import org.kne.concurrent.ThreadParker;
public class ReceivePacketSlidingWindow<K, V extends NetworkPacket> implements Closeable, AutoCloseable{
private ThreadParker tp=new ThreadParker();
private Map<K,SendItem<V>> recvMap = new ConcurrentHashMap<K,SendItem<V>>();
private AtomicLong recvWindowUsed=new AtomicLong(0);
//private LongAdder recvWindowUsed=new LongAdder();
private long recvWindowSize;
private int headerCalibrate = 0;
private volatile boolean closed = false;
public ReceivePacketSlidingWindow(long recvWindowSize) {
super();
this.recvWindowSize = recvWindowSize;
}
public ReceivePacketSlidingWindow(long recvWindowSize, int headerCalibrate) {
super();
this.recvWindowSize = recvWindowSize;
this.headerCalibrate = headerCalibrate;
}
public void setHeaderCalibrate(int headerCalibrate) {
this.headerCalibrate = headerCalibrate;
}
public int getHeaderCalibrate() {
return headerCalibrate;
}
public void put(K sequence, V packet) {
SendItem<V> newitem = new SendItem<V>(packet,headerCalibrate);
SendItem<V> old = recvMap.putIfAbsent(sequence, newitem);
long dx=0;
if (old != null) {
dx=-old.getPacketLength() ;
}
recvWindowUsed.addAndGet(newitem.getPacketLength() +dx);
tp.unpark();
}
public V poll(K key) {
SendItem<V> dtp2 = recvMap.remove(key);
if (dtp2 != null) {
recvWindowUsed.addAndGet(-dtp2.getPacketLength());
return dtp2.getPacket();
}else {
return null;
}
}
public V take(K key) throws SocketException {
while(true) {
V val=poll(key);
if(val!=null) {
return val;
}
if(isClosed()) {
throw new SocketException("Receive window closed!");
}tp.parkNanos(1000000L);
}
}
public V take(K key,long timeoutNanos) throws SocketException, SocketTimeoutException {
long start=System.nanoTime();
while(true) {
if(System.nanoTime()-start>timeoutNanos) {
throw new SocketTimeoutException("Receive time out:"+(System.nanoTime()-start)+">"+timeoutNanos);
}
V val=poll(key);
if(val!=null) {
return val;
}
if(isClosed()) {
throw new SocketException("Receive window closed!");
}tp.parkNanos(1000000L);
}
}
public long getRecvWindowUsed() {
return recvWindowUsed.get();
}
public long getRecvWindowAvaliable() {
return Math.max(0, recvWindowSize-recvWindowUsed.get());
}
public Map<K, SendItem<V>> getRecvMap() {
return recvMap;
}
public int getWindowPacketCount() {
return recvMap.size();
}
public boolean isEmpty() {
return recvMap.isEmpty();
}
public boolean isClosed() {
return closed;
}
@Override
public void close() {
closed = true;
}
@Override
public String toString() {
return "ReceivePacketSlidingWindow [recvWindowUsed=" + recvWindowUsed + ", recvWindowSize=" + recvWindowSize
+ ", closed=" + closed + "]";
}
}
@@ -0,0 +1,20 @@
package org.kne.cloud.network.congestion;
public class Vegas2CongestionAlgorithmFactory implements CongestionAlgorithmFactory {
@Override
public CongestionAlgorithm create() {
return new Vegas2CongestionAlgorithm();
}
@Override
public String getName() {
return "Vegas2";
}
@Override
public String toString() {
return getName();
}
}