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

zyk pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 8224a70a41 [IOTDB-5398] Handle failure in schema query if there is 
exception occurs during the iteration (#8830)
8224a70a41 is described below

commit 8224a70a415f4b7ff31129e6c14d23f122e4c0d6
Author: Chen YZ <[email protected]>
AuthorDate: Wed Jan 11 20:31:01 2023 +0800

    [IOTDB-5398] Handle failure in schema query if there is exception occurs 
during the iteration (#8830)
---
 .../commons/schema/tree/AbstractTreeVisitor.java   |  2 +-
 .../db/metadata/mtree/MTreeBelowSGCachedImpl.java  | 30 +++++++
 .../db/metadata/mtree/MTreeBelowSGMemoryImpl.java  | 30 +++++++
 .../traverser/TraverserWithLimitOffsetWrapper.java | 10 +++
 .../db/metadata/query/reader/ISchemaReader.java    | 16 +++-
 .../schemaregion/SchemaRegionMemoryImpl.java       | 10 +++
 .../schemaregion/SchemaRegionSchemaFileImpl.java   | 10 +++
 .../schema/CountGroupByLevelScanOperator.java      |  3 +
 .../operator/schema/SchemaCountOperator.java       |  3 +
 .../operator/schema/SchemaQueryScanOperator.java   |  3 +
 .../schema/source/PathsUsingTemplateSource.java    | 26 ++++++
 .../schema/CountGroupByLevelMergeOperatorTest.java | 95 ++++++++++++++++------
 .../operator/schema/SchemaCountOperatorTest.java   | 88 ++++++++++----------
 .../operator/schema/SchemaOperatorTestUtil.java    | 66 +++++++++++++++
 .../schema/SchemaQueryScanOperatorTest.java        | 79 +++++++++---------
 15 files changed, 361 insertions(+), 110 deletions(-)

diff --git 
a/node-commons/src/main/java/org/apache/iotdb/commons/schema/tree/AbstractTreeVisitor.java
 
b/node-commons/src/main/java/org/apache/iotdb/commons/schema/tree/AbstractTreeVisitor.java
index 0255db8180..00db5f6d68 100644
--- 
a/node-commons/src/main/java/org/apache/iotdb/commons/schema/tree/AbstractTreeVisitor.java
+++ 
b/node-commons/src/main/java/org/apache/iotdb/commons/schema/tree/AbstractTreeVisitor.java
@@ -350,7 +350,7 @@ public abstract class AbstractTreeVisitor<N extends 
ITreeNode, R>
     this.throwable = e;
   }
 
-  protected Throwable getFailure() {
+  public Throwable getFailure() {
     return throwable;
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGCachedImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGCachedImpl.java
index dd1cad6cee..957019bd1b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGCachedImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGCachedImpl.java
@@ -998,6 +998,16 @@ public class MTreeBelowSGCachedImpl implements 
IMTreeBelowSG {
         new TraverserWithLimitOffsetWrapper<>(
             collector, showDevicesPlan.getLimit(), 
showDevicesPlan.getOffset());
     return new ISchemaReader<IDeviceSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return traverser.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return traverser.getFailure();
+      }
+
       @Override
       public void close() {
         traverser.close();
@@ -1041,6 +1051,16 @@ public class MTreeBelowSGCachedImpl implements 
IMTreeBelowSG {
         new TraverserWithLimitOffsetWrapper<>(
             collector, showTimeSeriesPlan.getLimit(), 
showTimeSeriesPlan.getOffset());
     return new ISchemaReader<ITimeSeriesSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return traverser.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return traverser.getFailure();
+      }
+
       @Override
       public void close() {
         traverser.close();
@@ -1071,6 +1091,16 @@ public class MTreeBelowSGCachedImpl implements 
IMTreeBelowSG {
         };
     collector.setTargetLevel(showNodesPlan.getLevel());
     return new ISchemaReader<INodeSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return collector.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return collector.getFailure();
+      }
+
       @Override
       public void close() {
         collector.close();
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
index f263dc3f7c..0872b26ac5 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
@@ -882,6 +882,16 @@ public class MTreeBelowSGMemoryImpl implements 
IMTreeBelowSG {
         new TraverserWithLimitOffsetWrapper<>(
             collector, showDevicesPlan.getLimit(), 
showDevicesPlan.getOffset());
     return new ISchemaReader<IDeviceSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return traverser.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return traverser.getFailure();
+      }
+
       @Override
       public void close() {
         traverser.close();
@@ -924,6 +934,16 @@ public class MTreeBelowSGMemoryImpl implements 
IMTreeBelowSG {
         new TraverserWithLimitOffsetWrapper<>(
             collector, showTimeSeriesPlan.getLimit(), 
showTimeSeriesPlan.getOffset());
     return new ISchemaReader<ITimeSeriesSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return traverser.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return traverser.getFailure();
+      }
+
       @Override
       public void close() {
         traverser.close();
@@ -955,6 +975,16 @@ public class MTreeBelowSGMemoryImpl implements 
IMTreeBelowSG {
     collector.setTargetLevel(showNodesPlan.getLevel());
 
     return new ISchemaReader<INodeSchemaInfo>() {
+      @Override
+      public boolean isSuccess() {
+        return collector.isSuccess();
+      }
+
+      @Override
+      public Throwable getFailure() {
+        return collector.getFailure();
+      }
+
       @Override
       public void close() {
         collector.close();
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/traverser/TraverserWithLimitOffsetWrapper.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/traverser/TraverserWithLimitOffsetWrapper.java
index f43468756e..d357e951e4 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/traverser/TraverserWithLimitOffsetWrapper.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/traverser/TraverserWithLimitOffsetWrapper.java
@@ -73,6 +73,16 @@ public class TraverserWithLimitOffsetWrapper<R> extends 
Traverser<R> {
     throw new UnsupportedOperationException();
   }
 
+  @Override
+  public boolean isSuccess() {
+    return traverser.isSuccess();
+  }
+
+  @Override
+  public Throwable getFailure() {
+    return traverser.getFailure();
+  }
+
   @Override
   protected boolean shouldVisitSubtreeOfInternalMatchedNode(IMNode node) {
     return false;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/query/reader/ISchemaReader.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/query/reader/ISchemaReader.java
index 21e00b102a..84e5a40987 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/query/reader/ISchemaReader.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/query/reader/ISchemaReader.java
@@ -23,4 +23,18 @@ import org.apache.iotdb.db.metadata.query.info.ISchemaInfo;
 
 import java.util.Iterator;
 
-public interface ISchemaReader<T extends ISchemaInfo> extends Iterator<T>, 
AutoCloseable {}
+public interface ISchemaReader<T extends ISchemaInfo> extends Iterator<T>, 
AutoCloseable {
+  /**
+   * Determines if the iteration is successful when it completes.
+   *
+   * @return False if an exception occurs during the iteration. Otherwise, 
return true.
+   */
+  boolean isSuccess();
+
+  /**
+   * Get throwable if there is exception occurs during the iteration.
+   *
+   * @return Throwable, null if no exception.
+   */
+  Throwable getFailure();
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
index c9d3633039..2696455c84 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
@@ -1190,6 +1190,16 @@ public class SchemaRegionMemoryImpl implements 
ISchemaRegion {
           showTimeseriesWithIndex(showTimeSeriesPlan);
       Iterator<ShowTimeSeriesResult> iterator = 
showTimeSeriesResultList.iterator();
       return new ISchemaReader<ITimeSeriesSchemaInfo>() {
+        @Override
+        public boolean isSuccess() {
+          return true;
+        }
+
+        @Override
+        public Throwable getFailure() {
+          return null;
+        }
+
         @Override
         public void close() {}
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
index eeabf63f1a..8ae626eae2 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
@@ -1273,6 +1273,16 @@ public class SchemaRegionSchemaFileImpl implements 
ISchemaRegion {
           showTimeseriesWithIndex(showTimeSeriesPlan);
       Iterator<ShowTimeSeriesResult> iterator = 
showTimeSeriesResultList.iterator();
       return new ISchemaReader<ITimeSeriesSchemaInfo>() {
+        @Override
+        public boolean isSuccess() {
+          return true;
+        }
+
+        @Override
+        public Throwable getFailure() {
+          return null;
+        }
+
         @Override
         public void close() {}
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelScanOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelScanOperator.java
index e1d4152ef2..0e3746e7b3 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelScanOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelScanOperator.java
@@ -125,6 +125,9 @@ public class CountGroupByLevelScanOperator<T extends 
ISchemaInfo> implements Sou
         break;
       }
     }
+    if (!schemaReader.isSuccess()) {
+      throw new RuntimeException(schemaReader.getFailure());
+    }
 
     TsBlockBuilder tsBlockBuilder = new TsBlockBuilder(OUTPUT_DATA_TYPES);
     for (Map.Entry<PartialPath, Long> entry : countMap.entrySet()) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperator.java
index 7de9660365..5e313e9fea 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperator.java
@@ -83,6 +83,9 @@ public class SchemaCountOperator<T extends ISchemaInfo> 
implements SourceOperato
       schemaReader.next();
       count++;
     }
+    if (!schemaReader.isSuccess()) {
+      throw new RuntimeException(schemaReader.getFailure());
+    }
 
     tsBlockBuilder.getTimeColumnBuilder().writeLong(0L);
     tsBlockBuilder.getColumnBuilder(0).writeLong(count);
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperator.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperator.java
index ab05353858..ee27182135 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperator.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperator.java
@@ -139,6 +139,9 @@ public class SchemaQueryScanOperator<T extends ISchemaInfo> 
implements SourceOpe
         break;
       }
     }
+    if (!schemaReader.isSuccess()) {
+      throw new RuntimeException(schemaReader.getFailure());
+    }
     return tsBlockBuilder.build();
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/source/PathsUsingTemplateSource.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/source/PathsUsingTemplateSource.java
index b3c121f322..ba9b450244 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/source/PathsUsingTemplateSource.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/execution/operator/schema/source/PathsUsingTemplateSource.java
@@ -67,12 +67,15 @@ public class PathsUsingTemplateSource implements 
ISchemaSource<IDeviceSchemaInfo
 
     final ISchemaRegion schemaRegion;
 
+    private Throwable throwable;
+
     ISchemaReader<IDeviceSchemaInfo> currentDeviceReader;
 
     DevicesUsingTemplateReader(
         Iterator<PartialPath> pathPatternIterator, ISchemaRegion schemaRegion) 
{
       this.pathPatternIterator = pathPatternIterator;
       this.schemaRegion = schemaRegion;
+      this.throwable = null;
     }
 
     @Override
@@ -85,11 +88,18 @@ public class PathsUsingTemplateSource implements 
ISchemaSource<IDeviceSchemaInfo
     @Override
     public boolean hasNext() {
       try {
+        if (throwable != null) {
+          return false;
+        }
         if (currentDeviceReader != null) {
           if (currentDeviceReader.hasNext()) {
             return true;
           } else {
             currentDeviceReader.close();
+            if (!currentDeviceReader.isSuccess()) {
+              throwable = currentDeviceReader.getFailure();
+              return false;
+            }
           }
         }
 
@@ -114,5 +124,21 @@ public class PathsUsingTemplateSource implements 
ISchemaSource<IDeviceSchemaInfo
     public IDeviceSchemaInfo next() {
       return currentDeviceReader.next();
     }
+
+    @Override
+    public boolean isSuccess() {
+      return throwable == null && (currentDeviceReader == null || 
currentDeviceReader.isSuccess());
+    }
+
+    @Override
+    public Throwable getFailure() {
+      if (throwable != null) {
+        return throwable;
+      } else if (currentDeviceReader != null) {
+        return currentDeviceReader.getFailure();
+      } else {
+        return null;
+      }
+    }
   }
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelMergeOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelMergeOperatorTest.java
index f5a50001b6..e503c460de 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelMergeOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/CountGroupByLevelMergeOperatorTest.java
@@ -22,7 +22,6 @@ import 
org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.metadata.query.info.ITimeSeriesSchemaInfo;
-import org.apache.iotdb.db.metadata.query.reader.ISchemaReader;
 import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.common.PlanFragmentId;
@@ -49,8 +48,10 @@ import java.util.Set;
 import java.util.concurrent.ExecutorService;
 
 import static 
org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext.createFragmentInstanceContext;
+import static 
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaOperatorTestUtil.EXCEPTION_MESSAGE;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
 
@@ -82,14 +83,14 @@ public class CountGroupByLevelMergeOperatorTest {
               planNodeId,
               driverContext.getOperatorContexts().get(0),
               2,
-              mockSchemaSource(schemaRegion, new PartialPath(OPERATOR_TEST_SG 
+ ".device2")));
+              mockSchemaSource(schemaRegion, new PartialPath(OPERATOR_TEST_SG 
+ ".device2"), true));
 
       CountGroupByLevelScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesCountOperator2 =
           new CountGroupByLevelScanOperator<>(
               planNodeId,
               driverContext.getOperatorContexts().get(0),
               2,
-              mockSchemaSource(schemaRegion, new 
PartialPath(OPERATOR_TEST_SG)));
+              mockSchemaSource(schemaRegion, new 
PartialPath(OPERATOR_TEST_SG), true));
 
       CountGroupByLevelMergeOperator mergeOperator =
           new CountGroupByLevelMergeOperator(
@@ -132,21 +133,80 @@ public class CountGroupByLevelMergeOperatorTest {
     }
   }
 
+  @Test
+  public void testCountScanOperator() {
+    ExecutorService instanceNotificationExecutor =
+        IoTDBThreadPoolFactory.newFixedThreadPool(1, 
"test-instance-notification");
+    try {
+      QueryId queryId = new QueryId("stub_query");
+      FragmentInstanceId instanceId =
+          new FragmentInstanceId(new PlanFragmentId(queryId, 0), 
"stub-instance");
+      FragmentInstanceStateMachine stateMachine =
+          new FragmentInstanceStateMachine(instanceId, 
instanceNotificationExecutor);
+      FragmentInstanceContext fragmentInstanceContext =
+          createFragmentInstanceContext(instanceId, stateMachine);
+      DriverContext driverContext = new DriverContext(fragmentInstanceContext, 
0);
+      PlanNodeId planNodeId = queryId.genPlanNodeId();
+      OperatorContext operatorContext =
+          driverContext.addOperatorContext(
+              1, planNodeId, 
CountGroupByLevelScanOperator.class.getSimpleName());
+      ISchemaRegion schemaRegion = Mockito.mock(ISchemaRegion.class);
+      operatorContext.setDriverContext(
+          new SchemaDriverContext(fragmentInstanceContext, schemaRegion));
+      CountGroupByLevelScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesCountOperator =
+          new CountGroupByLevelScanOperator<>(
+              planNodeId,
+              driverContext.getOperatorContexts().get(0),
+              2,
+              mockSchemaSource(schemaRegion, new 
PartialPath(OPERATOR_TEST_SG), true));
+      TsBlock tsBlock = null;
+      while (timeSeriesCountOperator.hasNext()) {
+        tsBlock = timeSeriesCountOperator.next();
+        for (int i = 0; i < tsBlock.getPositionCount(); i++) {
+          assertEquals(1, tsBlock.getColumn(1).getLong(i));
+        }
+      }
+      assertNotNull(tsBlock);
+
+      // Assert failure if exception occurs
+      CountGroupByLevelScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesCountOperatorFailure =
+          new CountGroupByLevelScanOperator<>(
+              planNodeId,
+              driverContext.getOperatorContexts().get(0),
+              2,
+              mockSchemaSource(schemaRegion, new 
PartialPath(OPERATOR_TEST_SG), false));
+      try {
+        while (timeSeriesCountOperatorFailure.hasNext()) {
+          timeSeriesCountOperatorFailure.next();
+        }
+      } catch (RuntimeException e) {
+        Assert.assertTrue(e.getMessage().contains(EXCEPTION_MESSAGE));
+      }
+    } catch (Exception e) {
+      e.printStackTrace();
+      fail();
+    } finally {
+      instanceNotificationExecutor.shutdown();
+    }
+  }
+
   private ISchemaSource<ITimeSeriesSchemaInfo> mockSchemaSource(
-      ISchemaRegion schemaRegion, PartialPath path) throws Exception {
+      ISchemaRegion schemaRegion, PartialPath path, boolean success) throws 
Exception {
     ISchemaSource<ITimeSeriesSchemaInfo> schemaSource = 
Mockito.mock(ISchemaSource.class);
     if (path.equals(new PartialPath(OPERATOR_TEST_SG + ".device2"))) {
-      ISchemaReader<ITimeSeriesSchemaInfo> schemaReader =
-          mockSchemaReader(10, OPERATOR_TEST_SG + ".device2");
-      
Mockito.when(schemaSource.getSchemaReader(schemaRegion)).thenReturn(schemaReader);
+      mockSchemaReader(schemaSource, schemaRegion, 10, OPERATOR_TEST_SG + 
".device2", success);
     } else if (path.equals(new PartialPath(OPERATOR_TEST_SG))) {
-      ISchemaReader<ITimeSeriesSchemaInfo> schemaReader = 
mockSchemaReader(2000, OPERATOR_TEST_SG);
-      
Mockito.when(schemaSource.getSchemaReader(schemaRegion)).thenReturn(schemaReader);
+      mockSchemaReader(schemaSource, schemaRegion, 2000, OPERATOR_TEST_SG, 
success);
     }
     return schemaSource;
   }
 
-  private ISchemaReader<ITimeSeriesSchemaInfo> mockSchemaReader(int 
expectedNum, String prefix)
+  private void mockSchemaReader(
+      ISchemaSource<ITimeSeriesSchemaInfo> schemaSource,
+      ISchemaRegion schemaRegion,
+      int expectedNum,
+      String prefix,
+      boolean success)
       throws IllegalPathException {
     List<ITimeSeriesSchemaInfo> timeSeriesSchemaInfoList = new 
ArrayList<>(expectedNum);
     for (int i = 0; i < expectedNum; i++) {
@@ -156,19 +216,6 @@ public class CountGroupByLevelMergeOperatorTest {
       timeSeriesSchemaInfoList.add(timeSeriesSchemaInfo);
     }
     Iterator<ITimeSeriesSchemaInfo> iterator = 
timeSeriesSchemaInfoList.iterator();
-    return new ISchemaReader<ITimeSeriesSchemaInfo>() {
-      @Override
-      public void close() throws Exception {}
-
-      @Override
-      public boolean hasNext() {
-        return iterator.hasNext();
-      }
-
-      @Override
-      public ITimeSeriesSchemaInfo next() {
-        return iterator.next();
-      }
-    };
+    SchemaOperatorTestUtil.mockGetSchemaReader(schemaSource, iterator, 
schemaRegion, success);
   }
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperatorTest.java
index c818b729cb..2e9fa0e2db 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaCountOperatorTest.java
@@ -19,11 +19,9 @@
 package org.apache.iotdb.db.mpp.execution.operator.schema;
 
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
-import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.metadata.query.info.ISchemaInfo;
 import org.apache.iotdb.db.metadata.query.info.ITimeSeriesSchemaInfo;
-import org.apache.iotdb.db.metadata.query.reader.ISchemaReader;
 import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.common.PlanFragmentId;
@@ -37,6 +35,7 @@ import 
org.apache.iotdb.db.mpp.execution.operator.schema.source.ISchemaSource;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId;
 import org.apache.iotdb.tsfile.read.common.block.TsBlock;
 
+import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.Mockito;
 
@@ -46,6 +45,7 @@ import java.util.List;
 import java.util.concurrent.ExecutorService;
 
 import static 
org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext.createFragmentInstanceContext;
+import static 
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaOperatorTestUtil.EXCEPTION_MESSAGE;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertTrue;
@@ -74,29 +74,14 @@ public class SchemaCountOperatorTest {
               1, planNodeId, SchemaCountOperator.class.getSimpleName());
       operatorContext.setDriverContext(
           new SchemaDriverContext(fragmentInstanceContext, schemaRegion));
-      ISchemaSource<?> schemaSource = Mockito.mock(ISchemaSource.class);
+      ISchemaSource<ISchemaInfo> schemaSource = 
Mockito.mock(ISchemaSource.class);
 
       List<ISchemaInfo> schemaInfoList = new ArrayList<>(10);
       for (int i = 0; i < 10; i++) {
         schemaInfoList.add(Mockito.mock(ISchemaInfo.class));
       }
-      Iterator<ISchemaInfo> iterator = schemaInfoList.iterator();
-      Mockito.when(schemaSource.getSchemaReader(schemaRegion))
-          .thenReturn(
-              new ISchemaReader() {
-                @Override
-                public void close() throws Exception {}
-
-                @Override
-                public boolean hasNext() {
-                  return iterator.hasNext();
-                }
-
-                @Override
-                public ISchemaInfo next() {
-                  return iterator.next();
-                }
-              });
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          schemaSource, schemaInfoList.iterator(), schemaRegion, true);
 
       SchemaCountOperator<?> schemaCountOperator =
           new SchemaCountOperator<>(
@@ -107,7 +92,22 @@ public class SchemaCountOperatorTest {
         tsBlock = schemaCountOperator.next();
       }
       assertNotNull(tsBlock);
-      assertEquals(tsBlock.getColumn(0).getLong(0), 10);
+      assertEquals(10, tsBlock.getColumn(0).getLong(0));
+
+      // Assert failure if exception occurs
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          schemaSource, schemaInfoList.iterator(), schemaRegion, false);
+      try {
+        SchemaCountOperator<?> schemaCountOperatorFailure =
+            new SchemaCountOperator<>(
+                planNodeId, driverContext.getOperatorContexts().get(0), 
schemaSource);
+        while (schemaCountOperatorFailure.hasNext()) {
+          schemaCountOperatorFailure.next();
+        }
+        Assert.fail();
+      } catch (RuntimeException e) {
+        Assert.assertTrue(e.getMessage().contains(EXCEPTION_MESSAGE));
+      }
     } finally {
       instanceNotificationExecutor.shutdown();
     }
@@ -139,8 +139,8 @@ public class SchemaCountOperatorTest {
               planNodeId,
               driverContext.getOperatorContexts().get(0),
               2,
-              mockSchemaSource(schemaRegion));
-      TsBlock tsBlock = null;
+              mockSchemaSource(schemaRegion, true));
+      TsBlock tsBlock;
       List<TsBlock> tsBlockList = collectResult(timeSeriesCountOperator);
       assertEquals(1, tsBlockList.size());
       tsBlock = tsBlockList.get(0);
@@ -156,12 +156,26 @@ public class SchemaCountOperatorTest {
               planNodeId,
               driverContext.getOperatorContexts().get(0),
               1,
-              mockSchemaSource(schemaRegion));
+              mockSchemaSource(schemaRegion, true));
       tsBlockList = collectResult(timeSeriesCountOperator2);
       assertEquals(1, tsBlockList.size());
       tsBlock = tsBlockList.get(0);
 
       assertEquals(100, tsBlock.getColumn(1).getLong(0));
+
+      // Assert failure if exception occurs
+      try {
+        CountGroupByLevelScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesCountOperatorFailure =
+            new CountGroupByLevelScanOperator<>(
+                planNodeId,
+                driverContext.getOperatorContexts().get(0),
+                1,
+                mockSchemaSource(schemaRegion, false));
+        collectResult(timeSeriesCountOperatorFailure);
+        Assert.fail();
+      } catch (RuntimeException e) {
+        Assert.assertTrue(e.getMessage().contains(EXCEPTION_MESSAGE));
+      }
     } catch (Exception e) {
       e.printStackTrace();
       fail();
@@ -182,15 +196,9 @@ public class SchemaCountOperatorTest {
     return tsBlocks;
   }
 
-  private ISchemaSource<ITimeSeriesSchemaInfo> mockSchemaSource(ISchemaRegion 
schemaRegion)
-      throws Exception {
+  private ISchemaSource<ITimeSeriesSchemaInfo> mockSchemaSource(
+      ISchemaRegion schemaRegion, boolean success) throws Exception {
     ISchemaSource<ITimeSeriesSchemaInfo> schemaSource = 
Mockito.mock(ISchemaSource.class);
-    ISchemaReader<ITimeSeriesSchemaInfo> schemaReader = mockSchemaReader();
-    
Mockito.when(schemaSource.getSchemaReader(schemaRegion)).thenReturn(schemaReader);
-    return schemaSource;
-  }
-
-  private ISchemaReader<ITimeSeriesSchemaInfo> mockSchemaReader() throws 
IllegalPathException {
     List<ITimeSeriesSchemaInfo> timeSeriesSchemaInfoList = new 
ArrayList<>(1000);
     for (int i = 0; i < 10; i++) {
       for (int j = 0; j < 10; j++) {
@@ -201,19 +209,7 @@ public class SchemaCountOperatorTest {
       }
     }
     Iterator<ITimeSeriesSchemaInfo> iterator = 
timeSeriesSchemaInfoList.iterator();
-    return new ISchemaReader<ITimeSeriesSchemaInfo>() {
-      @Override
-      public void close() throws Exception {}
-
-      @Override
-      public boolean hasNext() {
-        return iterator.hasNext();
-      }
-
-      @Override
-      public ITimeSeriesSchemaInfo next() {
-        return iterator.next();
-      }
-    };
+    SchemaOperatorTestUtil.mockGetSchemaReader(schemaSource, iterator, 
schemaRegion, success);
+    return schemaSource;
   }
 }
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaOperatorTestUtil.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaOperatorTestUtil.java
new file mode 100644
index 0000000000..882076f6fc
--- /dev/null
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaOperatorTestUtil.java
@@ -0,0 +1,66 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.db.mpp.execution.operator.schema;
+
+import org.apache.iotdb.commons.exception.MetadataException;
+import org.apache.iotdb.db.metadata.query.info.ISchemaInfo;
+import org.apache.iotdb.db.metadata.query.reader.ISchemaReader;
+import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion;
+import org.apache.iotdb.db.mpp.execution.operator.schema.source.ISchemaSource;
+
+import org.mockito.Mockito;
+
+import java.util.Iterator;
+
+public class SchemaOperatorTestUtil {
+  public static final String EXCEPTION_MESSAGE = "ExceptionMessage";
+
+  public static <T extends ISchemaInfo> void mockGetSchemaReader(
+      ISchemaSource<T> schemaSource,
+      Iterator<T> iterator,
+      ISchemaRegion schemaRegion,
+      boolean isSuccess) {
+    Mockito.when(schemaSource.getSchemaReader(schemaRegion))
+        .thenReturn(
+            new ISchemaReader<T>() {
+              @Override
+              public boolean isSuccess() {
+                return isSuccess;
+              }
+
+              @Override
+              public Throwable getFailure() {
+                return isSuccess ? null : new 
MetadataException(EXCEPTION_MESSAGE);
+              }
+
+              @Override
+              public void close() {}
+
+              @Override
+              public boolean hasNext() {
+                return iterator.hasNext();
+              }
+
+              @Override
+              public T next() {
+                return iterator.next();
+              }
+            });
+  }
+}
diff --git 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperatorTest.java
 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperatorTest.java
index 8a3914d8a7..a9b501324d 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperatorTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/mpp/execution/operator/schema/SchemaQueryScanOperatorTest.java
@@ -23,7 +23,6 @@ import org.apache.iotdb.commons.exception.MetadataException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.metadata.query.info.IDeviceSchemaInfo;
 import org.apache.iotdb.db.metadata.query.info.ITimeSeriesSchemaInfo;
-import org.apache.iotdb.db.metadata.query.reader.ISchemaReader;
 import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion;
 import org.apache.iotdb.db.mpp.common.FragmentInstanceId;
 import org.apache.iotdb.db.mpp.common.PlanFragmentId;
@@ -52,11 +51,11 @@ import org.mockito.Mockito;
 
 import java.util.ArrayList;
 import java.util.Collections;
-import java.util.Iterator;
 import java.util.List;
 import java.util.concurrent.ExecutorService;
 
 import static 
org.apache.iotdb.db.mpp.execution.fragment.FragmentInstanceContext.createFragmentInstanceContext;
+import static 
org.apache.iotdb.db.mpp.execution.operator.schema.SchemaOperatorTestUtil.EXCEPTION_MESSAGE;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
@@ -88,34 +87,22 @@ public class SchemaQueryScanOperatorTest {
       Mockito.when(deviceSchemaInfo.getFullPath())
           .thenReturn(META_SCAN_OPERATOR_TEST_SG + ".device0");
       Mockito.when(deviceSchemaInfo.isAligned()).thenReturn(false);
-      Iterator<IDeviceSchemaInfo> iterator = 
Collections.singletonList(deviceSchemaInfo).iterator();
       operatorContext.setDriverContext(
           new SchemaDriverContext(fragmentInstanceContext, schemaRegion));
       ISchemaSource<IDeviceSchemaInfo> deviceSchemaSource =
           SchemaSourceFactory.getDeviceSchemaSource(partialPath, false, 10, 0, 
true);
-      Mockito.when(deviceSchemaSource.getSchemaReader(schemaRegion))
-          .thenReturn(
-              new ISchemaReader<IDeviceSchemaInfo>() {
-                @Override
-                public void close() throws Exception {}
-
-                @Override
-                public boolean hasNext() {
-                  return iterator.hasNext();
-                }
-
-                @Override
-                public IDeviceSchemaInfo next() {
-                  return iterator.next();
-                }
-              });
-
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          deviceSchemaSource,
+          Collections.singletonList(deviceSchemaInfo).iterator(),
+          schemaRegion,
+          true);
+      //
       List<ColumnHeader> columns = 
deviceSchemaSource.getInfoQueryColumnHeaders();
 
       SchemaQueryScanOperator<IDeviceSchemaInfo> devicesSchemaScanOperator =
           new SchemaQueryScanOperator<>(
               planNodeId, driverContext.getOperatorContexts().get(0), 
deviceSchemaSource);
-
+      //
       while (devicesSchemaScanOperator.hasNext()) {
         TsBlock tsBlock = devicesSchemaScanOperator.next();
         assertEquals(3, tsBlock.getValueColumnCount());
@@ -143,6 +130,23 @@ public class SchemaQueryScanOperatorTest {
           }
         }
       }
+      // Assert failure if exception occurs
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          deviceSchemaSource,
+          Collections.singletonList(deviceSchemaInfo).iterator(),
+          schemaRegion,
+          false);
+      try {
+        SchemaQueryScanOperator<IDeviceSchemaInfo> 
devicesSchemaScanOperatorFailure =
+            new SchemaQueryScanOperator<>(
+                planNodeId, driverContext.getOperatorContexts().get(0), 
deviceSchemaSource);
+        while (devicesSchemaScanOperatorFailure.hasNext()) {
+          devicesSchemaScanOperatorFailure.next();
+        }
+        Assert.fail();
+      } catch (RuntimeException e) {
+        Assert.assertTrue(e.getMessage().contains(EXCEPTION_MESSAGE));
+      }
     } catch (MetadataException e) {
       e.printStackTrace();
       fail();
@@ -184,7 +188,6 @@ public class SchemaQueryScanOperatorTest {
         Mockito.when(timeSeriesSchemaInfo.getAttributes()).thenReturn(null);
         showTimeSeriesResults.add(timeSeriesSchemaInfo);
       }
-      Iterator<ITimeSeriesSchemaInfo> iterator = 
showTimeSeriesResults.iterator();
 
       ISchemaRegion schemaRegion = Mockito.mock(ISchemaRegion.class);
       
Mockito.when(schemaRegion.getStorageGroupFullPath()).thenReturn(META_SCAN_OPERATOR_TEST_SG);
@@ -194,22 +197,8 @@ public class SchemaQueryScanOperatorTest {
       ISchemaSource<ITimeSeriesSchemaInfo> timeSeriesSchemaSource =
           SchemaSourceFactory.getTimeSeriesSchemaSource(
               partialPath, false, 10, 0, null, null, false, 
Collections.emptyMap());
-      Mockito.when(timeSeriesSchemaSource.getSchemaReader(schemaRegion))
-          .thenReturn(
-              new ISchemaReader<ITimeSeriesSchemaInfo>() {
-                @Override
-                public void close() throws Exception {}
-
-                @Override
-                public boolean hasNext() {
-                  return iterator.hasNext();
-                }
-
-                @Override
-                public ITimeSeriesSchemaInfo next() {
-                  return iterator.next();
-                }
-              });
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          timeSeriesSchemaSource, showTimeSeriesResults.iterator(), 
schemaRegion, true);
 
       SchemaQueryScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesMetaScanOperator =
           new SchemaQueryScanOperator<>(
@@ -254,6 +243,20 @@ public class SchemaQueryScanOperatorTest {
           }
         }
       }
+      // Assert failure if exception occurs
+      SchemaOperatorTestUtil.mockGetSchemaReader(
+          timeSeriesSchemaSource, showTimeSeriesResults.iterator(), 
schemaRegion, false);
+      try {
+        SchemaQueryScanOperator<ITimeSeriesSchemaInfo> 
timeSeriesMetaScanOperatorFailure =
+            new SchemaQueryScanOperator<>(
+                planNodeId, driverContext.getOperatorContexts().get(0), 
timeSeriesSchemaSource);
+        while (timeSeriesMetaScanOperatorFailure.hasNext()) {
+          timeSeriesMetaScanOperatorFailure.next();
+        }
+        Assert.fail();
+      } catch (RuntimeException e) {
+        Assert.assertTrue(e.getMessage().contains(EXCEPTION_MESSAGE));
+      }
     } catch (MetadataException e) {
       e.printStackTrace();
       fail();


Reply via email to