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

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


The following commit(s) were added to refs/heads/master by this push:
     new 4e2ce58  add sessionpool example
4e2ce58 is described below

commit 4e2ce58888d202daf86031f97fec8838132192df
Author: qiaojialin <[email protected]>
AuthorDate: Wed May 13 11:42:55 2020 +0800

    add sessionpool example
---
 .../java/org/apache/iotdb/SessionPoolExample.java  | 88 +++++++++++-----------
 .../iotdb/session/pool/SessionDataSetWrapper.java  |  5 ++
 .../org/apache/iotdb/session/pool/SessionPool.java |  4 +-
 .../apache/iotdb/session/pool/SessionPoolTest.java |  7 +-
 4 files changed, 54 insertions(+), 50 deletions(-)

diff --git 
a/example/session/src/main/java/org/apache/iotdb/SessionPoolExample.java 
b/example/session/src/main/java/org/apache/iotdb/SessionPoolExample.java
index 607c280..96a660d 100644
--- a/example/session/src/main/java/org/apache/iotdb/SessionPoolExample.java
+++ b/example/session/src/main/java/org/apache/iotdb/SessionPoolExample.java
@@ -18,11 +18,14 @@
  */
 package org.apache.iotdb;
 
+import java.util.ArrayList;
 import java.util.Collections;
+import java.util.List;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import org.apache.iotdb.rpc.IoTDBConnectionException;
 import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.session.SessionDataSet.DataIterator;
 import org.apache.iotdb.session.pool.SessionDataSetWrapper;
 import org.apache.iotdb.session.pool.SessionPool;
 
@@ -32,79 +35,80 @@ public class SessionPoolExample {
   private static ExecutorService service;
 
   public static void main(String[] args)
-      throws StatementExecutionException, IoTDBConnectionException {
+      throws StatementExecutionException, IoTDBConnectionException, 
InterruptedException {
     pool = new SessionPool("127.0.0.1", 6667, "root", "root", 3);
+    service = Executors.newFixedThreadPool(10);
 
     insertRecord();
-    queryGettingAllData();
-    queryGettingPartData();
-    deleteData();
+    queryByRowRecord();
+    Thread.sleep(1000);
+    queryByIterator();
     pool.close();
+    service.shutdown();
   }
 
   // more insert example, see SessionExample.java
   private static void insertRecord() throws StatementExecutionException, 
IoTDBConnectionException {
-    for (int i = 0; i < 10; i++) {
-      pool.insertRecord("root.sg1.d1", i, Collections.singletonList("s" + i),
-          Collections.singletonList("" + i));
+    String deviceId = "root.sg1.d1";
+    List<String> measurements = new ArrayList<>();
+    measurements.add("s1");
+    measurements.add("s2");
+    measurements.add("s3");
+    for (long time = 0; time < 10; time++) {
+      List<String> values = new ArrayList<>();
+      values.add("1");
+      values.add("2");
+      values.add("3");
+      pool.insertRecord(deviceId, time, measurements, values);
     }
   }
 
-  // Query getting all data
-  private static void queryGettingAllData() {
-    service = Executors.newFixedThreadPool(10);
-    for (int i = 0; i < 10; i++) {
-      final int no = i;
+  private static void queryByRowRecord() {
+    for (int i = 0; i < 1; i++) {
       service.submit(() -> {
         SessionDataSetWrapper wrapper = null;
         try {
-          wrapper = pool
-              .executeQueryStatement("select * from root.sg1.d1 where time = " 
+ no);
-          // if you get all data, don't calling closeResultSet() is OK.
-          // but it's suggested.
+          wrapper = pool.executeQueryStatement("select * from root.sg1.d1");
+          System.out.println(wrapper.getColumnNames());
+          System.out.println(wrapper.getColumnTypes());
           while (wrapper.hasNext()) {
             System.out.println(wrapper.next());
           }
         } catch (IoTDBConnectionException | StatementExecutionException e) {
-          // if there is exception when you call 
SessionDataSetWrapper.hasNext() or next(),
-          // you have to call closeResultSet().
-          try {
-            pool.closeResultSet(wrapper);
-          } catch (StatementExecutionException ex) {
-            ex.printStackTrace();
-          }
           e.printStackTrace();
+        } finally {
+          // remember to close data set finally!
+          pool.closeResultSet(wrapper);
         }
       });
     }
-    service.shutdown();
   }
 
-  // Query getting part data
-  private static void queryGettingPartData() {
-    service = Executors.newFixedThreadPool(10);
-    for (int i = 0; i < 10; i++) {
-      final int no = i;
+  private static void queryByIterator() {
+    for (int i = 0; i < 1; i++) {
       service.submit(() -> {
+        SessionDataSetWrapper wrapper = null;
         try {
-          SessionDataSetWrapper wrapper = pool
-              .executeQueryStatement("select * from root.sg1.d1 where time = " 
+ no);
-          wrapper.next();
-          // REMEMBER to call closeResultSet() here.
-          // otherwise the query left will be blocked.
-          pool.closeResultSet(wrapper);
+          wrapper = pool.executeQueryStatement("select * from root.sg1.d1");
+          // get DataIterator like JDBC
+          DataIterator dataIterator = wrapper.iterator();
+          System.out.println(wrapper.getColumnNames());
+          System.out.println(wrapper.getColumnTypes());
+          while (dataIterator.next()) {
+            StringBuilder builder = new StringBuilder();
+            for (String columnName: wrapper.getColumnNames()) {
+              builder.append(dataIterator.getString(columnName) + " ");
+            }
+            System.out.println(builder.toString());
+          }
         } catch (IoTDBConnectionException | StatementExecutionException e) {
           e.printStackTrace();
+        } finally {
+          // remember to close data set finally!
+          pool.closeResultSet(wrapper);
         }
       });
     }
-    service.shutdown();
-  }
-
-  private static void deleteData() throws IoTDBConnectionException, 
StatementExecutionException {
-    String path = "root.sg1.d1";
-    long deleteTime = 99;
-    pool.deleteData(path, deleteTime);
   }
 
 }
diff --git 
a/session/src/main/java/org/apache/iotdb/session/pool/SessionDataSetWrapper.java
 
b/session/src/main/java/org/apache/iotdb/session/pool/SessionDataSetWrapper.java
index b319b86..53df42d 100644
--- 
a/session/src/main/java/org/apache/iotdb/session/pool/SessionDataSetWrapper.java
+++ 
b/session/src/main/java/org/apache/iotdb/session/pool/SessionDataSetWrapper.java
@@ -24,6 +24,7 @@ import org.apache.iotdb.rpc.StatementExecutionException;
 import org.apache.iotdb.session.Session;
 import org.apache.iotdb.session.SessionDataSet;
 import org.apache.iotdb.session.SessionDataSet.DataIterator;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.read.common.RowRecord;
 
 public class SessionDataSetWrapper {
@@ -86,4 +87,8 @@ public class SessionDataSetWrapper {
   public List<String> getColumnNames() {
     return sessionDataSet.getColumnNames();
   }
+
+  public List<TSDataType> getColumnTypes() {
+    return sessionDataSet.getColumnTypes();
+  }
 }
diff --git 
a/session/src/main/java/org/apache/iotdb/session/pool/SessionPool.java 
b/session/src/main/java/org/apache/iotdb/session/pool/SessionPool.java
index 62e0356..7ad90ce 100644
--- a/session/src/main/java/org/apache/iotdb/session/pool/SessionPool.java
+++ b/session/src/main/java/org/apache/iotdb/session/pool/SessionPool.java
@@ -189,11 +189,11 @@ public class SessionPool {
     occupied.clear();
   }
 
-  public void closeResultSet(SessionDataSetWrapper wrapper) throws 
StatementExecutionException {
+  public void closeResultSet(SessionDataSetWrapper wrapper) {
     boolean putback = true;
     try {
       wrapper.sessionDataSet.closeOperationHandle();
-    } catch (IoTDBConnectionException e) {
+    } catch (IoTDBConnectionException | StatementExecutionException e) {
       removeSession();
       putback = false;
     } finally {
diff --git 
a/session/src/test/java/org/apache/iotdb/session/pool/SessionPoolTest.java 
b/session/src/test/java/org/apache/iotdb/session/pool/SessionPoolTest.java
index 0e7ce31..7df3d1f 100644
--- a/session/src/test/java/org/apache/iotdb/session/pool/SessionPoolTest.java
+++ b/session/src/test/java/org/apache/iotdb/session/pool/SessionPoolTest.java
@@ -193,12 +193,7 @@ public class SessionPoolTest {
         wrapper.next();
       }
     } catch (IoTDBConnectionException e) {
-      try {
-        pool.closeResultSet(wrapper);
-      } catch (StatementExecutionException ex) {
-        ex.printStackTrace();
-        fail();
-      }
+      pool.closeResultSet(wrapper);
       EnvironmentUtils.reactiveDaemon();
       correctQuery(pool);
       pool.close();

Reply via email to