优化调度算法,避免路由震荡

This commit is contained in:
Administrator
2022-12-30 20:01:14 +08:00
parent c791720f95
commit 1af088a29d
16 changed files with 201 additions and 67 deletions
+2 -2
View File
@@ -1,8 +1,8 @@
package org.kne.cloud.network.klalb;
public class Consts {
public static final int BLOCKSIZE=32768;
public static final long PINGTIMENS=2000000000L;
public static final int BLOCKSIZE=16384;
public static final long PINGTIMENS=1000000000L;
public static final double A = 0.125;
public static final int SO_TIMEOUT = 20000;
public static final long SN_KEEP = 60000000000L;
@@ -9,9 +9,13 @@ import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Queue;
import java.util.UUID;
import java.util.Vector;
import java.util.WeakHashMap;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.atomic.AtomicBoolean;
import org.kne.cloud.network.mport.ServiceElement;
@@ -168,6 +172,15 @@ public class IOThreadManager {
e.printStackTrace();
} finally {
klc.getRemotetcps().remove(tc);
Queue<KLALBBlock> sq=tc.getSendDeque();
KLALBBlock tmp;
while((tmp=sq.poll())!=null) {
try {
klc.submitDataBlock(tmp);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
@@ -11,7 +11,7 @@ public class KLALBBlock implements Comparable<KLALBBlock>{
private static AtomicLong sng=new AtomicLong(1);
public long sn;//每个数据包的唯一编号
public long cuid;//用于识别数据包的stream ID号
public int cuid;//用于识别数据包的stream ID号
public long number;//数据包的编号,用于排序
public byte[]data;//数据内容
public int command;//功能标记 0=PING 1=PONG 2=SYN 3=RST
@@ -210,7 +210,7 @@ public class KLALBClientGUI extends XFrame {
for (Iterator<ServiceElement> iterator = lse.iterator(); iterator.hasNext();) {
ServiceElement serviceElement = (ServiceElement) iterator.next();
ProxyIPPortProperties pippp = new ProxyIPPortProperties(pdis, serviceElement.name,
"127.0.0.1", serviceElement.ipport.getPort());
"0.0.0.0", serviceElement.ipport.getPort());
try {
kc.open(pippp.getBindPort(), serviceElement, pippp.getBindIP());
} catch (BindException e) {
+27 -11
View File
@@ -26,6 +26,7 @@ import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
@@ -37,7 +38,7 @@ import org.kne.cloud.network.mport.ThreadTool;
public class KLALBCore {
private volatile int cacheblocks;
private Map<Long,LocalTCPConnection> localtcps = new ConcurrentHashMap<>();
private Map<Integer, LocalTCPConnection> localtcps = new ConcurrentHashMap<>();
private List<RemoteTCPConnection> remotetcps = new Vector<>();
private Predicate<KLALBBlock>acceptSYN;
@@ -81,7 +82,7 @@ public class KLALBCore {
this.acceptSYN = acceptSYN;
}
public Map<Long, LocalTCPConnection> getLocaltcps() {
public Map<Integer, LocalTCPConnection> getLocaltcps() {
return localtcps;
}
@@ -103,21 +104,28 @@ public class KLALBCore {
pdb.pingtime=System.nanoTime();
tc.sendBlock(pdb);
}
if(!tc.getSendDeque().isEmpty())
if(!tc.getNDSendDeque().isEmpty()||!tc.getSendDeque().isEmpty())
break;
Thread.sleep(1);
synchronized (tc.getSendDeque()) {
tc.getSendDeque().wait(10);
}
}
KLALBBlock k0=tc.getNDSendDeque().poll();
if(k0!=null) {
tc.sendBlock(k0,tc.getNDSendDeque().isEmpty());
}else {
KLALBBlock k=tc.getSendDeque().poll();
if(k!=null) {
//if(localtcps.containsKey(k.cuid)||k.cuid.equals(ZERO_UUID))
tc.sendBlock(k);
tc.sendBlock(k,tc.getSendDeque().isEmpty());
}
}
}
public void remoteReceive(RemoteTCPConnection tc) throws IOException {
KLALBBlock brc=null;
for(;;) {
brc=tc.receiveBlock();
if(brc.number!=0) {
if(brc.number!=0||brc.sn==0) {
break;
}
SN s=new SN(brc.sn);
@@ -148,11 +156,11 @@ public class KLALBCore {
pdb.number=0;
pdb.command=1;
pdb.pingtime=brc.pingtime;
tc.getSendDeque().addFirst(pdb);
tc.getNDSendDeque().offer(pdb);
}else if(brc.command==1) {
long cur=System.nanoTime();
long delay=(cur-brc.pingtime)/2;
tc.setDelay(delay);
tc.setDelayAvg(delay);
}
}
}else {
@@ -215,11 +223,15 @@ public class KLALBCore {
}
}
});
//System.out.println(remotetcps);
for (Iterator<RemoteTCPConnection> iterator = remotetcps.iterator(); iterator.hasNext();) {
RemoteTCPConnection remoteTCPConnection = (RemoteTCPConnection) iterator.next();
BlockingDeque<KLALBBlock> bdq=remoteTCPConnection.getSendDeque();
if(bdq.size()<2) {
Queue<KLALBBlock> bdq=remoteTCPConnection.getSendDeque();
if(bdq.size()<1) {
bdq.add(kb);
synchronized (bdq) {
bdq.notifyAll();
}
return;
}
}
@@ -232,7 +244,11 @@ public class KLALBCore {
synchronized (remotetcps) {
for (Iterator<RemoteTCPConnection> iterator = remotetcps.iterator(); iterator.hasNext();) {
RemoteTCPConnection rmt = (RemoteTCPConnection) iterator.next();
rmt.getSendDeque().addFirst(kb);
Queue<KLALBBlock> sq = rmt.getNDSendDeque();
sq.offer(kb);
synchronized (sq) {
sq.notifyAll();
}
}
}
}
@@ -5,15 +5,16 @@ import java.net.Socket;
import java.util.*;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import org.kne.cloud.network.mport.ServiceElement;
public class LocalTCPConnection extends TCPConnection {
private static AtomicLong sng=new AtomicLong(1);
private static AtomicInteger sng=new AtomicInteger(1);
private long cuid;
private int cuid;
private ServiceElement serviceElement;
@@ -26,14 +27,14 @@ private volatile long outputcount = 1;
super(s);
cuid=sng.getAndIncrement();
}
public LocalTCPConnection(Socket s,long uid) throws IOException {
public LocalTCPConnection(Socket s,int uid) throws IOException {
super(s);
cuid=uid;
}
public long getCuid() {
public int getCuid() {
return cuid;
}
public void setCuid(long cuid) {
public void setCuid(int cuid) {
this.cuid = cuid;
}
public ServiceElement getServiceElement() {
@@ -10,9 +10,13 @@ import java.io.OutputStream;
import java.io.StreamCorruptedException;
import java.net.Socket;
import java.net.UnknownHostException;
import java.util.Queue;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.zip.GZIPInputStream;
import java.util.zip.GZIPOutputStream;
@@ -54,10 +58,15 @@ public class RemoteTCPConnection extends TCPConnection {
}
dout = new DataOutputStream(new AddOutputStream(new BufferedOutputStream(outm)) );
din = new DataInputStream(new AddInputStream(new BufferedInputStream(inm)));
dout = new DataOutputStream(new BufferedOutputStream(outm,Consts.BLOCKSIZE*2+1) );
din = new DataInputStream(new BufferedInputStream(inm,Consts.BLOCKSIZE*2+1));
}
@Override
public String toString() {
return "RemoteTCPConnection [tunnel=" + tunnel + "]";
}
public RemoteTCPConnection(Tunnel t, Socket s) throws UnknownHostException, IOException {
super(s);
tunnel=t;
@@ -87,7 +96,17 @@ public class RemoteTCPConnection extends TCPConnection {
}
}
long odelay=Long.MAX_VALUE;
public void setDelayAvg(long delay) {
if(odelay==Long.MAX_VALUE) {
odelay=delay;
setDelay(odelay);
}else {
odelay=(delay+odelay*100)/101;
setDelay(odelay);
}
}
@Override
public void close() {
if (tunnel != null) {
@@ -103,20 +122,27 @@ public class RemoteTCPConnection extends TCPConnection {
}
super.finalize();
}
public void sendBlock(KLALBBlock data) throws IOException {
dout.writeLong(data.sn);
dout.writeLong(data.cuid);
sendBlock(data,true);
}
public void sendBlock(KLALBBlock data,boolean isflush) throws IOException {
dout.writeInt(data.cuid);
dout.writeLong(data.number);
if (data.number > 0) {
if (data.data == null) {
dout.writeInt(-1);
dout.writeShort(-1);
} else {
dout.writeInt(data.data.length);
dout.writeShort(data.data.length);
dout.write(data.data);
}
} else if (data.number == 0) {
dout.write(data.command);
if(data.command!=0&&data.command!=1) {
dout.writeLong(data.sn);
}else {
}
switch (data.command) {
case 0:
case 1:
@@ -126,11 +152,14 @@ public class RemoteTCPConnection extends TCPConnection {
dout.writeUTF(data.lservice);
dout.writeUTF(data.lipport.toString());
break;
default:
break;
}
} else {
dout.writeInt(data.cacheused);
}
if(isflush)
dout.flush();
//System.out.println("SEND:" + data);
@@ -138,11 +167,10 @@ public class RemoteTCPConnection extends TCPConnection {
public KLALBBlock receiveBlock() throws IOException {
KLALBBlock klb = new KLALBBlock();
klb.sn = din.readLong();
klb.cuid = din.readLong();
klb.cuid = din.readInt();
klb.number = din.readLong();
if (klb.number > 0) {
int size = din.readInt();
int size = din.readShort();
if (size == -1) {
klb.data = null;
} else {
@@ -152,6 +180,11 @@ public class RemoteTCPConnection extends TCPConnection {
}
} else if (klb.number == 0) {
klb.command = din.read();
if(klb.command!=0&&klb.command!=1) {
klb.sn = din.readLong();
}else {
klb.sn=0;
}
switch (klb.command) {
case 0:
case 1:
@@ -160,6 +193,9 @@ public class RemoteTCPConnection extends TCPConnection {
case 2:
klb.lservice = din.readUTF();
klb.lipport = new IPPort(din.readUTF());
break;
default:
break;
}
} else {
klb.cacheused = din.readInt();
@@ -180,10 +216,16 @@ public class RemoteTCPConnection extends TCPConnection {
return false;
}
}
private Queue<KLALBBlock> NDsendDeque = new ConcurrentLinkedQueue<>();
private BlockingDeque<KLALBBlock> sendDeque = new LinkedBlockingDeque<>();
public Queue<KLALBBlock> getNDSendDeque() {
return NDsendDeque;
}
private Queue<KLALBBlock> sendDeque = new ConcurrentLinkedQueue<>();
public BlockingDeque<KLALBBlock> getSendDeque() {
public Queue<KLALBBlock> getSendDeque() {
return sendDeque;
}
+16 -1
View File
@@ -4,7 +4,9 @@ import java.io.BufferedInputStream;
import java.io.BufferedOutputStream;
import java.io.DataInputStream;
import java.io.DataOutputStream;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.PrintStream;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.net.UnknownHostException;
@@ -83,7 +85,7 @@ public class Tunnel{
@Override
public String toString() {
return getName()+"$"+getIpport();
return getName()+"$"+getIpport()+" "+delay/1000000;
}
public Socket connectClientSocket() throws UnknownHostException, IOException {
Socket socket=new Socket();
@@ -138,8 +140,21 @@ public class Tunnel{
public AtomicLong getIM() {
return im;
}
PrintStream ndbg=null;
private void initcdbg() {
try {
ndbg=new PrintStream(name+".csv");
} catch (FileNotFoundException e) {
e.printStackTrace();
}
}
public void updateTraffics(long timems) {
tlr.setTraffic(om.get(), im.get(), om.get()-om1, im.get()-im1);
if(ndbg==null)
initcdbg();
ndbg.println((om.get()-om1)+ (im.get()-im1)+","+delay);
om1=om.get();
im1=im.get();
}
@@ -1,14 +1,25 @@
package org.kne.cloud.network.mport;
import java.lang.reflect.Method;
public class ThreadTool {
public static boolean first=true;
public static boolean forceD=false;
public static Thread makeVThreadIfSupport(String name,Runnable r) {
if(forceD)
return new Thread(r, name);
try { return new Thread(r, name);
try { //return new Thread(r, name);
Class<?> c=Class.forName("java.lang.Thread");
Method m= c.getDeclaredMethod("ofVirtual", null);
Object o=m.invoke(null, null);
Class<?> cb=Class.forName("java.lang.Thread$Builder");
Method mn= cb.getDeclaredMethod("name", String.class);
mn.invoke(o, name);
Method mu= cb.getDeclaredMethod("unstarted", Runnable.class);
return (Thread) mu.invoke(o, r);
//return Thread.ofVirtual().name(name).unstarted(r);
}catch(Throwable e) {
//e.printStackTrace();
if(first) {
System.out.println("请使用java19以上版本以提高性能!");
first=false;
+1
View File
@@ -568,5 +568,6 @@ public class XFrame extends JFrame {
}
});
}
}
+11
View File
@@ -12,6 +12,9 @@ import java.awt.event.ComponentEvent;
import java.awt.event.ComponentListener;
import java.awt.event.ContainerEvent;
import java.awt.event.ContainerListener;
import java.awt.event.MouseWheelEvent;
import java.awt.event.MouseWheelListener;
import javax.swing.JScrollBar;
public class YScrollPane extends JPanel {
@@ -109,6 +112,14 @@ settings.addContainerListener(new ContainerListener() {
}
});
settings.addMouseWheelListener(new MouseWheelListener() {
@Override
public void mouseWheelMoved(MouseWheelEvent e) {
int x=e.getWheelRotation()*20;
scrollBar.setValue(scrollBar.getValue()+x);
}
});
}
public JPanel getView() {