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; }
