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

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


The following commit(s) were added to refs/heads/mpp-ty by this push:
     new 4edf84e  add all current useful operators
4edf84e is described below

commit 4edf84e44beb8f5a97b30feb1c0933bfbe0005ea
Author: JackieTien97 <[email protected]>
AuthorDate: Fri Mar 18 22:21:11 2022 +0800

    add all current useful operators
---
 ...SinkOperator.java => FragmentSinkOperator.java} | 53 ++++++++++++++--------
 .../iotdb/db/mpp/operator/sink/SinkOperator.java   |  5 +-
 2 files changed, 37 insertions(+), 21 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/FragmentSinkOperator.java
similarity index 53%
copy from 
server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
copy to 
server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/FragmentSinkOperator.java
index b03e7ec..aca892f 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/FragmentSinkOperator.java
@@ -18,27 +18,44 @@
  */
 package org.apache.iotdb.db.mpp.operator.sink;
 
-import org.apache.iotdb.db.mpp.operator.Operator;
+import org.apache.iotdb.db.mpp.common.TsBlock;
+import org.apache.iotdb.db.mpp.operator.OperatorContext;
 
-import java.nio.ByteBuffer;
+import com.google.common.util.concurrent.ListenableFuture;
 
-public interface SinkOperator extends Operator {
+public class FragmentSinkOperator implements SinkOperator {
 
-  /**
-   * Sends a tsBlock to an unpartitioned buffer. If no-more-tsBlocks has been 
set, the send tsBlock
-   * call is ignored. This can happen with limit queries.
-   */
-  void send(ByteBuffer tsBlock);
+  @Override
+  public OperatorContext getOperatorContext() {
+    return null;
+  }
 
-  /**
-   * Notify SinkHandle that no more tsBlocks will be sent. Any future calls to 
send a tsBlock are
-   * ignored.
-   */
-  void setNoMoreTsBlocks();
+  @Override
+  public ListenableFuture<Void> isBlocked() {
+    return SinkOperator.super.isBlocked();
+  }
 
-  /**
-   * Abort the sink handle, discarding all tsBlocks which may still in memory 
buffer, but blocking
-   * readers. It is expected that readers will be unblocked when the failed 
query is cleaned up.
-   */
-  void abort();
+  @Override
+  public TsBlock next() {
+    return null;
+  }
+
+  @Override
+  public boolean hasNext() {
+    return false;
+  }
+
+  @Override
+  public void close() throws Exception {
+    SinkOperator.super.close();
+  }
+
+  @Override
+  public void send(TsBlock tsBlock) {}
+
+  @Override
+  public void setNoMoreTsBlocks() {}
+
+  @Override
+  public void abort() {}
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
index b03e7ec..c3f16ca 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/operator/sink/SinkOperator.java
@@ -18,17 +18,16 @@
  */
 package org.apache.iotdb.db.mpp.operator.sink;
 
+import org.apache.iotdb.db.mpp.common.TsBlock;
 import org.apache.iotdb.db.mpp.operator.Operator;
 
-import java.nio.ByteBuffer;
-
 public interface SinkOperator extends Operator {
 
   /**
    * Sends a tsBlock to an unpartitioned buffer. If no-more-tsBlocks has been 
set, the send tsBlock
    * call is ignored. This can happen with limit queries.
    */
-  void send(ByteBuffer tsBlock);
+  void send(TsBlock tsBlock);
 
   /**
    * Notify SinkHandle that no more tsBlocks will be sent. Any future calls to 
send a tsBlock are

Reply via email to