【问题标题】:How to avoid writing too much data to output and which blocks the thread?如何避免写入太多数据输出以及哪些阻塞线程?
【发布时间】:2015-07-12 06:42:14
【问题描述】:

在一个由线程池运行的任务中,我想向远程写入一堆字符串,并且有一个标志指示该任务是否已被取消。

我正在使用以下代码来确保我可以尽快停止:

public void sendToRemote(Iterator<String> data, OutputStream remote) {
    try {
        System.out.println("##### staring sending")
        while(!this.cancelled && data.hasNext()) {
            remote.write(data.next())
        }
        System.out.println("##### finished sending")
        System.out.println("")
    } catch(Throwable e) {
        e.printStackTrace();
    } finally {
        remote.close();
    }
}

我发现有时候,如果我给这个方法一个非常大的数据(或无限迭代器),即使我稍后将this.cancelled设置为true,它也无法及时完成。代码好像被阻塞了,过了很长时间(1分钟左右),会出现如下错误:

java.net.SocketTimeoutException: write blocked too long

所以我猜可能是remote.write 方法可以在有太多数据要发送的情况下自行阻塞,但远程没有及时使用它。虽然我将this.cancelled设置为true,但是该方法在remote.write(data.next())行中被阻塞了很长时间,所以它没有机会检查this.cancelled的值并跳过循环。相反,它终于在很长一段时间后抛出了SocketTimeoutException

我的理解正确吗?如果是,如果要发送的数据过多,如何避免阻塞?

【问题讨论】:

  • 您可以使用 Channel(来自 java.nio.channels)而不是 OutputStream。那么write方法很容易被打断。
  • 谢谢,但就我而言,它必须是OutputStreamremote 参数来自第三方方法,我们无法更改它
  • Channels 类包含将 Streams 转换为 Channels 的方法,反之亦然。
  • 该频道仍会被屏蔽
  • 我从未见过java.net.SocketTimeoutException: write blocked too long。你在哪个平台?当套接字发送缓冲区已满时,写入应该无限期地阻塞。 真正的问题在于阅读端。为什么跟不上?

标签: java sockets blocking


【解决方案1】:

尝试简单地关闭远程 OutputStream。您的线程将以异常结束。

  1. 线程#1 忙于执行 sendToRemote();
  2. 线程#2 认为足够了,并关闭了遥控器。 (假设 OutPutStream 对象不是线程本地的,就像在全局中一样 参考某处)
  3. Thread#1 因异常而死亡 :)

编辑我在互联网上找到this

启用 linger 并将超时设置为特定秒数 将导致对 Socket.Close 的后续调用阻塞,直到所有数据 发送缓冲区中已发送或超时。

【讨论】:

  • 设置延迟超时确实会产生这种效果,但这有什么帮助呢?
  • 我们可以关闭套接字,而不必等待发送缓冲区清空。所以OP可以成功停止他的线程。
  • 糟糕。我错过了他提到的第三方限制。我的坏
【解决方案2】:

正确的解决方案可能是以某种方式使用 NIO。我已经评论了 Hadoop 是如何使用nio underneathhere 做到的。

但更简单的解决方案是在 Dexter 的answer 中。我还遇到了answer from EJP,他建议使用BufferedOutputStream 来控制数据何时流出。所以我将两者结合起来得到了如下所示的TimedOutputStream。它不能完全控制远程的输出缓冲(大部分由操作系统完成),但结合适当的缓冲区大小和写入超时至少可以提供一些控制(请参阅第二个程序以测试TimedOutputStream)。

我还没有完全测试过TimedOutputStream,所以请自己进行尽职调查。
编辑: 更新了写入方法以更好地关联缓冲区大小和写入超时,还调整了测试程序。添加了有关套接字输出流的非安全异步关闭的 cmets。

import java.io.*;
import java.util.concurrent.*;

/**
 * A {@link BufferedOutputStream} that sets time-out tasks on write operations 
 * (typically when the buffer is flushed). If a write timeout occurs, the underlying outputstream is closed
 * (which may not be appropriate when sockets are used, see also comments on {@link TimedOutputStream#interruptWriteOut}).
 * A {@link ScheduledThreadPoolExecutor} is required to schedule the time-out tasks.
 * This {@code ScheduledThreadPoolExecutor} should have {@link ScheduledThreadPoolExecutor#setRemoveOnCancelPolicy(boolean)}
 * set to {@code true} to prevent a huge task queue.
 * If no {@code ScheduledThreadPoolExecutor} is provided in the constructor, 
 * the executor is created and shutdown with the {@link #close()} method. 
 * @author vanOekel
 *
 */
public class TimedOutputStream extends FilterOutputStream {

    protected int timeoutMs = 50_000;
    protected final boolean closeExecutor;
    protected final ScheduledExecutorService executor;
    protected ScheduledFuture<?> timeoutTask;
    protected volatile boolean writeTimedout;
    protected volatile IOException writeTimeoutCloseException;

    /* *** new methods not in BufferedOutputStream *** */

    /**
     * Default timeout is 50 seconds.
     */
    public void setTimeoutMs(int timeoutMs) {
        this.timeoutMs = timeoutMs;
    }

    public int getTimeoutMs() {
        return timeoutMs;
    }

    public boolean isWriteTimeout() {
        return writeTimedout;
    }

    /**
     * If a write timeout occurs and closing the underlying output-stream caused an exception,
     * then this method will return a non-null IOException.
     */
    public IOException getWriteTimeoutCloseException() {
        return writeTimeoutCloseException;
    }

    public ScheduledExecutorService getScheduledExecutor() {
        return executor;
    }

    /**
     * See {@link BufferedOutputStream#close()}.
     */
    @Override
    public void close() throws IOException {

        try {
            super.close(); // calls flush via FilterOutputStream.
        } finally {
            if (closeExecutor) {
                executor.shutdownNow();
            }
        }
    }

    /* ** Mostly a copy of java.io.BufferedOutputStream and updated with time-out options. *** */

    protected byte buf[];
    protected int count;

    public TimedOutputStream(OutputStream out) {
        this(out, null);
    }

    public TimedOutputStream(OutputStream out, ScheduledExecutorService executor) {
        this(out, 8192, executor);
    }

    public TimedOutputStream(OutputStream out, int size) {
        this(out, size, null);
    }

    public TimedOutputStream(OutputStream out, int size, ScheduledExecutorService executor) {
        super(out);
        if (size <= 0) {
            throw new IllegalArgumentException("Buffer size <= 0");
        }
        if (executor == null) {
            this.executor = Executors.newScheduledThreadPool(1);
            ScheduledThreadPoolExecutor stp = (ScheduledThreadPoolExecutor) this.executor;
            stp.setRemoveOnCancelPolicy(true);
            closeExecutor = true;
        } else {
            this.executor = executor;
            closeExecutor = false;
        }
        buf = new byte[size];
    }

    /**
     * Flushbuffer is called by all the write-methods and "flush()".
     */
    protected void flushBuffer(boolean flushOut) throws IOException {

        if (count > 0 || flushOut) {
            timeoutTask = executor.schedule(new TimeoutTask(this), getTimeoutMs(), TimeUnit.MILLISECONDS);
            try {
                // long start = System.currentTimeMillis(); int len = count;
                if (count > 0) {
                    out.write(buf, 0, count);
                    count = 0;
                }
                if (flushOut) {
                    out.flush(); // in case out is also buffered, this will do the actual write.
                }
                // System.out.println(Thread.currentThread().getName() + " Write [" + len + "] " + (flushOut ? "and flush " : "") + "time: " + (System.currentTimeMillis() - start));
            } finally {
                timeoutTask.cancel(false);
            }
        }
    }

    protected class TimeoutTask implements Runnable {

        protected final TimedOutputStream tout;
        public TimeoutTask(TimedOutputStream tout) {
            this.tout = tout;
        }

        @Override public void run() {
            tout.interruptWriteOut();
        }
    }

    /**
     * Closes the outputstream after a write timeout. 
     * If sockets are used, calling {@link java.net.Socket#shutdownOutput()} is probably safer
     * since the behavior of an async close of the outputstream is undefined. 
     */
    protected void interruptWriteOut() {

        try {
            writeTimedout = true;
            out.close();
        } catch (IOException e) {
            writeTimeoutCloseException = e;
        }
    }

    /**
     * See {@link BufferedOutputStream#write(int b)}
     */
    @Override
    public void write(int b) throws IOException {

        if (count >= buf.length) {
            flushBuffer(false);
        }
        buf[count++] = (byte)b;
    }

    /**
     * Like {@link BufferedOutputStream#write(byte[], int, int)}
     * but with one big difference: the full buffer is always written
     * to the underlying outputstream. Large byte-arrays are chopped
     * into buffer-size pieces and writtten out piece by piece.
     * <br>This provides a closer relation to the write timeout
     * and the maximum (buffer) size of the write-operation to wait on. 
     */
    @Override
    public void write(byte b[], int off, int len) throws IOException {

        if (count >= buf.length) {
            flushBuffer(false);
        }
        if (len <= buf.length - count) {
            System.arraycopy(b, off, buf, count, len);
            count += len;
        } else {
            final int fill = buf.length - count;
            System.arraycopy(b, off, buf, count, fill);
            count += fill;
            flushBuffer(false);
            final int remaining = len - fill;
            int start = off + fill;
            for (int i = 0; i < remaining / buf.length; i++) {
                System.arraycopy(b, start, buf, count, buf.length);
                count = buf.length;
                flushBuffer(false);
                start += buf.length;
            }
            count = remaining % buf.length;
            System.arraycopy(b, start, buf, 0, count);
        }
    }

    /**
     * See {@link BufferedOutputStream#flush()}
     * <br>If a write timeout occurred (i.e. {@link #isWriteTimeout()} returns {@code true}),
     * then this method does nothing. 
     */
    @Override
    public void flush() throws IOException {

        // Protect against flushing before closing after a write-timeout.
        // If that happens, then "out" is already closed in interruptWriteOut.
        if (!isWriteTimeout()) {
            flushBuffer(true);
        }
    }

}

以及测试程序:

import java.io.*;
import java.net.*;
import java.util.concurrent.*;

public class TestTimedSocketOut implements Runnable, Closeable {

    public static void main(String[] args) {

        TestTimedSocketOut m = new TestTimedSocketOut();
        try {
            m.run();
        } finally {
            m.close();
        }
    }

    final int clients = 3; // 2 is minimum, client 1 is expected to fail.
    final int timeOut = 1000;
    final int bufSize = 4096;
    final long maxWait = 5000L;
    // need a large array to write, else the OS just buffers everything and makes it work
    byte[] largeMsg = new byte[28_602];
    final ThreadPoolExecutor tp = (ThreadPoolExecutor) Executors.newCachedThreadPool();
    final ScheduledThreadPoolExecutor stp = (ScheduledThreadPoolExecutor) Executors.newScheduledThreadPool(1);
    final ConcurrentLinkedQueue<Closeable> closeables = new ConcurrentLinkedQueue<Closeable>();
    final CountDownLatch[] serversReady = new CountDownLatch[clients];
    final CountDownLatch clientsDone = new CountDownLatch(clients);
    final CountDownLatch serversDone = new CountDownLatch(clients);

    ServerSocket ss;
    int port;

    @Override public void run()  {

        stp.setRemoveOnCancelPolicy(true);
        try {
            ss = new ServerSocket();
            ss.bind(null);
            port = ss.getLocalPort();
            tp.execute(new SocketAccept());
            for (int i = 0; i < clients; i++) {
                serversReady[i] = new CountDownLatch(1);
                ClientSideSocket css = new ClientSideSocket(i);
                closeables.add(css);
                tp.execute(css);
                // need sleep to ensure client 0 connects first.
                Thread.sleep(50L);
            }
            if (!clientsDone.await(maxWait, TimeUnit.MILLISECONDS)) {
                println("CLIENTS DID NOT FINISH");
            } else {
                if (!serversDone.await(maxWait, TimeUnit.MILLISECONDS)) {
                    println("SERVERS DID NOT FINISH");
                } else {
                    println("Finished");
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    } 

    @Override public void close() {

        try { if (ss != null) ss.close(); } catch (Exception ignored) {}
        Closeable c = null;
        while ((c = closeables.poll()) != null) {
            try { c.close(); } catch (Exception ignored) {}
        }
        tp.shutdownNow();
        println("Scheduled tasks executed: " + stp.getTaskCount() + ", max. threads: " + stp.getLargestPoolSize());
        stp.shutdownNow();
    }

    class SocketAccept implements Runnable {

        @Override public void run() {
            try {
                for (int i = 0; i < clients; i++) {
                    SeverSideSocket sss = new SeverSideSocket(ss.accept(), i);
                    closeables.add(sss);
                    tp.execute(sss);
                }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }

    class SeverSideSocket implements Runnable, Closeable {

        Socket s;
        int number, cnumber;
        boolean completed;

        public SeverSideSocket(Socket s, int number) {
            this.s = s;
            this.number = number;
            cnumber = -1;
        }

        @Override public void run() {

            String t = "nothing";
            try {
                DataInputStream in = new DataInputStream(s.getInputStream());
                DataOutputStream out = new DataOutputStream(s.getOutputStream());
                serversReady[number].countDown();
                Thread.sleep(timeOut);
                t = in.readUTF();
                in.readFully(new byte[largeMsg.length], 0, largeMsg.length);
                t += in.readUTF();
                out.writeByte(1);
                out.flush();
                cnumber = in.readInt();
                completed = true;
            } catch (Exception e) {
                println("server side " + number + " stopped after " + e);
                // e.printStackTrace();
            } finally {
                println("server side " + number + " received: " + t);
                if (completed && cnumber != number) {
                    println("server side " + number + " expected client number " + number + " but got " + cnumber);
                }
                close();
                serversDone.countDown();
            }
        }

        @Override public void close() {
            TestTimedSocketOut.close(s);
            s = null;
        }
    }

    class ClientSideSocket implements Runnable, Closeable {

        Socket s;
        int number;

        public ClientSideSocket(int number) {
            this.number = number;
        }

        @SuppressWarnings("resource")
        @Override public void run() {

            Byte b = -1;
            TimedOutputStream tout = null;
            try {
                s = new Socket();
                s.connect(new InetSocketAddress(port));
                DataInputStream in = new DataInputStream(s.getInputStream());
                tout = new TimedOutputStream(s.getOutputStream(), bufSize, stp);
                if (number == 1) {
                    // expect fail
                    tout.setTimeoutMs(timeOut / 2);
                } else {
                    // expect all OK
                    tout.setTimeoutMs(timeOut * 2);
                }
                DataOutputStream out = new DataOutputStream(tout);
                if (!serversReady[number].await(maxWait, TimeUnit.MILLISECONDS)) {
                    throw new RuntimeException("Server side for client side " + number + " not ready.");
                }
                out.writeUTF("client side " + number + " starting transfer");
                out.write(largeMsg);
                out.writeUTF(" - client side " + number + " completed transfer");
                out.flush();
                b = in.readByte();
                out.writeInt(number);
                out.flush();
            } catch (Exception e) {
                println("client side " + number + " stopped after " + e);
                // e.printStackTrace();
            } finally {
                println("client side " + number + " result: " + b);
                if (tout != null) {
                    if (tout.isWriteTimeout()) {
                        println("client side " + number + " had write timeout, close exception: " + tout.getWriteTimeoutCloseException());
                    } else {
                        println("client side " + number + " had no write timeout");
                    }
                }
                close();
                clientsDone.countDown();
            }
        }

        @Override public void close() {
            TestTimedSocketOut.close(s);
            s = null;
        }
    }

    private static void close(Socket s) {
        try { if (s != null) s.close(); } catch (Exception ignored) {}
    }

    private static final long START_TIME = System.currentTimeMillis(); 

    private static void println(String msg) {
        System.out.println((System.currentTimeMillis() - START_TIME) + "\t " + msg);
    }

}

【讨论】:

  • OMG 每次写入一个线程。可怕。 java.net 中没有任何内容表明流可以异步关闭。并重新实现BufferedOutputStream 而不仅仅是使用一个。还有一个额外的flush 参数来破坏API,当调用者可以在需要时调用flush() 时完全多余。
  • 我忘记了异步关闭,这确实很糟糕。每次写入都有一个未来的任务,但我不希望为每个计划任务(立即)启动一个线程。我确实尝试过直接使用BufferedOutputStream,但我找不到一种方法来只使用超时任务写入底层输出流。额外的flush-paramater只存在于内部(受保护的)方法中,内部方法最好是私有的。
  • @EJP 我能够异步关闭一个阻塞的 .accept() 调用。我认为它也应该适用于 .write()。
猜你喜欢
  • 2016-07-06
  • 1970-01-01
  • 2014-04-12
  • 2011-10-18
  • 2012-06-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-07-20
相关资源
最近更新 更多