diff --git a/common/src/main/java/org/apache/uniffle/common/netty/TransportFrameDecoder.java b/common/src/main/java/org/apache/uniffle/common/netty/TransportFrameDecoder.java index a0aef7d1a0..9d8ecd7f9d 100644 --- a/common/src/main/java/org/apache/uniffle/common/netty/TransportFrameDecoder.java +++ b/common/src/main/java/org/apache/uniffle/common/netty/TransportFrameDecoder.java @@ -27,6 +27,8 @@ import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.apache.uniffle.common.exception.RssException; import org.apache.uniffle.common.netty.protocol.Message; @@ -47,6 +49,7 @@ * method. */ public class TransportFrameDecoder extends ChannelInboundHandlerAdapter implements FrameDecoder { + private static final Logger LOG = LoggerFactory.getLogger(TransportFrameDecoder.class); private int msgSize = -1; private int bodySize = -1; private Message.Type curType = Message.Type.UNKNOWN_TYPE; @@ -199,6 +202,12 @@ public void handlerRemoved(ChannelHandlerContext ctx) throws Exception { // - When the Channel becomes inactive // - When the decoder is removed from the ChannelPipeline for (ByteBuf b : buffers) { + LOG.warn("Check NettyManagedBuffer release"); + LOG.warn(Thread.currentThread().getName()); + LOG.warn(this.toString()); + LOG.warn(b.toString()); + LOG.warn("size: " + b.readableBytes()); + LOG.warn("Check decoder release stack tree", new Throwable()); b.release(); } buffers.clear(); diff --git a/common/src/main/java/org/apache/uniffle/common/netty/buffer/NettyManagedBuffer.java b/common/src/main/java/org/apache/uniffle/common/netty/buffer/NettyManagedBuffer.java index 547a4e891a..a3c346a115 100644 --- a/common/src/main/java/org/apache/uniffle/common/netty/buffer/NettyManagedBuffer.java +++ b/common/src/main/java/org/apache/uniffle/common/netty/buffer/NettyManagedBuffer.java @@ -21,8 +21,11 @@ import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; public class NettyManagedBuffer extends ManagedBuffer { + private static final Logger LOG = LoggerFactory.getLogger(NettyManagedBuffer.class); public static final NettyManagedBuffer EMPTY_BUFFER = new NettyManagedBuffer(Unpooled.EMPTY_BUFFER); @@ -56,6 +59,12 @@ public ManagedBuffer retain() { @Override public ManagedBuffer release() { + LOG.warn("Check NettyManagedBuffer release"); + LOG.warn(Thread.currentThread().getName()); + LOG.warn(this.toString()); + LOG.warn(this.buf.toString()); + LOG.warn("size: " + this.buf.readableBytes()); + LOG.warn("Check NettyManagedBuffer release stack tree", new Throwable()); buf.release(); return this; }