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

yuyuankang pushed a commit to branch IoTDB-1499
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 0681eb96587687c8f6163443125e87be49e77c9c
Author: Ring-k <[email protected]>
AuthorDate: Thu Jul 15 14:09:41 2021 +0800

    remove series registeration
---
 .../UserGuide/Ecosystem Integration/Flink IoTDB.md |  3 ++-
 .../org/apache/iotdb/flink/FlinkIoTDBSink.java     | 26 +++++++++++++---------
 .../java/org/apache/iotdb/flink/IoTDBSink.java     | 26 +---------------------
 .../iotdb/flink/options/IoTDBSinkOptions.java      | 11 ---------
 4 files changed, 19 insertions(+), 47 deletions(-)

diff --git a/docs/UserGuide/Ecosystem Integration/Flink IoTDB.md 
b/docs/UserGuide/Ecosystem Integration/Flink IoTDB.md
index 41cfbce..5a5520a 100644
--- a/docs/UserGuide/Ecosystem Integration/Flink IoTDB.md       
+++ b/docs/UserGuide/Ecosystem Integration/Flink IoTDB.md       
@@ -35,6 +35,8 @@ This example shows a case that sends data to a IoTDB server 
from a Flink job:
 - A simulated Source `SensorSource` generates data points per 1 second.
 - Flink uses `IoTDBSink` to consume the generated data points and write the 
data into IoTDB.
 
+It is noteworthy that to use IoTDBSink, schema auto-creation in IoTDB should 
be enabled. 
+
 ```java
 import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
@@ -59,7 +61,6 @@ public class FlinkIoTDBSink {
     options.setPort(6667);
     options.setUser("root");
     options.setPassword("root");
-    options.setStorageGroup("root.sg");
 
     // If the server enables auto_create_schema, then we do not need to 
register all timeseries
     // here.
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkIoTDBSink.java 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkIoTDBSink.java
index 5620949..9b72438 100644
--- a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkIoTDBSink.java
+++ b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkIoTDBSink.java
@@ -36,19 +36,25 @@ public class FlinkIoTDBSink {
     // run the flink job on local mini cluster
     StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
 
-    IoTDBSinkOptions options = new IoTDBSinkOptions();
-    options.setHost("127.0.0.1");
-    options.setPort(6667);
-    options.setUser("root");
-    options.setPassword("root");
-    options.setStorageGroup("root.sg");
+    String host = "127.0.0.1";
+    int port = 6667;
+    String user = "root";
+    String password = "root";
 
     // If the server enables auto_create_schema, then we do not need to 
register all timeseries
     // here.
-    options.setTimeseriesOptionList(
-        Lists.newArrayList(
-            new IoTDBSinkOptions.TimeseriesOption(
-                "root.sg.d1.s1", TSDataType.DOUBLE, TSEncoding.GORILLA, 
CompressionType.SNAPPY)));
+    IoTDBSinkOptions options =
+        new IoTDBSinkOptions(
+            host,
+            port,
+            user,
+            password,
+            Lists.newArrayList(
+                new IoTDBSinkOptions.TimeseriesOption(
+                    "root.sg.d1.s1",
+                    TSDataType.DOUBLE,
+                    TSEncoding.GORILLA,
+                    CompressionType.SNAPPY)));
 
     IoTSerializationSchema serializationSchema = new 
DefaultIoTSerializationSchema();
     IoTDBSink ioTDBSink =
diff --git 
a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java 
b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java
index 95e40a2..8ab06c0 100644
--- a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java
+++ b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/IoTDBSink.java
@@ -19,8 +19,6 @@
 package org.apache.iotdb.flink;
 
 import org.apache.iotdb.flink.options.IoTDBSinkOptions;
-import org.apache.iotdb.rpc.StatementExecutionException;
-import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.session.pool.SessionPool;
 import org.apache.iotdb.tsfile.common.constant.TsFileConstant;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
@@ -78,7 +76,7 @@ public class IoTDBSink<IN> extends RichSinkFunction<IN> {
     initScheduler();
   }
 
-  void initSession() throws Exception {
+  void initSession() {
     pool =
         new SessionPool(
             options.getHost(),
@@ -86,28 +84,6 @@ public class IoTDBSink<IN> extends RichSinkFunction<IN> {
             options.getUser(),
             options.getPassword(),
             sessionPoolSize);
-
-    try {
-      pool.setStorageGroup(options.getStorageGroup());
-    } catch (StatementExecutionException e) {
-      if (e.getStatusCode() != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-        throw e;
-      }
-    }
-
-    for (IoTDBSinkOptions.TimeseriesOption option : 
options.getTimeseriesOptionList()) {
-      if (!pool.checkTimeseriesExists(option.getPath())) {
-        try {
-          pool.createTimeseries(
-              option.getPath(), option.getDataType(), option.getEncoding(), 
option.getCompressor());
-        } catch (StatementExecutionException e) {
-          // path could have been created by the other process here
-          if (e.getStatusCode() != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-            throw e;
-          }
-        }
-      }
-    }
   }
 
   void initScheduler() {
diff --git 
a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/options/IoTDBSinkOptions.java
 
b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/options/IoTDBSinkOptions.java
index 60bab60..1a709b0 100644
--- 
a/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/options/IoTDBSinkOptions.java
+++ 
b/flink-iotdb-connector/src/main/java/org/apache/iotdb/flink/options/IoTDBSinkOptions.java
@@ -28,7 +28,6 @@ import java.util.List;
 /** IoTDBOptions describes the configuration related information for IoTDB and 
timeseries. */
 public class IoTDBSinkOptions extends IoTDBOptions {
 
-  private String storageGroup;
   private List<TimeseriesOption> timeseriesOptionList;
 
   public IoTDBSinkOptions() {}
@@ -38,21 +37,11 @@ public class IoTDBSinkOptions extends IoTDBOptions {
       int port,
       String user,
       String password,
-      String storageGroup,
       List<TimeseriesOption> timeseriesOptionList) {
     super(host, port, user, password);
-    this.storageGroup = storageGroup;
     this.timeseriesOptionList = timeseriesOptionList;
   }
 
-  public String getStorageGroup() {
-    return storageGroup;
-  }
-
-  public void setStorageGroup(String storageGroup) {
-    this.storageGroup = storageGroup;
-  }
-
   public List<TimeseriesOption> getTimeseriesOptionList() {
     return timeseriesOptionList;
   }

Reply via email to