【问题标题】:Tcp server WritePendingException although thread locksTcp 服务器 WritePendingException 虽然线程锁
【发布时间】:2016-07-30 22:40:07
【问题描述】:

我用 nio 编写了一个简单的异步 tcp 服务器。 服务器应该能够为每个客户端同时读取和写入。 这是我用一个简单的数据包队列实现的。

public class TcpJobHandler {

private BlockingQueue<TcpJob> _packetQueue = new LinkedBlockingQueue<TcpJob>();
private Thread _jobThread;
private final ReentrantLock _lock = new ReentrantLock();

public TcpJobHandler(){
    _jobThread = new Thread(new Runnable() {
        @Override
        public void run() {
            jobLoop();
        }       
    });

    _jobThread.start();
}

private void jobLoop(){
    while(true){
        try {
            _lock.lock();
            TcpJob job = _packetQueue.take();
            if(job == null){
                continue;
            }
            job.execute();
        } catch (Exception e) {
            AppLogger.error("Failed to dequeue packet from job queue.", e);
        }finally{
            _lock.unlock();
        }
    }
}

public void insertJob(TcpJob job){
    try{
        _packetQueue.put(job);
    }catch(InterruptedException e){
        AppLogger.error("Failed to queue packet to the tcp job queue.", e);
    }
}
}

这段代码的作用只是检查一个新的数据包。如果有新的数据包可用,该数据包将被发送到客户端。 在类 tcp 作业中,只有要发送的数据包和一个将数据包写入客户端流的写入类。 如您所见,只有一个线程应该能够将数据包写入客户端流。

这就是重点,为什么我不明白,为什么我会收到这个错误?如果我是对的,这个异常表示,我尝试将数据发送到流中,但是已经有一个线程正在将数据写入该流中。但为什么?

//编辑: 我遇到了这个异常:

19:18:41.468 [ERROR] - [mufisync.server.data.tcp.handler.TcpJobHandler] : Failed to dequeue packet from job queue. Exception: java.nio.channels.WritePendingException
at sun.nio.ch.AsynchronousSocketChannelImpl.write(Unknown Source)
at sun.nio.ch.AsynchronousSocketChannelImpl.write(Unknown Source)
at mufisync.server.data.tcp.stream.OutputStreamAdapter.write(OutputStreamAdapter.java:35)
at mufisync.server.data.tcp.stream.OutputStreamAdapter.write(OutputStreamAdapter.java:26)
at mufisync.server.data.tcp.stream.BinaryWriter.write(BinaryWriter.java:21)
at mufisync.server.data.tcp.TcpJob.execute(TcpJob.java:29)
at mufisync.server.data.tcp.handler.TcpJobHandler.jobLoop(TcpJobHandler.java:40)
at mufisync.server.data.tcp.handler.TcpJobHandler.access$0(TcpJobHandler.java:32)
at mufisync.server.data.tcp.handler.TcpJobHandler$1.run(TcpJobHandler.java:25)
at java.lang.Thread.run(Unknown Source)

TcpJob 如下所示:

public class TcpJob {

private BasePacket _packet;
private BinaryWriter _writer;

public TcpJob(BasePacket packet, BinaryWriter writer){
    _packet = packet;
    _writer = writer;
}

public void execute(){
    try {
        if(_packet == null){
            AppLogger.warn("Tcp job packet is null");
            return;
        }

        _writer.write(_packet.toByteArray());
    } catch (IOException e) {
        AppLogger.error("Failed to write packet into the stream.", e);
    }
}

public BasePacket get_packet() {
    return _packet;
}   
}

BinaryStream 只是耦合到 AsynchronousSocketChannel,它从套接字通道调用 write(byte[]) 方法。

【问题讨论】:

  • 您遇到了什么异常?你能把它包括在问题中吗?我们看不到您的计算机。
  • "只有一个线程应该能够将数据包写入客户端流。"在您提供的代码中没有任何地方可以看到这一点。你能包括你正在谈论的代码吗?
  • _lock 打算做什么?其他线程如何使用它?
  • 好的,我已经编辑了我的问题。抱歉没有信息。

标签: java multithreading sockets networking tcp


【解决方案1】:

您正在使用异步 NIO2。当您使用异步 IO 时,您不能在最后一次写入完成之前调用 write()。来自 Javadoc

 * @throws  WritePendingException
 *          If a write operation is already in progress on this channel

例如如果你用过

public abstract Future<Integer> write(ByteBuffer src);

在 Future.get() 返回之前,您不能再次写入。

如果你使用

public abstract <A> void write(ByteBuffer src,
                               long timeout,
                               TimeUnit unit,
                               A attachment,
                               CompletionHandler<Integer,? super A> handler);

在调用 CompletionHandler 之前,您不能再次写入。

注意:您也不能同时执行两次读取。

在你的情况下,你想要类似的东西

ByteBuffer lastBuffer = null;
Future<Integer> future = null;

public void execute(){
    try {
        if(_packet == null){
            AppLogger.warn("Tcp job packet is null");
            return;
        }
        // need to wait until the last buffer was written completely.
        while (future != null) {
           future.get();
           if (lastBuffer.remaining() > 0)
              future = _writer.write(lasBuffer);
           else
              break;
        }
        // start another write.
        future = _writer.write(lastBuffer = _packet.toByteArray());
    } catch (IOException e) {
        AppLogger.error("Failed to write packet into the stream.", e);
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-01-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-03-04
    • 1970-01-01
    • 1970-01-01
    • 2022-01-25
    相关资源
    最近更新 更多