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