This is an automated email from the ASF dual-hosted git repository.

jackietien pushed a commit to branch ty/queryopt
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 7ec19a11ec9af2e88a398534e978973f4752bcff
Author: JackieTien97 <[email protected]>
AuthorDate: Tue Jun 13 17:03:48 2023 +0800

    Opt time batch load
---
 .../iotdb/tsfile/encoding/decoder/Decoder.java     |   8 +
 .../encoding/decoder/DeltaBinaryDecoder.java       | 258 +++++++++++++--------
 .../tsfile/read/reader/page/TimePageReader.java    |  10 +-
 3 files changed, 176 insertions(+), 100 deletions(-)

diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/Decoder.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/Decoder.java
index 4c4eda9dfcc..3d96d820aef 100644
--- a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/Decoder.java
+++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/Decoder.java
@@ -172,6 +172,10 @@ public abstract class Decoder {
     throw new TsFileDecodingException("Method readInt is not supported by 
Decoder");
   }
 
+  public void readInt(ByteBuffer buffer, int[] data, int length) {
+    throw new TsFileDecodingException("Method readInt is not supported by 
Decoder");
+  }
+
   public boolean readBoolean(ByteBuffer buffer) {
     throw new TsFileDecodingException("Method readBoolean is not supported by 
Decoder");
   }
@@ -184,6 +188,10 @@ public abstract class Decoder {
     throw new TsFileDecodingException("Method readLong is not supported by 
Decoder");
   }
 
+  public void readLong(ByteBuffer buffer, long[] data, int length) {
+    throw new TsFileDecodingException("Method readInt is not supported by 
Decoder");
+  }
+
   public float readFloat(ByteBuffer buffer) {
     throw new TsFileDecodingException("Method readFloat is not supported by 
Decoder");
   }
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/DeltaBinaryDecoder.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/DeltaBinaryDecoder.java
index 7f50cd4bb6a..61e1245153c 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/DeltaBinaryDecoder.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/encoding/decoder/DeltaBinaryDecoder.java
@@ -21,7 +21,6 @@ package org.apache.iotdb.tsfile.encoding.decoder;
 
 import org.apache.iotdb.tsfile.encoding.encoder.DeltaBinaryEncoder;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
-import org.apache.iotdb.tsfile.utils.BytesUtils;
 import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
 
 import java.io.IOException;
@@ -36,11 +35,10 @@ import java.nio.ByteBuffer;
  */
 public abstract class DeltaBinaryDecoder extends Decoder {
 
-  protected long count = 0;
-  protected byte[] deltaBuf;
+  protected static final int[] MASK =
+      new int[] {0xFF, 0xFF >> 1, 0xFF >> 2, 0xFF >> 3, 0xFF >> 4, 0xFF >> 5, 
0xFF >> 6, 0xFF >> 7};
 
-  /** the first value in one pack. */
-  protected int readIntTotalCount = 0;
+  protected byte[] deltaBuf;
 
   protected int nextReadIndex = 0;
   /** max bit length of all value in a pack. */
@@ -51,16 +49,10 @@ public abstract class DeltaBinaryDecoder extends Decoder {
   /** how many bytes data takes after encoding. */
   protected int encodingLength;
 
-  public DeltaBinaryDecoder() {
+  protected DeltaBinaryDecoder() {
     super(TSEncoding.TS_2DIFF);
   }
 
-  protected abstract void readHeader(ByteBuffer buffer) throws IOException;
-
-  protected abstract void allocateDataArray();
-
-  protected abstract void readValue(int i);
-
   /**
    * calculate the bytes length containing v bits.
    *
@@ -68,12 +60,37 @@ public abstract class DeltaBinaryDecoder extends Decoder {
    * @return number of bytes
    */
   protected int ceil(int v) {
-    return (int) Math.ceil((double) (v) / 8.0);
+    return (int) Math.ceil(v / 8.0);
+  }
+
+  /**
+   * if remaining data has been run out, load next pack from InputStream.
+   *
+   * @param buffer ByteBuffer
+   */
+  protected void loadBatch(ByteBuffer buffer) {
+    packNum = ReadWriteIOUtils.readInt(buffer);
+    packWidth = ReadWriteIOUtils.readInt(buffer);
+    readHeader(buffer);
+
+    encodingLength = ceil(packNum * packWidth);
+    deltaBuf = new byte[encodingLength];
+    buffer.get(deltaBuf);
+    allocateDataArray();
+
+    nextReadIndex = 0;
+    readPack();
   }
 
+  protected abstract void readHeader(ByteBuffer buffer);
+
+  protected abstract void allocateDataArray();
+
+  protected abstract void readPack();
+
   @Override
   public boolean hasNext(ByteBuffer buffer) throws IOException {
-    return (nextReadIndex < readIntTotalCount) || buffer.remaining() > 0;
+    return (nextReadIndex < packNum) || buffer.remaining() > 0;
   }
 
   public static class IntDeltaDecoder extends DeltaBinaryDecoder {
@@ -94,46 +111,51 @@ public abstract class DeltaBinaryDecoder extends Decoder {
      * @param buffer ByteBuffer
      * @return int
      */
-    protected int readT(ByteBuffer buffer) {
-      if (nextReadIndex == readIntTotalCount) {
-        return loadIntBatch(buffer);
+    @Override
+    public int readInt(ByteBuffer buffer) {
+      if (nextReadIndex == packNum) {
+        loadBatch(buffer);
+        return firstValue;
       }
       return data[nextReadIndex++];
     }
 
     @Override
-    public int readInt(ByteBuffer buffer) {
-      return readT(buffer);
-    }
-
-    /**
-     * if remaining data has been run out, load next pack from InputStream.
-     *
-     * @param buffer ByteBuffer
-     * @return int
-     */
-    protected int loadIntBatch(ByteBuffer buffer) {
-      packNum = ReadWriteIOUtils.readInt(buffer);
-      packWidth = ReadWriteIOUtils.readInt(buffer);
-      count++;
-      readHeader(buffer);
-
-      encodingLength = ceil(packNum * packWidth);
-      deltaBuf = new byte[encodingLength];
-      buffer.get(deltaBuf);
-      allocateDataArray();
-
-      previous = firstValue;
-      readIntTotalCount = packNum;
-      nextReadIndex = 0;
-      readPack();
-      return firstValue;
-    }
-
-    private void readPack() {
-      for (int i = 0; i < packNum; i++) {
-        readValue(i);
-        previous = data[i];
+    public void readInt(ByteBuffer buffer, int[] data, int length) {
+      int offset = 0;
+      while (offset < length) {
+        packNum = ReadWriteIOUtils.readInt(buffer);
+        packWidth = ReadWriteIOUtils.readInt(buffer);
+        readHeader(buffer);
+        data[offset++] = firstValue;
+        encodingLength = ceil(packNum * packWidth);
+        deltaBuf = new byte[encodingLength];
+        buffer.get(deltaBuf);
+
+        int value;
+        int index = 0;
+        int currentByteOffset = 0;
+        int width;
+
+        for (int i = 0; i < packNum; i++) {
+          value = 0;
+          width = packWidth;
+          while (width > 0) {
+            int m = width + currentByteOffset >= 8 ? 8 - currentByteOffset : 
width;
+            width -= m;
+            value = value << m;
+            int y = deltaBuf[index] & MASK[currentByteOffset];
+            currentByteOffset += m;
+            y >>>= (8 - currentByteOffset);
+            value |= y;
+            if (currentByteOffset == 8) {
+              currentByteOffset = 0;
+              ++index;
+            }
+          }
+          data[offset] = data[offset - 1] + minDeltaBase + value;
+          ++offset;
+        }
       }
     }
 
@@ -141,6 +163,7 @@ public abstract class DeltaBinaryDecoder extends Decoder {
     protected void readHeader(ByteBuffer buffer) {
       minDeltaBase = ReadWriteIOUtils.readInt(buffer);
       firstValue = ReadWriteIOUtils.readInt(buffer);
+      previous = firstValue;
     }
 
     @Override
@@ -149,9 +172,31 @@ public abstract class DeltaBinaryDecoder extends Decoder {
     }
 
     @Override
-    protected void readValue(int i) {
-      int v = BytesUtils.bytesToInt(deltaBuf, packWidth * i, packWidth);
-      data[i] = previous + minDeltaBase + v;
+    protected void readPack() {
+      int value;
+      int index = 0;
+      int currentByteOffset = 0;
+      int width;
+
+      for (int i = 0; i < packNum; i++) {
+        value = 0;
+        width = packWidth;
+        while (width > 0) {
+          int m = width + currentByteOffset >= 8 ? 8 - currentByteOffset : 
width;
+          width -= m;
+          value = value << m;
+          int y = deltaBuf[index] & MASK[currentByteOffset];
+          currentByteOffset += m;
+          y >>>= (8 - currentByteOffset);
+          value |= y;
+          if (currentByteOffset == 8) {
+            currentByteOffset = 0;
+            index++;
+          }
+        }
+        data[i] = previous + minDeltaBase + value;
+        previous = data[i];
+      }
     }
 
     @Override
@@ -178,65 +223,90 @@ public abstract class DeltaBinaryDecoder extends Decoder {
      * @param buffer ByteBuffer
      * @return long value
      */
-    protected long readT(ByteBuffer buffer) {
-      if (nextReadIndex == readIntTotalCount) {
-        return loadIntBatch(buffer);
+    @Override
+    public long readLong(ByteBuffer buffer) {
+      if (nextReadIndex == packNum) {
+        loadBatch(buffer);
+        return firstValue;
       }
       return data[nextReadIndex++];
     }
 
-    /**
-     * if remaining data has been run out, load next pack from InputStream.
-     *
-     * @param buffer ByteBuffer
-     * @return long value
-     */
-    protected long loadIntBatch(ByteBuffer buffer) {
-      packNum = ReadWriteIOUtils.readInt(buffer);
-      packWidth = ReadWriteIOUtils.readInt(buffer);
-      count++;
-      readHeader(buffer);
-
-      encodingLength = ceil(packNum * packWidth);
-      deltaBuf = new byte[encodingLength];
-      buffer.get(deltaBuf);
-      allocateDataArray();
-
-      previous = firstValue;
-      readIntTotalCount = packNum;
-      nextReadIndex = 0;
-      readPack();
-      return firstValue;
-    }
-
-    private void readPack() {
-      for (int i = 0; i < packNum; i++) {
-        readValue(i);
-        previous = data[i];
-      }
-    }
-
     @Override
-    public long readLong(ByteBuffer buffer) {
-
-      return readT(buffer);
+    public void readLong(ByteBuffer buffer, long[] data, int length) {
+      int offset = 0;
+      while (offset < length) {
+        packNum = ReadWriteIOUtils.readInt(buffer);
+        packWidth = ReadWriteIOUtils.readInt(buffer);
+        readHeader(buffer);
+        data[offset++] = firstValue;
+        encodingLength = ceil(packNum * packWidth);
+        deltaBuf = new byte[encodingLength];
+        buffer.get(deltaBuf);
+
+        long value;
+        int index = 0;
+        int currentByteOffset = 0;
+        int width;
+
+        for (int i = 0; i < packNum; i++) {
+          value = 0;
+          width = packWidth;
+          while (width > 0) {
+            int m = width + currentByteOffset >= 8 ? 8 - currentByteOffset : 
width;
+            width -= m;
+            value = value << m;
+            int y = deltaBuf[index] & MASK[currentByteOffset];
+            currentByteOffset += m;
+            y >>>= (8 - currentByteOffset);
+            value |= y;
+            if (currentByteOffset == 8) {
+              currentByteOffset = 0;
+              ++index;
+            }
+          }
+          data[offset] = data[offset - 1] + minDeltaBase + value;
+          ++offset;
+        }
+      }
     }
 
-    @Override
     protected void readHeader(ByteBuffer buffer) {
       minDeltaBase = ReadWriteIOUtils.readLong(buffer);
       firstValue = ReadWriteIOUtils.readLong(buffer);
     }
 
-    @Override
     protected void allocateDataArray() {
       data = new long[packNum];
     }
 
     @Override
-    protected void readValue(int i) {
-      long v = BytesUtils.bytesToLong(deltaBuf, packWidth * i, packWidth);
-      data[i] = previous + minDeltaBase + v;
+    protected void readPack() {
+      long value;
+      int index = 0;
+      int currentByteOffset = 0;
+      int width;
+
+      for (int i = 0; i < packNum; i++) {
+        value = 0;
+        width = packWidth;
+        while (width > 0) {
+          int m = width + currentByteOffset >= 8 ? 8 - currentByteOffset : 
width;
+          width -= m;
+          value = value << m;
+          int y = deltaBuf[index] & MASK[currentByteOffset];
+          currentByteOffset += m;
+          y >>>= (8 - currentByteOffset);
+          value |= y;
+          if (currentByteOffset == 8) {
+            currentByteOffset = 0;
+            index++;
+          }
+        }
+
+        data[i] = previous + minDeltaBase + value;
+        previous = data[i];
+      }
     }
 
     @Override
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/TimePageReader.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/TimePageReader.java
index fc05e2b0ba2..982aab0d266 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/TimePageReader.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/TimePageReader.java
@@ -61,12 +61,10 @@ public class TimePageReader {
     return timeDecoder.readLong(timeBuffer);
   }
 
-  public long[] nextTimeBatch() throws IOException {
-    long[] timeBatch = new long[(int) pageHeader.getStatistics().getCount()];
-    int index = 0;
-    while (timeDecoder.hasNext(timeBuffer)) {
-      timeBatch[index++] = timeDecoder.readLong(timeBuffer);
-    }
+  public long[] nextTimeBatch() {
+    int length = (int) pageHeader.getStatistics().getCount();
+    long[] timeBatch = new long[length];
+    timeDecoder.readLong(timeBuffer, timeBatch, length);
     return timeBatch;
   }
 

Reply via email to