This is an automated email from the ASF dual-hosted git repository. jiangtian pushed a commit to branch cluster_add_snappy in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 938ec8b0d08bc7cec661f505c7b65b178ed1ba04 Author: jt <[email protected]> AuthorDate: Tue Dec 1 19:50:12 2020 +0800 fix thrift dependency --- .../iotdb/rpc/AutoScalingBufferReadTransport.java | 61 ++++++++++++++++++++-- .../iotdb/rpc/TSnappyElasticFramedTransport.java | 1 - 2 files changed, 58 insertions(+), 4 deletions(-) diff --git a/service-rpc/src/main/java/org/apache/iotdb/rpc/AutoScalingBufferReadTransport.java b/service-rpc/src/main/java/org/apache/iotdb/rpc/AutoScalingBufferReadTransport.java index 085e9e8..d94dbf8 100644 --- a/service-rpc/src/main/java/org/apache/iotdb/rpc/AutoScalingBufferReadTransport.java +++ b/service-rpc/src/main/java/org/apache/iotdb/rpc/AutoScalingBufferReadTransport.java @@ -19,17 +19,72 @@ package org.apache.iotdb.rpc; -import org.apache.thrift.transport.AutoExpandingBufferReadTransport; +import org.apache.thrift.transport.TTransport; +import org.apache.thrift.transport.TTransportException; -public class AutoScalingBufferReadTransport extends AutoExpandingBufferReadTransport { +public class AutoScalingBufferReadTransport extends TTransport { private final AutoExpandingBuffer buf; + private int pos = 0; + private int limit = 0; public AutoScalingBufferReadTransport(int initialCapacity) { - super(initialCapacity); this.buf = new AutoExpandingBuffer(initialCapacity); } + public void fill(TTransport inTrans, int length) throws TTransportException { + buf.resizeIfNecessary(length); + inTrans.readAll(buf.array(), 0, length); + pos = 0; + limit = length; + } + + @Override + public void close() { + // do nothing + } + + @Override + public boolean isOpen() { return true; } + + @Override + public void open() { + // do nothing + } + + @Override + public final int read(byte[] target, int off, int len) { + int amtToRead = Math.min(len, getBytesRemainingInBuffer()); + System.arraycopy(buf.array(), pos, target, off, amtToRead); + consumeBuffer(amtToRead); + return amtToRead; + } + + @Override + public void write(byte[] buf, int off, int len) { + throw new UnsupportedOperationException(); + } + + @Override + public final void consumeBuffer(int len) { + pos += len; + } + + @Override + public final byte[] getBuffer() { + return buf.array(); + } + + @Override + public final int getBufferPosition() { + return pos; + } + + @Override + public final int getBytesRemainingInBuffer() { + return limit - pos; + } + public void resizeIfNecessary(int size) { buf.resizeIfNecessary(size); } diff --git a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSnappyElasticFramedTransport.java b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSnappyElasticFramedTransport.java index 307b363..426b759 100644 --- a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSnappyElasticFramedTransport.java +++ b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSnappyElasticFramedTransport.java @@ -132,7 +132,6 @@ public class TSnappyElasticFramedTransport extends TFastFramedTransport { @Override public void flush() throws TTransportException { int length = writeBuffer.getPos(); - TFramedTransport.encodeFrameSize(length, i32buf); try { int maxCompressedLength = Snappy.maxCompressedLength(length); writeCompressBuffer = resizeCompressBuf(maxCompressedLength, writeCompressBuffer);
