(@NotNull SocketChannel sc)
| 633 | * 考虑清楚以后去掉while(true)? |
| 634 | */ |
| 635 | private void doWrite(@NotNull SocketChannel sc) throws Exception { // 只在selector线程调用 |
| 636 | sendCount++; |
| 637 | int blockSize = selector.getSelectors().getBbPoolBlockSize(); |
| 638 | int bufSize = outputBuffer.size(); |
| 639 | while (true) { |
| 640 | for (Action0 op; /*bufSize < blockSize * 2 &&*/ (op = operates.poll()) != null; ) { |
| 641 | op.run(); |
| 642 | bufSize = outputBuffer.size(); |
| 643 | } |
| 644 | var flushed = true; |
| 645 | var codec = outputCodecChain; |
| 646 | if (codec != null) { |
| 647 | // 减慢flush频率, |
| 648 | // 在保持底层outputBuffer.writeTo能满载的情况下,尽量缓冲住Chain里面的数据。 |
| 649 | // 这使得某些Chain算法比如Zstd能大块的工作,具有更高的效率。 |
| 650 | if (bufSize < blockSize) { |
| 651 | codec.flush(); |
| 652 | int newBufSize = outputBuffer.size(); |
| 653 | int deltaLen = newBufSize - bufSize; |
| 654 | if (deltaLen != 0) { |
| 655 | bufSize = newBufSize; |
| 656 | outputBufferSizeHandle.getAndAdd(this, (long)deltaLen); |
| 657 | } |
| 658 | } else |
| 659 | flushed = false; |
| 660 | } |
| 661 | |
| 662 | if (bufSize > 0) { |
| 663 | var rc = outputBuffer.writeTo(sc); |
| 664 | if (rc < 0) { |
| 665 | close(); // 很罕见的正常关闭, 不设置异常, 其实write抛异常的可能性更大 |
| 666 | return; |
| 667 | } |
| 668 | sendSize += rc; |
| 669 | outputBufferSizeHandle.getAndAdd(this, -rc); |
| 670 | bufSize = outputBuffer.size(); |
| 671 | if (bufSize > 0) { |
| 672 | // 有数据正在发送,此时可以安全退出执行,写完以后Selector会再次触发doWrite。 |
| 673 | // add write event,里面判断了事件没有变化时不做操作,严格来说,再次注册事件是不需要的。 |
| 674 | return; |
| 675 | } |
| 676 | // 全部都写出去了,继续尝试看看有没有新的操作。 |
| 677 | } |
| 678 | // 时间窗口 |
| 679 | // 必须和把Operate加入队列同步!否则可能会出现,刚加入操作没有被处理,但是OP_WRITE又被Remove的问题。 |
| 680 | if (operates.isEmpty() && flushed) { // 此时bufSize=0,下次循环会触发flush |
| 681 | // 真的没有等待处理的操作了,去掉事件,返回。以后新的操作在下一次doWrite时处理。 |
| 682 | removeInterestOps(SelectionKey.OP_WRITE); |
| 683 | if (operates.isEmpty()) { // 再判断一次,避免跟submitAction的并发竞争问题 |
| 684 | if (closePending) |
| 685 | realClose(); |
| 686 | return; |
| 687 | } |
| 688 | addInterestOps(SelectionKey.OP_WRITE); |
| 689 | } |
| 690 | // 发现数据,继续尝试处理。 |
| 691 | } |
| 692 | } |
no test coverage detected