You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于CompletableFuture的Java TCP客户端异常行为排查求助

问题:基于CompletableFuture的TCP回显客户端读取操作未等待服务器返回就执行

我尝试实现一个用CompletableFuture链式调用的TCP回显客户端,但遇到异常:socket读取操作在写入完成后立即执行,没有等待服务器返回数据。

客户端代码

package tcp;

import utilities.CompletableFutureFactory;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.AsynchronousSocketChannel;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;

public class ClientTxRx {
    final private AsynchronousSocketChannel m_sockChannel;
    final private InetSocketAddress m_socketAddress;
    final private int m_dataBlockSizeInBytes;

    public ClientTxRx( String host, int port, int dataBlockSizeInBytes ) throws IOException {
        //create a socket channel
        m_sockChannel = AsynchronousSocketChannel.open();
        m_socketAddress = new InetSocketAddress(host, port);
        m_dataBlockSizeInBytes = dataBlockSizeInBytes;
    }

    public void sendRequestReceiveResponse() throws InterruptedException, ExecutionException {
        ByteBuffer buf = ByteBuffer.allocate(m_dataBlockSizeInBytes);

        CompletableFutureFactory.create(m_sockChannel.connect(m_socketAddress))
        .thenCompose(s -> {
            System.out.println("client: start writing");
            buf.put(0, (byte)1);
            buf.put(1, (byte)2);

            return CompletableFutureFactory.create(m_sockChannel.write(buf));
        }).thenCompose(s -> {
            System.out.println("client: finished writing, start reading");
            return CompletableFutureFactory.create(m_sockChannel.read(buf));
        }).thenRun(() -> {
            final int test0 = buf.get(0);
            final int test1 = buf.get(1);

            System.out.println(String.format("client Data: %d %d", test0, test1));

            System.out.println("client: finished!");
        }).orTimeout(1, TimeUnit.SECONDS).get();
    }
}

服务器代码

package tcp;

import utilities.CompletableFutureFactory;

import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.AsynchronousServerSocketChannel;
import java.util.concurrent.TimeUnit;

public class ServerRxTx {
    final private AsynchronousServerSocketChannel m_serverSock;
    final private int m_dataBlockSizeInBytes;

    public ServerRxTx(String bindAddr,
                      int bindPort,
                      int dataBlockSizeInBytes) throws IOException {
        InetSocketAddress sockAddr = new InetSocketAddress(bindAddr, bindPort);

        m_serverSock = AsynchronousServerSocketChannel.open().bind(sockAddr);
        m_dataBlockSizeInBytes = dataBlockSizeInBytes;
    }

    public void startRxTx() {
        final ByteBuffer buf = ByteBuffer.allocate(m_dataBlockSizeInBytes);

        CompletableFutureFactory.create(m_serverSock.accept())
                .thenAccept(s -> {
                    System.out.println("server: accepted");
                    try {
                        CompletableFutureFactory
                                .create(s.read(buf))
                                .thenCompose(readBytes -> {
                                    System.out.println("server: read passed, ready to write");
                                    final int test0 = buf.get(0);
                                    final int test1 = buf.get(1);

                                    System.out.println(String.format("server Data: %d %d", test0, test1));

                                    if (readBytes == -1) {
                                        throw new RuntimeException();
                                    }
                                    try {
                                        Thread.sleep(2_000);
                                    }catch (Exception e) {
                                        e.printStackTrace();
                                    }

                                    buf.put(0, (byte)10);
                                    buf.put(1, (byte)20);

                                    return CompletableFutureFactory.create(
                                            s.write(buf)
                                    );

//                                    return CompletableFutureFactory.create(
//                                            s.write(executer.execute(buf))
//                                    );
                                }).orTimeout(1, TimeUnit.SECONDS).get();
                    } catch (Exception e) {
                        throw new RuntimeException(e.getMessage());
                    }
                }).thenRun(() -> {
                    System.out.println("server: finished");
                }).orTimeout(2, TimeUnit.SECONDS);
    }

    public void stopRxTx() {
        try {
            m_serverSock.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

Future转CompletableFuture辅助代码

package utilities;

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.Future;

public class CompletableFutureFactory {
    public static <T> CompletableFuture<T> create(Future<T> future) {
        if (future.isDone())
            return transformDoneFuture(future);
        return CompletableFuture.supplyAsync(() -> {
            try {
                if (!future.isDone())
                    awaitFutureIsDoneInForkJoinPool(future);
                return future.get();
            } catch (ExecutionException e) {
                throw new RuntimeException(e);
            } catch (InterruptedException e) {
                // Normally, this should never happen inside ForkJoinPool
                Thread.currentThread().interrupt();
                // Add the following statement if the future doesn't have side effects
                // future.cancel(true);
                throw new RuntimeException(e);
            }
        });
    }

    private static <T> CompletableFuture<T> transformDoneFuture(Future<T> future) {
        CompletableFuture<T> cf = new CompletableFuture<>();
        T result;
        try {
            result = future.get();
        } catch (Throwable ex) {
            cf.completeExceptionally(ex);
            return cf;
        }
        cf.complete(result);
        return cf;
    }

    private static void awaitFutureIsDoneInForkJoinPool(Future<?> future)
            throws InterruptedException {
        ForkJoinPool.managedBlock(new ForkJoinPool.ManagedBlocker() {
            @Override public boolean block() throws InterruptedException {
                try {
                    future.get();
                } catch (ExecutionException e) {
                    throw new RuntimeException(e);
                }
                return true;
            }
            @Override public boolean isReleasable() {
                return future.isDone();
            }
        });
    }
}

测试场景代码

import tcp.ClientTxRx;
import tcp.ServerRxTx;

public class Main {
    public static void main(String[] args) {
        try {
            final var server = new ServerRxTx("127.0.0.1", 10000, 256);
            final var client = new ClientTxRx("127.0.0.1", 10000, 256);

            server.startRxTx();
            client.sendRequestReceiveResponse();
            server.stopRxTx();
        }catch (Exception e) {
            e.printStackTrace();
        }

        try {
            Thread.sleep(5_000);
        }catch (Exception e) {
            e.printStackTrace();
        }
    }
}

问题细节

客户端的读取操作在写入完成后立刻触发,没有实际等待服务器返回响应数据,最终读取到的还是本地写入的原始值,而非服务器返回的修改后数据。


问题根源与修复方案

1. ByteBuffer读写模式未切换(核心问题)

NIO的ByteBuffer有写模式和读模式两种状态:

  • 写入数据后,position会指向写入后的下一个位置,此时直接调用read/write会从当前position开始操作,而非从头开始;
  • 服务器和客户端都没有在读写操作之间切换模式,导致数据写入/读取的位置错误,甚至没有实际传输有效数据。

客户端修复

写入完成前切换到读模式(让socket从缓冲区开头读取数据发送),读取前清空缓冲区切换回写模式:

.thenCompose(s -> {
    System.out.println("client: start writing");
    buf.put(0, (byte)1);
    buf.put(1, (byte)2);
    buf.flip(); // 切换到读模式,准备向socket写入数据
    return CompletableFutureFactory.create(m_sockChannel.write(buf));
}).thenCompose(s -> {
    System.out.println("client: finished writing, start reading");
    buf.clear(); // 清空缓冲区,切换回写模式准备接收服务器数据
    return CompletableFutureFactory.create(m_sockChannel.read(buf));
}).thenRun(() -> {
    buf.flip(); // 切换到读模式,读取服务器返回的有效数据
    final int test0 = buf.get(0);
    final int test1 = buf.get(1);
    System.out.println(String.format("client Data: %d %d", test0, test1));
    System.out.println("client: finished!");
})

服务器端修复

读取客户端数据后切换模式读取内容,写入响应前切换写模式,写入前再切换读模式:

.thenCompose(readBytes -> {
    System.out.println("server: read passed, ready to write");
    buf.flip(); // 切换到读模式,读取客户端发送的数据
    final int test0 = buf.get(0);
    final int test1 = buf.get(1);
    System.out.println(String.format("server Data: %d %d", test0, test1));

    if (readBytes == -1) {
        throw new RuntimeException();
    }
    try {
        Thread.sleep(2_000);
    }catch (Exception e) {
        e.printStackTrace();
    }

    buf.clear(); // 清空缓冲区,切换到写模式
    buf.put(0, (byte)10);
    buf.put(1, (byte)20);
    buf.flip(); // 切换到读模式,准备向客户端写入响应
    return CompletableFutureFactory.create(s.write(buf));
})

2. 服务器端阻塞调用问题

服务器的thenAccept中调用get()会阻塞线程,破坏CompletableFuture的异步特性,建议改用链式调用替代:

.thenAccept(s -> {
    System.out.println("server: accepted");
    CompletableFutureFactory
            .create(s.read(buf))
            .thenCompose(readBytes -> {
                // 原有逻辑...
                return CompletableFutureFactory.create(s.write(buf));
            }).orTimeout(1, TimeUnit.SECONDS)
            .exceptionally(e -> {
                throw new RuntimeException(e.getMessage());
            });
})

额外注意事项

  • AsynchronousSocketChannel的write/read操作可能不会一次性处理完所有数据,生产环境需要循环操作直到全部数据写入/读取完成,避免数据截断;
  • 客户端的orTimeout设置为1秒,而服务器有2秒的休眠,会导致客户端超时,需要调整超时时间匹配服务器逻辑。

内容的提问来源于stack exchange,提问作者Dmitry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 00:05:53