From 1240099c815b7a673f782904b8ad64c36c253018 Mon Sep 17 00:00:00 2001 From: Andrew Volostnykh Date: Fri, 12 Mar 2021 18:39:17 +0200 Subject: [PATCH] Full refactoring and code cleaning --- tools/pom.xml | 2 + .../tools/migrator/DictionaryParser.java | 14 +- .../client/tools/migrator/MigratorTool.java | 15 +- .../tools/migrator/PgCaLatestMigrator.java | 183 ------------- .../client/tools/migrator/PgCaMigrator.java | 254 ++++++++++++++++++ .../PostgresToCassandraTelemetryMigrator.java | 224 --------------- .../tools/migrator/RelatedEntitiesParser.java | 68 ++--- 7 files changed, 309 insertions(+), 451 deletions(-) delete mode 100644 tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaLatestMigrator.java create mode 100644 tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java delete mode 100644 tools/src/main/java/org/thingsboard/client/tools/migrator/PostgresToCassandraTelemetryMigrator.java diff --git a/tools/pom.xml b/tools/pom.xml index 3b163b44b3..e3df0e08d0 100644 --- a/tools/pom.xml +++ b/tools/pom.xml @@ -54,6 +54,7 @@ org.apache.cassandra cassandra-all + 3.11.10 com.datastax.cassandra @@ -63,6 +64,7 @@ commons-io commons-io + 2.5 diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java index 276daf003e..53caeaeb0d 100644 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/DictionaryParser.java @@ -44,13 +44,17 @@ public class DictionaryParser { } private void parseDictionaryDump(LineIterator iterator) { - String tempLine; - while(iterator.hasNext()) { - tempLine = iterator.nextLine(); + try { + String tempLine; + while (iterator.hasNext()) { + tempLine = iterator.nextLine(); - if(isBlockStarted(tempLine)) { - processBlock(iterator); + if (isBlockStarted(tempLine)) { + processBlock(iterator); + } } + } finally { + iterator.close(); } } diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java index 83ca54624a..526d59f936 100644 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/MigratorTool.java @@ -33,23 +33,24 @@ public class MigratorTool { try { boolean castEnable = Boolean.parseBoolean(cmd.getOptionValue("castEnable")); File allTelemetrySource = new File(cmd.getOptionValue("telemetryFrom")); + File tsSaveDir = null; + File partitionsSaveDir = null; + File latestSaveDir = null; RelatedEntitiesParser allEntityIdsAndTypes = new RelatedEntitiesParser(new File(cmd.getOptionValue("relatedEntities"))); DictionaryParser dictionaryParser = new DictionaryParser(allTelemetrySource); if(cmd.getOptionValue("latestTelemetryOut") != null) { - File latestSaveDir = new File(cmd.getOptionValue("latestTelemetryOut")); - PgCaLatestMigrator.migrateLatest(allTelemetrySource, latestSaveDir, allEntityIdsAndTypes, dictionaryParser, castEnable); + latestSaveDir = new File(cmd.getOptionValue("latestTelemetryOut")); } if(cmd.getOptionValue("telemetryOut") != null) { - File tsSaveDir = new File(cmd.getOptionValue("telemetryOut")); - File partitionsSaveDir = new File(cmd.getOptionValue("partitionsOut")); - PostgresToCassandraTelemetryMigrator.migrateTs( - allTelemetrySource, tsSaveDir, partitionsSaveDir, allEntityIdsAndTypes, dictionaryParser, castEnable - ); + tsSaveDir = new File(cmd.getOptionValue("telemetryOut")); + partitionsSaveDir = new File(cmd.getOptionValue("partitionsOut")); } + new PgCaMigrator(allTelemetrySource, tsSaveDir, partitionsSaveDir, latestSaveDir, allEntityIdsAndTypes, dictionaryParser, castEnable).migrate(); + } catch (Throwable th) { th.printStackTrace(); throw new IllegalStateException("failed", th); diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaLatestMigrator.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaLatestMigrator.java deleted file mode 100644 index 5d4a2a6e81..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaLatestMigrator.java +++ /dev/null @@ -1,183 +0,0 @@ -/** - * Copyright © 2016-2021 The Thingsboard Authors - * - * Licensed 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.thingsboard.client.tools.migrator; - -import com.google.common.collect.Lists; -import org.apache.cassandra.io.sstable.CQLSSTableWriter; -import org.apache.commons.io.FileUtils; -import org.apache.commons.io.LineIterator; -import org.apache.commons.lang3.StringUtils; -import org.apache.commons.lang3.math.NumberUtils; - -import java.io.File; -import java.io.IOException; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Date; -import java.util.List; -import java.util.UUID; -import java.util.stream.Collectors; - -public class PgCaLatestMigrator { - - private static final long LOG_BATCH = 1000000; - private static final long rowPerFile = 1000000; - - - private static long linesProcessed = 0; - private static long linesMigrated = 0; - private static long castErrors = 0; - private static long castedOk = 0; - - private static long currentWriterCount = 1; - private static RelatedEntitiesParser allIdsAndTypes; - private static DictionaryParser keyPairs; - - public static void migrateLatest(File sourceFile, - File outDir, - RelatedEntitiesParser allEntityIdsAndTypes, - DictionaryParser dictionaryParser, - boolean castStringsIfPossible) throws IOException { - long startTs = System.currentTimeMillis(); - long stepLineTs = System.currentTimeMillis(); - long stepOkLineTs = System.currentTimeMillis(); - LineIterator iterator = FileUtils.lineIterator(sourceFile); - CQLSSTableWriter currentTsWriter = WriterBuilder.getLatestWriter(outDir); - allIdsAndTypes = allEntityIdsAndTypes; - keyPairs = dictionaryParser; - - boolean isBlockStarted = false; - boolean isBlockFinished = false; - - String line; - while (iterator.hasNext()) { - if (linesProcessed++ % LOG_BATCH == 0) { - System.out.println(new Date() + " linesProcessed = " + linesProcessed + " in " + (System.currentTimeMillis() - stepLineTs) + " castOk " + castedOk + " castErr " + castErrors); - stepLineTs = System.currentTimeMillis(); - } - - line = iterator.nextLine(); - - if (isBlockFinished) { - break; - } - - if (!isBlockStarted) { - if (isBlockStarted(line)) { - System.out.println(); - System.out.println(); - System.out.println(line); - System.out.println(); - System.out.println(); - isBlockStarted = true; - } - continue; - } - - if (isBlockFinished(line)) { - isBlockFinished = true; - } else { - try { - List raw = Arrays.stream(line.trim().split("\t")) - .map(String::trim) - .filter(StringUtils::isNotEmpty) - .collect(Collectors.toList()); - List values = toValues(raw); - - if (currentWriterCount == 0) { - System.out.println(new Date() + " close writer " + new Date()); - currentTsWriter.close(); - currentTsWriter = WriterBuilder.getLatestWriter(outDir); - } - - if (castStringsIfPossible) { - currentTsWriter.addRow(castToNumericIfPossible(values)); - } else { - currentTsWriter.addRow(values); - } - currentWriterCount++; - if (currentWriterCount >= rowPerFile) { - currentWriterCount = 0; - } - - if (linesMigrated++ % LOG_BATCH == 0) { - System.out.println(new Date() + " migrated = " + linesMigrated + " in " + (System.currentTimeMillis() - stepOkLineTs) + " ms."); - stepOkLineTs = System.currentTimeMillis(); - } - } catch (Exception ex) { - System.out.println(ex.getMessage() + " -> " + line); - } - - } - } - - long endTs = System.currentTimeMillis(); - System.out.println(); - System.out.println(new Date() + " Migrated rows " + linesMigrated + " in " + (endTs - startTs) + " ts"); - - currentTsWriter.close(); - System.out.println(); - System.out.println("Finished migrate Latest Telemetry"); - } - - - private static List castToNumericIfPossible(List values) { - try { - if (values.get(6) != null && NumberUtils.isNumber(values.get(6).toString())) { - Double casted = NumberUtils.createDouble(values.get(6).toString()); - List numeric = Lists.newArrayList(); - numeric.addAll(values); - numeric.set(6, null); - numeric.set(8, casted); - castedOk++; - return numeric; - } - } catch (Throwable th) { - castErrors++; - } - return values; - } - - private static List toValues(List raw) { - //expected Table structure: - //COPY public.ts_kv_latest (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) FROM stdin; - - List result = new ArrayList<>(); - result.add(allIdsAndTypes.getEntityType(raw.get(0))); - result.add(UUID.fromString(raw.get(0))); - result.add(keyPairs.getKeyByKeyId(raw.get(1))); - - long ts = Long.parseLong(raw.get(2)); - result.add(3, ts); - - result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE); - result.add(raw.get(4).equals("\\N") ? null : raw.get(4)); - result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5))); - result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6))); - result.add(raw.get(7).equals("\\N") ? null : raw.get(7)); - - return result; - } - - private static boolean isBlockStarted(String line) { - return line.startsWith("COPY public.ts_kv_latest ("); - } - - private static boolean isBlockFinished(String line) { - return StringUtils.isBlank(line) || line.equals("\\."); - } - -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java new file mode 100644 index 0000000000..afa341f880 --- /dev/null +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/PgCaMigrator.java @@ -0,0 +1,254 @@ +package org.thingsboard.client.tools.migrator; + +import com.google.common.collect.Lists; +import org.apache.cassandra.io.sstable.CQLSSTableWriter; +import org.apache.commons.io.FileUtils; +import org.apache.commons.io.LineIterator; +import org.apache.commons.lang3.StringUtils; +import org.apache.commons.lang3.math.NumberUtils; + +import java.io.File; +import java.io.IOException; +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneOffset; +import java.time.temporal.ChronoUnit; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Date; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.UUID; +import java.util.function.Function; +import java.util.stream.Collectors; + +public class PgCaMigrator { + + private final long LOG_BATCH = 1000000; + private final long rowPerFile = 1000000; + + private long linesProcessed = 0; + private long linesTsMigrated = 0; + private long linesLatestMigrated = 0; + private long castErrors = 0; + private long castedOk = 0; + + private long currentWriterCount = 1; + + private final File sourceFile; + private final boolean castStringIfPossible; + + private final RelatedEntitiesParser entityIdsAndTypes; + private final DictionaryParser keyParser; + private CQLSSTableWriter currentTsWriter; + private CQLSSTableWriter currentPartitionsWriter; + private CQLSSTableWriter currentTsLatestWriter; + private final Set partitions = new HashSet<>(); + + private File outTsDir; + private File outTsLatestDir; + + public PgCaMigrator(File sourceFile, + File ourTsDir, + File outTsPartitionDir, + File outTsLatestDir, + RelatedEntitiesParser allEntityIdsAndTypes, + DictionaryParser dictionaryParser, + boolean castStringsIfPossible) { + this.sourceFile = sourceFile; + this.entityIdsAndTypes = allEntityIdsAndTypes; + this.keyParser = dictionaryParser; + this.castStringIfPossible = castStringsIfPossible; + if(outTsLatestDir != null) { + this.currentTsLatestWriter = WriterBuilder.getLatestWriter(outTsLatestDir); + this.outTsLatestDir = outTsLatestDir; + } + if(ourTsDir != null) { + this.currentTsWriter = WriterBuilder.getTsWriter(ourTsDir); + this.currentPartitionsWriter = WriterBuilder.getPartitionWriter(outTsPartitionDir); + this.outTsDir = ourTsDir; + } + } + + public void migrate() throws IOException { + boolean isTsDone = false; + boolean isLatestDone = false; + String line; + LineIterator iterator = FileUtils.lineIterator(this.sourceFile); + + try { + while(iterator.hasNext()) { + line = iterator.nextLine(); + if(!isLatestDone && isBlockLatestStarted(line)) { + System.out.println("START TO MIGRATE LATEST"); + long start = System.currentTimeMillis(); + processBlock(iterator, currentTsLatestWriter, outTsLatestDir, this::toValuesLatest); + System.out.println("FORMING OF SSL FOR LATEST TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms."); + isLatestDone = true; + } + + if(!isTsDone && isBlockTsStarted(line)) { + System.out.println("START TO MIGRATE TS"); + long start = System.currentTimeMillis(); + processBlock(iterator, currentTsWriter, outTsDir, this::toValuesTs); + System.out.println("FORMING OF SSL FOR TS FINISHED WITH TIME: " + (System.currentTimeMillis() - start) + " ms."); + isTsDone = true; + } + } + + System.out.println("Partitions collected " + partitions.size()); + long startTs = System.currentTimeMillis(); + for (String partition : partitions) { + String[] split = partition.split("\\|"); + List values = Lists.newArrayList(); + values.add(split[0]); + values.add(UUID.fromString(split[1])); + values.add(split[2]); + values.add(Long.parseLong(split[3])); + currentPartitionsWriter.addRow(values); + } + + System.out.println(new Date() + " Migrated partitions " + partitions.size() + " in " + (System.currentTimeMillis() - startTs)); + + System.out.println(); + System.out.println("Finished migrate Telemetry"); + + } finally { + iterator.close(); + currentTsLatestWriter.close(); + currentTsWriter.close(); + currentPartitionsWriter.close(); + } + } + + private void logLinesProcessed() { + if (linesProcessed++ % LOG_BATCH == 0) { + System.out.println(new Date() + " linesProcessed = " + linesProcessed + " in, castOk " + castedOk + " castErr " + castErrors); + } + } + + private List toValuesTs(List raw) { + linesTsMigrated++; + List result = new ArrayList<>(); + result.add(entityIdsAndTypes.getEntityType(raw.get(0))); + result.add(UUID.fromString(raw.get(0))); + result.add(keyParser.getKeyByKeyId(raw.get(1))); + + long ts = Long.parseLong(raw.get(2)); + long partition = toPartitionTs(ts); + result.add(partition); + result.add(ts); + + result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE); + result.add(raw.get(4).equals("\\N") ? null : raw.get(4)); + result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5))); + result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6))); + result.add(raw.get(7).equals("\\N") ? null : raw.get(7)); + + processPartitions(result); + + return result; + } + + private List toValuesLatest(List raw) { + linesLatestMigrated++; + List result = new ArrayList<>(); + result.add(this.entityIdsAndTypes.getEntityType(raw.get(0))); + result.add(UUID.fromString(raw.get(0))); + result.add(this.keyParser.getKeyByKeyId(raw.get(1))); + + long ts = Long.parseLong(raw.get(2)); + result.add(3, ts); + + result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE); + result.add(raw.get(4).equals("\\N") ? null : raw.get(4)); + result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5))); + result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6))); + result.add(raw.get(7).equals("\\N") ? null : raw.get(7)); + + return result; + } + + private long toPartitionTs(long ts) { + LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); + return time.truncatedTo(ChronoUnit.DAYS).withDayOfMonth(1).toInstant(ZoneOffset.UTC).toEpochMilli(); + } + + private void processPartitions(List values) { + String key = values.get(0) + "|" + values.get(1) + "|" + values.get(2) + "|" + values.get(3); + partitions.add(key); + } + + private void processBlock(LineIterator iterator, CQLSSTableWriter writer, File outDir, Function, List> function) { + String currentLine; + linesProcessed = 0; + while(iterator.hasNext()) { + logLinesProcessed(); + currentLine = iterator.nextLine(); + if(isBlockFinished(currentLine)) { + return; + } + + try { + List raw = Arrays.stream(currentLine.trim().split("\t")) + .map(String::trim) + .filter(StringUtils::isNotEmpty) + .collect(Collectors.toList()); + List values = function.apply(raw); + + if (this.currentWriterCount == 0) { + System.out.println(new Date() + " close writer " + new Date()); + writer.close(); + writer = WriterBuilder.getLatestWriter(outDir); + } + + if (this.castStringIfPossible) { + writer.addRow(castToNumericIfPossible(values)); + } else { + writer.addRow(values); + } + + currentWriterCount++; + if (currentWriterCount >= rowPerFile) { + currentWriterCount = 0; + } + } catch (Exception ex) { + System.out.println(ex.getMessage() + " -> " + currentLine); + } + } + } + + private List castToNumericIfPossible(List values) { + try { + if (values.get(6) != null && NumberUtils.isNumber(values.get(6).toString())) { + Double casted = NumberUtils.createDouble(values.get(6).toString()); + List numeric = Lists.newArrayList(); + numeric.addAll(values); + numeric.set(6, null); + numeric.set(8, casted); + castedOk++; + return numeric; + } + } catch (Throwable th) { + castErrors++; + } + + processPartitions(values); + + return values; + } + + private boolean isBlockFinished(String line) { + return StringUtils.isBlank(line) || line.equals("\\."); + } + + private boolean isBlockTsStarted(String line) { + return line.startsWith("COPY public.ts_kv ("); + } + + private boolean isBlockLatestStarted(String line) { + return line.startsWith("COPY public.ts_kv_latest ("); + } + +} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/PostgresToCassandraTelemetryMigrator.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/PostgresToCassandraTelemetryMigrator.java deleted file mode 100644 index 80e617f505..0000000000 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/PostgresToCassandraTelemetryMigrator.java +++ /dev/null @@ -1,224 +0,0 @@ -/** - * Copyright © 2016-2021 The Thingsboard Authors - * - * Licensed 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.thingsboard.client.tools.migrator; - -import com.google.common.collect.Lists; -import org.apache.cassandra.io.sstable.CQLSSTableWriter; -import org.apache.commons.io.FileUtils; -import org.apache.commons.io.LineIterator; -import org.apache.commons.lang3.StringUtils; -import org.apache.commons.lang3.math.NumberUtils; - -import java.io.File; -import java.io.IOException; -import java.time.Instant; -import java.time.LocalDateTime; -import java.time.ZoneOffset; -import java.time.temporal.ChronoUnit; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.Date; -import java.util.HashSet; -import java.util.List; -import java.util.Set; -import java.util.UUID; -import java.util.stream.Collectors; - -public class PostgresToCassandraTelemetryMigrator { - - private static final long LOG_BATCH = 1000000; - private static final long rowPerFile = 1000000; - - private static long linesProcessed = 0; - private static long linesMigrated = 0; - private static long castErrors = 0; - private static long castedOk = 0; - - private static long currentWriterCount = 1; - private static CQLSSTableWriter currentTsWriter = null; - private static CQLSSTableWriter currentPartitionWriter = null; - - private static Set partitions = new HashSet<>(); - private static RelatedEntitiesParser entityIdsAndTypes; - private static DictionaryParser keyParser; - - public static void migrateTs(File sourceFile, - File outTsDir, - File outPartitionDir, - RelatedEntitiesParser allEntityIdsAndTypes, - DictionaryParser dictionaryParser, - boolean castStringsIfPossible) throws IOException { - long startTs = System.currentTimeMillis(); - long stepLineTs = System.currentTimeMillis(); - long stepOkLineTs = System.currentTimeMillis(); - LineIterator iterator = FileUtils.lineIterator(sourceFile); - currentTsWriter = WriterBuilder.getTsWriter(outTsDir); - currentPartitionWriter = WriterBuilder.getPartitionWriter(outPartitionDir); - entityIdsAndTypes = allEntityIdsAndTypes; - keyParser = dictionaryParser; - - boolean isBlockStarted = false; - boolean isBlockFinished = false; - - String line; - while (iterator.hasNext()) { - if (linesProcessed++ % LOG_BATCH == 0) { - System.out.println(new Date() + " linesProcessed = " + linesProcessed + " in " + (System.currentTimeMillis() - stepLineTs) + " castOk " + castedOk + " castErr " + castErrors); - stepLineTs = System.currentTimeMillis(); - } - - line = iterator.nextLine(); - - if (isBlockFinished) { - break; - } - - if (!isBlockStarted) { - if (isBlockStarted(line)) { - System.out.println(); - System.out.println(); - System.out.println(line); - System.out.println(); - System.out.println(); - isBlockStarted = true; - } - continue; - } - - if (isBlockFinished(line)) { - isBlockFinished = true; - } else { - try { - List raw = Arrays.stream(line.trim().split("\t")) - .map(String::trim) - .filter(StringUtils::isNotEmpty) - .collect(Collectors.toList()); - List values = toValues(raw); - - if (currentWriterCount == 0) { - System.out.println(new Date() + " close writer " + new Date()); - currentTsWriter.close(); - currentTsWriter = WriterBuilder.getTsWriter(outTsDir); - } - - if (castStringsIfPossible) { - currentTsWriter.addRow(castToNumericIfPossible(values)); - } else { - currentTsWriter.addRow(values); - } - processPartitions(values); - currentWriterCount++; - if (currentWriterCount >= rowPerFile) { - currentWriterCount = 0; - } - - if (linesMigrated++ % LOG_BATCH == 0) { - System.out.println(new Date() + " migrated = " + linesMigrated + " in " + (System.currentTimeMillis() - stepOkLineTs) + " partitions = " + partitions.size()); - stepOkLineTs = System.currentTimeMillis(); - } - } catch (Exception ex) { - System.out.println(ex.getMessage() + " -> " + line); - } - - } - } - - long endTs = System.currentTimeMillis(); - System.out.println(); - System.out.println(new Date() + " Migrated rows " + linesMigrated + " in " + (endTs - startTs)); - System.out.println("Partitions collected " + partitions.size()); - - startTs = System.currentTimeMillis(); - for (String partition : partitions) { - String[] split = partition.split("\\|"); - List values = Lists.newArrayList(); - values.add(split[0]); - values.add(UUID.fromString(split[1])); - values.add(split[2]); - values.add(Long.parseLong(split[3])); - currentPartitionWriter.addRow(values); - } - currentPartitionWriter.close(); - endTs = System.currentTimeMillis(); - System.out.println(); - System.out.println(); - System.out.println(new Date() + " Migrated partitions " + partitions.size() + " in " + (endTs - startTs)); - - - currentTsWriter.close(); - System.out.println(); - System.out.println("Finished migrate Telemetry"); - } - - private static List castToNumericIfPossible(List values) { - try { - if (values.get(6) != null && NumberUtils.isNumber(values.get(6).toString())) { - Double casted = NumberUtils.createDouble(values.get(6).toString()); - List numeric = Lists.newArrayList(); - numeric.addAll(values); - numeric.set(6, null); - numeric.set(8, casted); - castedOk++; - return numeric; - } - } catch (Throwable th) { - castErrors++; - } - return values; - } - - private static void processPartitions(List values) { - String key = values.get(0) + "|" + values.get(1) + "|" + values.get(2) + "|" + values.get(3); - partitions.add(key); - } - - private static List toValues(List raw) { - //expected Table structure: -// COPY public.ts_kv (entity_type, entity_id, key, ts, bool_v, str_v, long_v, dbl_v) FROM stdin; - - List result = new ArrayList<>(); - result.add(entityIdsAndTypes.getEntityType(raw.get(0))); - result.add(UUID.fromString(raw.get(0))); - result.add(keyParser.getKeyByKeyId(raw.get(1))); - - long ts = Long.parseLong(raw.get(2)); - long partition = toPartitionTs(ts); - result.add(partition); - result.add(ts); - - result.add(raw.get(3).equals("\\N") ? null : raw.get(3).equals("t") ? Boolean.TRUE : Boolean.FALSE); - result.add(raw.get(4).equals("\\N") ? null : raw.get(4)); - result.add(raw.get(5).equals("\\N") ? null : Long.parseLong(raw.get(5))); - result.add(raw.get(6).equals("\\N") ? null : Double.parseDouble(raw.get(6))); - result.add(raw.get(7).equals("\\N") ? null : raw.get(7)); - return result; - } - - private static long toPartitionTs(long ts) { - LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); - return time.truncatedTo(ChronoUnit.DAYS).withDayOfMonth(1).toInstant(ZoneOffset.UTC).toEpochMilli(); -// return TsPartitionDate.MONTHS.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); - } - - private static boolean isBlockStarted(String line) { - return line.startsWith("COPY public.ts_kv ("); - } - - private static boolean isBlockFinished(String line) { - return StringUtils.isBlank(line) || line.equals("\\."); - } - -} diff --git a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java b/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java index ecd5164c70..fac2609797 100644 --- a/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java +++ b/tools/src/main/java/org/thingsboard/client/tools/migrator/RelatedEntitiesParser.java @@ -42,39 +42,43 @@ public class RelatedEntitiesParser { private void processAllTables(LineIterator lineIterator) { String currentLine; - while(lineIterator.hasNext()) { - currentLine = lineIterator.nextLine(); - if(currentLine.startsWith("COPY public.alarm")) { - processBlock(lineIterator, EntityType.ALARM); - } else if (currentLine.startsWith("COPY public.asset")) { - processBlock(lineIterator, EntityType.ASSET); - } else if (currentLine.startsWith("COPY public.customer")) { - processBlock(lineIterator, EntityType.CUSTOMER); - } else if (currentLine.startsWith("COPY public.dashboard")) { - processBlock(lineIterator, EntityType.DASHBOARD); - } else if (currentLine.startsWith("COPY public.device")) { - processBlock(lineIterator, EntityType.DEVICE); - } else if (currentLine.startsWith("COPY public.rule_chain")) { - processBlock(lineIterator, EntityType.RULE_CHAIN); - } else if (currentLine.startsWith("COPY public.rule_node")) { - processBlock(lineIterator, EntityType.RULE_NODE); - } else if (currentLine.startsWith("COPY public.tenant")) { - processBlock(lineIterator, EntityType.TENANT); - } else if (currentLine.startsWith("COPY public.tb_user")) { - processBlock(lineIterator, EntityType.USER); - } else if (currentLine.startsWith("COPY public.entity_view")) { - processBlock(lineIterator, EntityType.ENTITY_VIEW); - } else if (currentLine.startsWith("COPY public.widgets_bundle")) { - processBlock(lineIterator, EntityType.WIDGETS_BUNDLE); - } else if (currentLine.startsWith("COPY public.widget_type")) { - processBlock(lineIterator, EntityType.WIDGET_TYPE); - } else if (currentLine.startsWith("COPY public.tenant_profile")) { - processBlock(lineIterator, EntityType.TENANT_PROFILE); - } else if (currentLine.startsWith("COPY public.device_profile")) { - processBlock(lineIterator, EntityType.DEVICE_PROFILE); - } else if (currentLine.startsWith("COPY public.api_usage_state")) { - processBlock(lineIterator, EntityType.API_USAGE_STATE); + try { + while (lineIterator.hasNext()) { + currentLine = lineIterator.nextLine(); + if (currentLine.startsWith("COPY public.alarm")) { + processBlock(lineIterator, EntityType.ALARM); + } else if (currentLine.startsWith("COPY public.asset")) { + processBlock(lineIterator, EntityType.ASSET); + } else if (currentLine.startsWith("COPY public.customer")) { + processBlock(lineIterator, EntityType.CUSTOMER); + } else if (currentLine.startsWith("COPY public.dashboard")) { + processBlock(lineIterator, EntityType.DASHBOARD); + } else if (currentLine.startsWith("COPY public.device")) { + processBlock(lineIterator, EntityType.DEVICE); + } else if (currentLine.startsWith("COPY public.rule_chain")) { + processBlock(lineIterator, EntityType.RULE_CHAIN); + } else if (currentLine.startsWith("COPY public.rule_node")) { + processBlock(lineIterator, EntityType.RULE_NODE); + } else if (currentLine.startsWith("COPY public.tenant")) { + processBlock(lineIterator, EntityType.TENANT); + } else if (currentLine.startsWith("COPY public.tb_user")) { + processBlock(lineIterator, EntityType.USER); + } else if (currentLine.startsWith("COPY public.entity_view")) { + processBlock(lineIterator, EntityType.ENTITY_VIEW); + } else if (currentLine.startsWith("COPY public.widgets_bundle")) { + processBlock(lineIterator, EntityType.WIDGETS_BUNDLE); + } else if (currentLine.startsWith("COPY public.widget_type")) { + processBlock(lineIterator, EntityType.WIDGET_TYPE); + } else if (currentLine.startsWith("COPY public.tenant_profile")) { + processBlock(lineIterator, EntityType.TENANT_PROFILE); + } else if (currentLine.startsWith("COPY public.device_profile")) { + processBlock(lineIterator, EntityType.DEVICE_PROFILE); + } else if (currentLine.startsWith("COPY public.api_usage_state")) { + processBlock(lineIterator, EntityType.API_USAGE_STATE); + } } + } finally { + lineIterator.close(); } }