У меня есть один поток, который делает сокет записи пакетов, считываемых из LinkedBlockquequeue, с другой веткой: < /p>
import java.io.ByteArrayOutputStream;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
public class SimpleSender implements Runnable {
private final Socket socket;
private final LinkedBlockingQueue queue = new LinkedBlockingQueue();
private volatile boolean running = true;
public SimpleSender(Socket socket) {
this.socket = socket;
}
@Override
public void run() {
while (running) {
try {
byte[] msg = queue.poll(1, TimeUnit.SECONDS);
if (msg != null) {
sendPacket(msg);
}
} catch (Exception e) {
e.printStackTrace();
running = false;
}
}
}
private void sendPacket(byte[] data) throws Exception {
try (ByteArrayOutputStream buffer = new ByteArrayOutputStream()) {
// prefix with length (2 bytes, little endian)
buffer.write(data.length & 0xff);
buffer.write((data.length >> 8) & 0xff);
buffer.write(data);
OutputStream out = socket.getOutputStream();
out.write(buffer.toByteArray());
out.flush();
}
}
public boolean offer(byte[] data) {
return queue.offer(data);
}
public void stop() {
running = false;
}
public static void main(String[] args) throws Exception {
try (Socket socket = new Socket()) {
socket.connect(new InetSocketAddress("127.0.0.1", 12345));
SimpleSender sender = new SimpleSender(socket);
Thread senderThread = new Thread(sender, "SenderThread");
senderThread.start();
sender.offer("Hello".getBytes());
sender.offer("World".getBytes());
Thread.sleep(2000);
sender.stop();
}
}
}
< /code>
Когда сеть медленная или медленно читается, мой поток может увядать. Тем временем мой LinkedBlockingQueue становится слишком большим, и мне нужно быстро вызвать запас. Я знаю только розетку и Socketcheannel, но последний кажется более сложным, и мне просто нужно отправлять один пакет за раз.
Подробнее здесь: https://stackoverflow.com/questions/797 ... -exception