
x33g5p2x  于2022-01-17 转载在 其他  





代码示例来源:origin: apache/flink

private ByteBuf readChunk() {
  if (isClosed) {
    return null;
  } else if (buf.readableBytes() <= chunkSize) {
    isEndOfInput = true;
    // Don't retain as the consumer is responsible to release it
    return buf.slice();
  } else {
    // Return a chunk sized slice of the buffer. The ref count is
    // shared with the original buffer. That's why we need to retain
    // a reference here.
    return buf.readSlice(chunkSize).retain();

代码示例来源:origin: org.apache.flink/flink-queryable-state-client-java_2.11

private ByteBuf readChunk() {
  if (isClosed) {
    return null;
  } else if (buf.readableBytes() <= chunkSize) {
    isEndOfInput = true;
    // Don't retain as the consumer is responsible to release it
    return buf.slice();
  } else {
    // Return a chunk sized slice of the buffer. The ref count is
    // shared with the original buffer. That's why we need to retain
    // a reference here.
    return buf.readSlice(chunkSize).retain();


public ByteBuf readChunk(ChannelHandlerContext ctx) throws Exception {
  if (isClosed) {
    return null;
  } else if (buf.readableBytes() <= chunkSize) {
    isEndOfInput = true;
    // Don't retain as the consumer is responsible to release it
    return buf.slice();
  } else {
    // Return a chunk sized slice of the buffer. The ref count is
    // shared with the original buffer. That's why we need to retain
    // a reference here.
    return buf.readSlice(chunkSize).retain();

代码示例来源:origin: org.apache.flink/flink-runtime

static BufferResponse readFrom(ByteBuf buffer) {
    InputChannelID receiverId = InputChannelID.fromByteBuf(buffer);
    int sequenceNumber = buffer.readInt();
    int backlog = buffer.readInt();
    boolean isBuffer = buffer.readBoolean();
    int size = buffer.readInt();
    ByteBuf retainedSlice = buffer.readSlice(size).retain();
    return new BufferResponse(retainedSlice, isBuffer, sequenceNumber, receiverId, backlog);

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

static BufferResponse readFrom(ByteBuf buffer) {
    InputChannelID receiverId = InputChannelID.fromByteBuf(buffer);
    int sequenceNumber = buffer.readInt();
    int backlog = buffer.readInt();
    boolean isBuffer = buffer.readBoolean();
    int size = buffer.readInt();
    ByteBuf retainedSlice = buffer.readSlice(size).retain();
    return new BufferResponse(retainedSlice, isBuffer, sequenceNumber, receiverId, backlog);


   * Parses the whole BufferResponse message and composes a new BufferResponse with both header parsed and
   * data buffer filled in. This method is used in non-credit-based network stack.
   * @param buffer the whole serialized BufferResponse message.
   * @return a BufferResponse object with the header parsed and the data buffer filled in.
  static BufferResponse readFrom(ByteBuf buffer) {
    InputChannelID receiverId = InputChannelID.fromByteBuf(buffer);
    int sequenceNumber = buffer.readInt();
    int backlog = buffer.readInt();
    boolean isBuffer = buffer.readBoolean();
    int size = buffer.readInt();
    ByteBuf retainedSlice = buffer.readSlice(size).retain();
    return new BufferResponse(retainedSlice, isBuffer, sequenceNumber, receiverId, backlog, size);
