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

Reactor Netty TCP服务器仅接收536字节,剩余内容被截断

问题

使用Project Reactor Netty开发TCP服务器,接收客户端字节数组请求并返回响应。当前问题:无论客户端发送多少数据,服务器最多仅接收536字节,剩余内容被截断。客户端消息前4字节标识实际数据长度,该长度总是大于Netty实际接收的数据量。

本地测试时添加SO_RCVBUF/SO_SNDBUF并设置为4096字节后可正常接收,但真实客户端仍出现截断。

服务器代码:

TcpServer tcpServer = TcpServer.create();
Optional.of(someTcpServerConfigObject)
        .filter(config -> config.isLoopResourcesEnabled())
        .ifPresent(enabled -> {
            LoopResources loopResources = LoopResources.create("prefix", 1, 4,true);
            tcpServer.runOn(loopResources);
        });
DisposableServer someTcpServer = tcpServer
        .host("12.123.456.789")
        .wiretap(true)
        .doOnBind(server -> log.info("Starting listener..."))
        .doOnBound(server -> log.info("Listener started on host: {}, port: {}", server.host(), server.port()))
        .port(12345)
        .option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
        .option(ChannelOption.AUTO_CLOSE, false)
        .childOption(ChannelOption.TCP_NODELAY,true)
        .childOption(ChannelOption.AUTO_CLOSE,false)
        .childOption(ChannelOption.SO_KEEPALIVE, true)
        .childOption(ChannelOption.SO_RCVBUF, 4096)
        .childOption(ChannelOption.SO_SNDBUF, 4096)
        .doOnConnection(connection -> {
            InetSocketAddress socketAddress = (InetSocketAddress) connection.channel().remoteAddress();
            log.info("Client has connected. Host: {}, Port: {}",
                        socketAddress.getAddress().getHostAddress(), socketAddress.getPort());
        })
        .doOnChannelInit((observer, channel, remoteAddress) ->
                channel.pipeline()
                        .addFirst(new LoggingHandler(MyTcpServer.class))
                        .addFirst(new TcpServerHandler())
        )
        .handle((inbound, outbound) ->
                inbound
                        .receive()
                        .asByteArray()
                        .flatMap(req -> processRequest(req))
                        //above processRequest() returns a java.nio.ByteBuffer
                        //doing rsp.array() to convert to byte[]
                        .flatMap(rsp -> outbound.sendByteArray(Flux.just(rsp.array()))
                        .doOnError(throwable -> log.error("Error processing the request: {}", throwable.getMessage(),throwable))
        ).bindNow();
someTcpServer.onDispose().block();

自定义TcpServerHandler代码:

@Slf4j
public class TcpServerHandler extends ChannelDuplexHandler {

    private final AtomicLong startTime = new AtomicLong(0L);
    private final AtomicLong endTime = new AtomicLong(0L);

    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) {
        byte[] data = HexFormat.of().formatHex(ByteBufUtil.getBytes((ByteBuf) msg));
        InetSocketAddress socketAddress = (InetSocketAddress) ctx.channel().remoteAddress();
        log.info("Receiving message from: Host: {}, Port: {}. Data: {}", socketAddress.getAddress().getHostAddress(),
                socketAddress.getPort(), data);
        byte[] byteArrayContainingFourByteLength = new byte[4];
        System.arraycopy(data, 0, byteArrayContainingFourByteLength, 0, 4);
        ByteBuffer wrapped = ByteBuffer.wrap(byteArrayContainingFourByteLength);
        short actualLength = wrapped.getShort();
        log.info("# of bytes received from netty: {}", data.length);
        log.info("# of bytes client actually sent: {}", actualLength);
        startTime.set(System.nanoTime());
        ctx.fireChannelRead(msg);
    }

    @Override
    public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
        endTime.set(System.nanoTime());
        log.info("Took {} ms to process", Duration.ofNanos(endTime.get() - startTime.get()).toMillis());
        super.write(ctx, msg, promise);
    }
}

怀疑问题与MSS的536字节限制有关,询问可配置的Netty选项、处理器来解决。

解决方案

1. 核心问题:TCP流式协议的分片特性

TCP是流式传输,inbound.receive()每次读取的是内核缓冲区中的片段(默认对应MSS大小,比如536字节),而非完整的应用层消息。必须通过基于长度的帧解码器将分片拼接成完整消息。

2. 添加Netty的LengthFieldBasedFrameDecoder

这是解决此类问题的标准方案,它会根据消息头中的长度字段自动拼接所有分片,输出完整的ByteBuf。

修改doOnChannelInit中的pipeline配置,将解码器添加到最前面(在LoggingHandler和自定义Handler之前):

.doOnChannelInit((observer, channel, remoteAddress) ->
        channel.pipeline()
                .addFirst(new LengthFieldBasedFrameDecoder(
                        Integer.MAX_VALUE, // 最大帧长度,根据实际业务设置
                        0, // 长度字段的偏移量(前4字节就是长度)
                        4, // 长度字段的字节数
                        0, // 长度调整值(如果长度字段包含自身长度则调整,这里不需要)
                        4 // 跳过的字节数(处理完长度字段后,跳过前4字节,直接取业务数据)
                ))
                .addFirst(new LoggingHandler(MyTcpServer.class))
                .addFirst(new TcpServerHandler())
)

注意:自定义Handler中用getShort()解析4字节长度是错误的,4字节对应int类型,需修正为wrapped.getInt(),否则会导致长度解析错误。

3. 调整缓冲区分配器

可以配置RCVBUF_ALLOCATOR让Netty使用更大的接收缓冲区,避免频繁触发分片读取:

.childOption(ChannelOption.RCVBUF_ALLOCATOR, new FixedRecvByteBufAllocator(8192))

或者使用自适应分配器(默认是自适应,但可以调整参数):

.childOption(ChannelOption.RCVBUF_ALLOCATOR, AdaptiveRecvByteBufAllocator.DEFAULT)

4. 简化handle逻辑

添加解码器后,inbound.receive()每次拿到的都是完整的业务消息(已跳过前4字节长度),无需再处理分片:

.handle((inbound, outbound) ->
        inbound
                .receive()
                .asByteArray()
                .flatMap(req -> processRequest(req))
                .flatMap(rsp -> outbound.sendByteArray(Flux.just(rsp.array())))
                .doOnError(throwable -> log.error("Error processing the request: {}", throwable.getMessage(), throwable))
)

5. 关于MSS的说明

MSS是TCP层的分段限制,应用层无需直接修改。帧解码器会自动处理TCP分多次发送的片段,将其拼接成完整的应用层消息,所以不需要调整MSS相关参数。

6. 修正自定义Handler的长度解析

在自定义Handler中,修正长度解析逻辑,用getInt()代替getShort():

ByteBuffer wrapped = ByteBuffer.wrap(byteArrayContainingFourByteLength);
int actualLength = wrapped.getInt(); // 4字节对应int类型

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 07:17:33