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