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);

Reply via email to