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 b4f698d  fix a class name error in flink (#1202)
b4f698d is described below

commit b4f698ddd3f1d94e0badc17a43b4607d7580fad0
Author: Haonan <[email protected]>
AuthorDate: Wed May 13 13:55:51 2020 +0800

    fix a class name error in flink (#1202)
    
    * fix TsFileUtils name error in flink
---
 .../apache/iotdb/flink/FlinkTsFileBatchSink.java   |   2 +-
 .../apache/iotdb/flink/FlinkTsFileBatchSource.java |   2 +-
 .../apache/iotdb/flink/FlinkTsFileStreamSink.java  |   2 +-
 .../iotdb/flink/FlinkTsFileStreamSource.java       |   2 +-
 .../java/org/apache/iotdb/flink/TsFileUtils.java   | 103 +++++++++++++++++++++
 .../java/org/apache/iotdb/flink/TsFlieUtils.java   |  98 --------------------
 6 files changed, 107 insertions(+), 102 deletions(-)

diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSink.java 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSink.java
index 2394824..f1174c9 100644
--- 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSink.java
+++ 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSink.java
@@ -106,7 +106,7 @@ public class FlinkTsFileBatchSink {
                        .filter(s -> !s.equals(QueryConstant.RESERVED_TIME))
                        .map(Path::new)
                        .collect(Collectors.toList());
-               String[] result = TsFlieUtils.readTsFile(path, paths);
+               String[] result = TsFileUtils.readTsFile(path, paths);
                for (String row : result) {
                        System.out.println(row);
                }
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSource.java
 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSource.java
index 6f01b75..54d31a4 100644
--- 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSource.java
+++ 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileBatchSource.java
@@ -41,7 +41,7 @@ public class FlinkTsFileBatchSource {
 
        public static void main(String[] args) throws Exception {
                String path = "test.tsfile";
-               TsFlieUtils.writeTsFile(path);
+               TsFileUtils.writeTsFile(path);
                new File(path).deleteOnExit();
                String[] filedNames = {
                        QueryConstant.RESERVED_TIME,
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSink.java 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSink.java
index cdf1390..0a5759c 100644
--- 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSink.java
+++ 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSink.java
@@ -107,7 +107,7 @@ public class FlinkTsFileStreamSink {
                        .filter(s -> !s.equals(QueryConstant.RESERVED_TIME))
                        .map(Path::new)
                        .collect(Collectors.toList());
-               String[] result = TsFlieUtils.readTsFile(path, paths);
+               String[] result = TsFileUtils.readTsFile(path, paths);
                for (String row : result) {
                        System.out.println(row);
                }
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSource.java
 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSource.java
index 83f7de9..d17d3c5 100644
--- 
a/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSource.java
+++ 
b/example/flink/src/main/java/org/apache/iotdb/flink/FlinkTsFileStreamSource.java
@@ -44,7 +44,7 @@ public class FlinkTsFileStreamSource {
 
        public static void main(String[] args) throws IOException {
                String path = "test.tsfile";
-               TsFlieUtils.writeTsFile(path);
+               TsFileUtils.writeTsFile(path);
                new File(path).deleteOnExit();
                String[] filedNames = {
                        QueryConstant.RESERVED_TIME,
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/TsFileUtils.java 
b/example/flink/src/main/java/org/apache/iotdb/flink/TsFileUtils.java
new file mode 100644
index 0000000..7da0eab
--- /dev/null
+++ b/example/flink/src/main/java/org/apache/iotdb/flink/TsFileUtils.java
@@ -0,0 +1,103 @@
+/*
+ * 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.flink;
+
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
+import org.apache.iotdb.tsfile.read.ReadOnlyTsFile;
+import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
+import org.apache.iotdb.tsfile.read.common.Path;
+import org.apache.iotdb.tsfile.read.common.RowRecord;
+import org.apache.iotdb.tsfile.read.expression.QueryExpression;
+import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
+import org.apache.iotdb.tsfile.write.TsFileWriter;
+import org.apache.iotdb.tsfile.write.record.TSRecord;
+import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
+import org.apache.iotdb.tsfile.write.record.datapoint.LongDataPoint;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
+import org.apache.iotdb.tsfile.write.schema.Schema;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.File;
+import java.io.IOException;
+import java.nio.file.Files;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.stream.Collectors;
+
+/**
+ * Utils used to prepare source TsFiles for the examples.
+ */
+public class TsFileUtils {
+
+  private static final Logger logger = 
LoggerFactory.getLogger(TsFileUtils.class);
+  private static final String DEFAULT_TEMPLATE = "template";
+
+  private TsFileUtils() {
+  }
+
+  public static void writeTsFile(String path) {
+    try {
+      File f = FSFactoryProducer.getFSFactory().getFile(path);
+                       Files.delete(f.toPath());
+      Schema schema = new Schema();
+      schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_1", TSDataType.FLOAT, TSEncoding.RLE));
+      schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_2", TSDataType.INT32, TSEncoding.TS_2DIFF));
+      schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_3", TSDataType.INT32, TSEncoding.TS_2DIFF));
+
+      try (TsFileWriter tsFileWriter = new TsFileWriter(f, schema)) {
+
+        // construct TSRecord
+        for (int i = 0; i < 100; i++) {
+          TSRecord tsRecord = new TSRecord(i, "device_" + (i % 4));
+          DataPoint dPoint1 = new LongDataPoint("sensor_1", i);
+          DataPoint dPoint2 = new LongDataPoint("sensor_2", i);
+          DataPoint dPoint3 = new LongDataPoint("sensor_3", i);
+          tsRecord.addTuple(dPoint1);
+          tsRecord.addTuple(dPoint2);
+          tsRecord.addTuple(dPoint3);
+          
+          // write TSRecord
+          tsFileWriter.write(tsRecord);
+        }
+      }
+
+    } catch (Exception e) {
+      logger.error("Write {} failed. ", path, e);
+    }
+  }
+
+  public static String[] readTsFile(String tsFilePath, List<Path> paths) 
throws IOException {
+    QueryExpression expression = QueryExpression.create(paths, null);
+    TsFileSequenceReader reader = new TsFileSequenceReader(tsFilePath);
+    try (ReadOnlyTsFile readTsFile = new ReadOnlyTsFile(reader)) {
+      QueryDataSet queryDataSet = readTsFile.query(expression);
+      List<String> result = new ArrayList<>();
+      while (queryDataSet.hasNext()) {
+        RowRecord rowRecord = queryDataSet.next();
+        String row = rowRecord.getFields().stream()
+            .map(f -> f == null ? "null" : f.getStringValue())
+            .collect(Collectors.joining(","));
+        result.add(rowRecord.getTimestamp() + "," + row);
+      }
+      return result.toArray(new String[0]);
+    }
+  }
+}
diff --git 
a/example/flink/src/main/java/org/apache/iotdb/flink/TsFlieUtils.java 
b/example/flink/src/main/java/org/apache/iotdb/flink/TsFlieUtils.java
deleted file mode 100644
index f97b683..0000000
--- a/example/flink/src/main/java/org/apache/iotdb/flink/TsFlieUtils.java
+++ /dev/null
@@ -1,98 +0,0 @@
-/*
- * 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.flink;
-
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
-import org.apache.iotdb.tsfile.fileSystem.FSFactoryProducer;
-import org.apache.iotdb.tsfile.read.ReadOnlyTsFile;
-import org.apache.iotdb.tsfile.read.TsFileSequenceReader;
-import org.apache.iotdb.tsfile.read.common.Path;
-import org.apache.iotdb.tsfile.read.common.RowRecord;
-import org.apache.iotdb.tsfile.read.expression.QueryExpression;
-import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
-import org.apache.iotdb.tsfile.write.TsFileWriter;
-import org.apache.iotdb.tsfile.write.record.TSRecord;
-import org.apache.iotdb.tsfile.write.record.datapoint.DataPoint;
-import org.apache.iotdb.tsfile.write.record.datapoint.LongDataPoint;
-import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
-import org.apache.iotdb.tsfile.write.schema.Schema;
-
-import java.io.File;
-import java.io.IOException;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.stream.Collectors;
-
-/**
- * Utils used to prepare source TsFiles for the examples.
- */
-public class TsFlieUtils {
-
-       private static final String DEFAULT_TEMPLATE = "template";
-
-       public static void writeTsFile(String path) {
-               try {
-                       File f = FSFactoryProducer.getFSFactory().getFile(path);
-                       if (f.exists()) {
-                               f.delete();
-                       }
-                       Schema schema = new Schema();
-                       schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_1", TSDataType.FLOAT, TSEncoding.RLE));
-                       schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_2", TSDataType.INT32, TSEncoding.TS_2DIFF));
-                       schema.extendTemplate(DEFAULT_TEMPLATE, new 
MeasurementSchema("sensor_3", TSDataType.INT32, TSEncoding.TS_2DIFF));
-
-                       TsFileWriter tsFileWriter = new TsFileWriter(f, schema);
-
-                       // construct TSRecord
-                       for (int i = 0; i < 100; i++) {
-                               TSRecord tsRecord = new TSRecord(i, "device_" + 
(i % 4));
-                               DataPoint dPoint1 = new 
LongDataPoint("sensor_1", i);
-                               DataPoint dPoint2 = new 
LongDataPoint("sensor_2", i);
-                               DataPoint dPoint3 = new 
LongDataPoint("sensor_3", i);
-                               tsRecord.addTuple(dPoint1);
-                               tsRecord.addTuple(dPoint2);
-                               tsRecord.addTuple(dPoint3);
-
-                               // write TSRecord
-                               tsFileWriter.write(tsRecord);
-                       }
-
-                       tsFileWriter.close();
-               } catch (Throwable e) {
-                       e.printStackTrace();
-                       System.out.println(e.getMessage());
-               }
-       }
-
-       public static String[] readTsFile(String tsFilePath, List<Path> paths) 
throws IOException {
-               QueryExpression expression = QueryExpression.create(paths, 
null);
-               TsFileSequenceReader reader = new 
TsFileSequenceReader(tsFilePath);
-               ReadOnlyTsFile readTsFile = new ReadOnlyTsFile(reader);
-               QueryDataSet queryDataSet = readTsFile.query(expression);
-               List<String> result = new ArrayList<>();
-               while (queryDataSet.hasNext()) {
-                       RowRecord rowRecord = queryDataSet.next();
-                       String row = rowRecord.getFields().stream()
-                               .map(f -> f == null ? "null" : 
f.getStringValue())
-                               .collect(Collectors.joining(","));
-                       result.add(rowRecord.getTimestamp() + "," + row);
-               }
-               return result.toArray(new String[0]);
-       }
-}

Reply via email to