基于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
相关产品推荐
相关产品推荐

