Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,240 @@
/*
* 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.db.it;

import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.it.env.EnvFactory;
import org.apache.iotdb.it.framework.IoTDBTestRunner;
import org.apache.iotdb.it.utils.TsFileGenerator;
import org.apache.iotdb.itbase.category.ClusterIT;
import org.apache.iotdb.jdbc.Config;

import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.enums.TSEncoding;
import org.apache.tsfile.write.schema.IMeasurementSchema;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;

import java.io.File;
import java.nio.file.Files;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.Statement;
import java.util.Arrays;
import java.util.List;

@RunWith(IoTDBTestRunner.class)
@Category({ClusterIT.class})
public class IoTDBLoadTsFileClusterIT {

private static final long PARTITION_INTERVAL = 10_000L;
private static final String DATABASE = "root.load_cluster";
private static final List<String> DEVICES =
Arrays.asList(
DATABASE + ".d1", DATABASE + ".d2", DATABASE + ".d3", DATABASE + ".d4", DATABASE + ".d5");
private static final List<String> ALIGNED_DEVICES =
Arrays.asList(DATABASE + ".a1", DATABASE + ".a2");
private static final List<String> MEASUREMENTS = Arrays.asList("s1", "s2", "s3", "s4");
private static final int POINT_COUNT_PER_DEVICE = 10_000;
private static final int REPLICATION_FACTOR = 3;
private static final int CONFIG_NODE_NUM = 3;
private static final int DATA_NODE_NUM = 3;

private File tmpDir;

@Before
public void setUp() throws Exception {
tmpDir = new File(Files.createTempDirectory("load-cluster-it").toUri());
EnvFactory.getEnv().getConfig().getCommonConfig().setTimePartitionInterval(PARTITION_INTERVAL);
EnvFactory.getEnv().getConfig().getCommonConfig().setEnforceStrongPassword(false);
EnvFactory.getEnv().getConfig().getCommonConfig().setPipeMemoryManagementEnabled(false);
EnvFactory.getEnv().getConfig().getCommonConfig().setDatanodeMemoryProportion("1:10:1:1:1:0");
EnvFactory.getEnv().getConfig().getCommonConfig().setTargetChunkPointNum(1_000);
EnvFactory.getEnv().getConfig().getCommonConfig().setMaxNumberOfPointsInPage(500);
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
.setConfigNodeConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
.setSchemaReplicationFactor(REPLICATION_FACTOR)
.setDataReplicationFactor(REPLICATION_FACTOR);
EnvFactory.getEnv().initClusterEnvironment(CONFIG_NODE_NUM, DATA_NODE_NUM);
}

@After
public void tearDown() throws Exception {
try (final Connection connection = EnvFactory.getEnv().getConnection();
final Statement statement = connection.createStatement()) {
statement.execute("delete database " + DATABASE);
} catch (final Exception ignored) {
}

EnvFactory.getEnv().cleanClusterEnvironment();

final File[] files = tmpDir.listFiles();
if (files != null) {
for (final File file : files) {
Assert.assertTrue(file.delete());
}
}
Assert.assertTrue(tmpDir.delete());
}

@Test
public void testLoadTsFileInCluster() throws Exception {
final long writtenPointCount;
try (final TsFileGenerator generator =
new TsFileGenerator(new File(tmpDir, "load-cluster-1-0-0-0.tsfile"))) {
final List<IMeasurementSchema> schemas =
Arrays.asList(
new MeasurementSchema(MEASUREMENTS.get(0), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(1), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(2), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(3), TSDataType.INT64, TSEncoding.PLAIN));
for (final String device : DEVICES) {
generator.registerTimeseries(device, schemas);
generator.generateData(device, POINT_COUNT_PER_DEVICE, 1L, false);
}
for (final String device : ALIGNED_DEVICES) {
generator.registerAlignedTimeseries(device, schemas);
generator.generateData(device, POINT_COUNT_PER_DEVICE, 1L, true, 1_000_000L);
}
writtenPointCount = generator.getTotalNumber();
}

try (final Connection connection = EnvFactory.getEnv().getConnection();
final Statement statement = connection.createStatement()) {
statement.execute("create database " + DATABASE);
for (final String device : DEVICES) {
for (final String measurement : MEASUREMENTS) {
statement.execute(
"create timeseries " + device + "." + measurement + " " + TSDataType.INT64.name());
}
}
for (final String device : ALIGNED_DEVICES) {
statement.execute(
"create aligned timeseries " + device + "(s1 INT64, s2 INT64, s3 INT64, s4 INT64)");
}
statement.execute("load \"" + tmpDir.getAbsolutePath() + "\"");

long actualPointCount = 0;
for (final String device : DEVICES) {
try (final ResultSet resultSet =
statement.executeQuery(
"select count(s1), count(s2), count(s3), count(s4) from " + device)) {
Assert.assertTrue(resultSet.next());
for (int columnIndex = 1; columnIndex <= MEASUREMENTS.size(); columnIndex++) {
actualPointCount += resultSet.getLong(columnIndex);
}
}
}
Assert.assertEquals(writtenPointCount, actualPointCount);
}
}

@Test
public void testLoadAfterOneDataNodeDown() throws Exception {
final long firstPointCount = generateTsFile("load-first.tsfile", 5_000, 0L);

try (final Connection connection = EnvFactory.getEnv().getConnection();
final Statement statement = connection.createStatement()) {
statement.execute("create database " + DATABASE);
for (final String device : DEVICES) {
for (final String measurement : MEASUREMENTS) {
statement.execute(
"create timeseries " + device + "." + measurement + " " + TSDataType.INT64.name());
}
}
for (final String device : ALIGNED_DEVICES) {
statement.execute(
"create aligned timeseries " + device + "(s1 INT64, s2 INT64, s3 INT64, s4 INT64)");
}

statement.execute("load \"" + tmpDir.getAbsolutePath() + "\"");
}

EnvFactory.getEnv().getDataNodeWrapper(0).stop();
Thread.sleep(1_000L);

final long secondPointCount = generateTsFile("load-second.tsfile", 3_000, 100_000L);
try (final Connection connection =
DriverManager.getConnection(
Config.IOTDB_URL_PREFIX
+ EnvFactory.getEnv().getDataNodeWrapper(1).getIpAndPortString(),
"root",
"root");
final Statement statement = connection.createStatement()) {
statement.execute("load \"" + tmpDir.getAbsolutePath() + "\"");

final long expectedPointCount = firstPointCount + secondPointCount;
long actualPointCount = queryPointCount(statement);
for (int retry = 0; actualPointCount < expectedPointCount && retry < 30; retry++) {
actualPointCount = queryPointCount(statement);
}
Assert.assertEquals(expectedPointCount, actualPointCount);
}
}

private long generateTsFile(
final String fileName, final int pointCountPerDevice, final long startTimestamp)
throws Exception {
final List<IMeasurementSchema> schemas =
Arrays.asList(
new MeasurementSchema(MEASUREMENTS.get(0), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(1), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(2), TSDataType.INT64, TSEncoding.PLAIN),
new MeasurementSchema(MEASUREMENTS.get(3), TSDataType.INT64, TSEncoding.PLAIN));
try (final TsFileGenerator generator = new TsFileGenerator(new File(tmpDir, fileName))) {
for (final String device : DEVICES) {
generator.registerTimeseries(device, schemas);
generator.generateData(device, pointCountPerDevice, 1L, false, startTimestamp);
}
final long alignedStartTimestamp = startTimestamp + 1_000_000L;
for (final String device : ALIGNED_DEVICES) {
generator.registerAlignedTimeseries(device, schemas);
generator.generateData(device, pointCountPerDevice, 1L, true, alignedStartTimestamp);
}
return generator.getTotalNumber();
}
}

private long queryPointCount(final Statement statement) throws Exception {
long actualPointCount = 0;
final List<String> allDevices = new java.util.ArrayList<>(DEVICES);
allDevices.addAll(ALIGNED_DEVICES);
for (final String device : allDevices) {
for (final String measurement : MEASUREMENTS) {
try (final ResultSet resultSet =
statement.executeQuery("select count(" + measurement + ") from " + device)) {
Assert.assertTrue(resultSet.next());
actualPointCount += resultSet.getLong(1);
}
}
}
return actualPointCount;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -1175,6 +1175,9 @@ public final class DataNodeQueryMessages {
"Start load TsFile {} locally.";
public static final String LOAD_ALL_FAILED_TSFILES_ARE_CONVERTED_TO_TABLETS =
"Load: all failed TsFiles are converted to tablets and inserted.";
public static final String
LOG_LOAD_FAILED_TO_LOAD_SOME_TSFILES_BY_CONVERTING_THEM_INTO_TABLETS_FAILED_TSFILES_ARG_7D9DB9C3 =
"Load: failed to load some TsFiles by converting them into tablets. Failed TsFiles: %s";

// --- Plan / Statement ---

Expand Down Expand Up @@ -2703,6 +2706,8 @@ public final class DataNodeQueryMessages {
"LoadTsFileScheduler: Region migration was detected during loading TsFile {}, will convert to insertion to avoid data loss";
public static final String LOAD_TSFILE_ARG_SUCCESSFULLY_LOAD_PROCESS_ARG_ARG =
"Load TsFile {} Successfully, load process [{}/{}]";
public static final String LOG_LOAD_BATCH_FINISHED_DELETING_ARG_SOURCE_TSFILES_AFTER_LOAD_D5EE56E9 =
"LOAD batch finished, deleting %d source TsFiles after load";
public static final String CAN_NOT_LOAD_TSFILE_ARG_LOAD_PROCESS_ARG_ARG =
"Can not Load TsFile {}, load process [{}/{}]";
public static final String LOAD_TSFILE_S_FAILED_WILL_TRY_TO_CONVERT_TO_TABLETS_AND_INSERT_FAILED_TSFILES_ARG =
Expand All @@ -2713,6 +2718,8 @@ public final class DataNodeQueryMessages {
"Parse or send TsFile %s error.";
public static final String DISPATCH_ONE_PIECE_TO_REPLICASET_ARG_ERROR_RESULT_STATUS_CODE_ARG =
"Dispatch one piece to ReplicaSet {} error. Result status code {}. ";
public static final String LOG_LOAD_CONSENSUS_SUBMIT_TRANSIENT_FAILURE_RETRY_D7E1D9A6 =
"Transient failure while submitting LOAD consensus {} (load {}) to {}, will retry ({}/{}): {}";
public static final String RESULT_STATUS_MESSAGE_ARG_DISPATCH_PIECE_NODE_ERROR_PERCENT_NARG =
"Result status message {}. Dispatch piece node error:%n{}";
public static final String SUB_STATUS_CODE_ARG_SUB_STATUS_MESSAGE_ARG =
Expand Down Expand Up @@ -3816,6 +3823,7 @@ private DataNodeQueryMessages() {}
public static final String EXCEPTION_THE_SECOND_ARGUMENT_OF_PERCENTILE_FUNCTION_PERCENTAGE_MUST_BE_A_DOUBLE_LITERAL_D9464B46 = "The second argument of 'percentile' function percentage must be a double literal";
public static final String EXCEPTION_DATA_TYPE_MISMATCH_FOR_MEASUREMENT_ARGARGARG_TYPE_IN_TSFILE_ARG_TYPE_IN_IOTDB_ARG_C5BA7DBD = "Data type mismatch for measurement %s%s%s, type in TsFile: %s, type in IoTDB: %s";
public static final String MESSAGE_FAILED_TO_RELEASE_EXTERNAL_TSFILE_QUERY_RESOURCE_712EE978 = "Failed to release external TsFile query resource";
public static final String EXCEPTION_UNKNOWN_LOADTSFILECONSENSUSOP_ORDINAL_ARG_62848FC2 = "Unknown LoadTsFileConsensusOp ordinal: ";
public static final String EXCEPTION_OUTER_QUERY_TIMEOUT_EXCEEDED_BEFORE_IOTDBLOCAL_QUERY_STARTS_800BFA63 = "Outer query timeout exceeded before IoTDBLocal query starts";
public static final String MESSAGE_FAILED_TO_CLOSE_UDF_RESULT_SET_AT_INDEX_ARG_A293B7EC = "Failed to close UDF result set at index {}";
public static final String EXCEPTION_INTERNAL_QUERY_EXECUTION_NOT_FOUND_62642542 = "Internal query execution not found";
Expand Down Expand Up @@ -3879,5 +3887,7 @@ private DataNodeQueryMessages() {}
"No more DeviceEntry records are available";
public static final String EXCEPTION_ONLY_INMEMORYDEVICEENTRYDATASET_SUPPORTS_GET_INLINE_DEVICE_ENTRIES_07A52CAB =
"Only InMemoryDeviceEntryDataSet supports get inline device entries";
public static final String EXCEPTION_LOAD_CONSENSUS_INVALID_PIECE_REF_F3498507 =
"Invalid LOAD consensus piece ref: path %s, offset %d, size %d";

}
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,54 @@ private StorageEngineMessages() {}
// ======================== StorageEngine ========================

public static final String FAIL_TO_RECOVER_WAL = "Fail to recover wal.";
public static final String LOG_LOAD_CONSENSUS_WRITE_TO_REGION_ARG_VIA_PROTOCOL_ARG_EBB55042 =
"Write LOAD consensus node to region {} via protocol {}";
public static final String LOG_LOAD_CONSENSUS_REFRESH_REPLICA_SET_FAILED_7C244C63 =
"Failed to refresh LOAD consensus replica set for region {}, using cached set: {}";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_CHECKSUM_MISMATCH_CF261675 =
"LOAD consensus piece checksum mismatch, loadId: %s, pieceIndex: %d";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_EOF_8743387D =
"Unexpected end of file when reading staged piece %s at offset %s.";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_DATA_MISSING_AFTER_PULL_8269CB0B =
"LOAD piece %d data of load %s is still missing after pulling from the write node.";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_PIECE_NOT_SUPPORTED_71EC4B46 =
"Fetching piece from leader is not supported.";
public static final String EXCEPTION_LOAD_CONSENSUS_PIECE_DATA_MISSING_OR_CHECKSUM_MISMATCH_AFTER_PULL_35F4972E =
"LOAD task %s piece %d data is missing or its checksum mismatches after pull.";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_READ_FAILED_0659D19B =
"Failed to read retained LOAD piece {} of load {} from {}: {}";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_WRITE_FAILED_99697608 =
"Failed to write retained LOAD piece {} of load {} to {}: {}";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_NOT_CONTINUOUS_F9408C19 =
"Staged file %s of load %s is not continuous: expected offset %d but current file length is %d.";
public static final String MESSAGE_LOAD_CONSENSUS_PREPARE_WITHOUT_STAGED_DATA_FE8ADC37 =
"Cannot prepare load %s because no staged data exists on this node.";
public static final String LOG_LOAD_CONSENSUS_RECOVERED_TASK_02824CE6 =
"Recovered in-progress LOAD task {} from disk; staged data is kept until COMMIT or ABORT.";
public static final String LOG_LOAD_CONSENSUS_RECOVER_TASK_META_FAILED_C39E04BB =
"Failed to recover the task meta of LOAD task {}: {}";
public static final String LOG_LOAD_CONSENSUS_TASK_META_WRITE_FAILED_5D2420BF =
"Failed to persist the task meta of LOAD task {} to {}: {}";
public static final String LOG_LOAD_CONSENSUS_TERMINAL_MARKER_WRITE_FAILED_4D6D7433 =
"Failed to write the terminal marker of LOAD task {} to {}: {}";
public static final String LOG_LOAD_CONSENSUS_RECOVER_TASK_UNRESUMABLE_A159436C =
"Cannot resume LOAD task {} from disk: its durable task meta is missing or corrupt; the next "
+ "command will re-create it.";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_SHORT_WRITE_E7392FAD =
"Staged file %s of load %s was not fully written: expected %d bytes at offset %d, wrote %d";
public static final String EXCEPTION_LOAD_CONSENSUS_PROGRESS_SERIALIZE_FAILED_28EFD091 =
"Failed to serialize LOAD progress index for time partition %d";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_INCOMPLETE_1CDE954B =
"Staged file %s of load %s is incomplete and cannot be committed.";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_TAKEN_09A7DD4C =
"Snapshotted %d in-progress LOAD task(s) with %d staged file(s) for region %s into %s.";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_RESTORED_90ABC1BF =
"Restored %d in-progress LOAD task(s) with %d staged file(s) from snapshot %s.";
public static final String EXCEPTION_LOAD_CONSENSUS_SNAPSHOT_RESTORE_FAILED_F8C29C64 =
"Failed to restore LOAD snapshot from %s: %s";
public static final String EXCEPTION_LOAD_TSFILE_ALIGNED_VALUE_CHUNK_TIME_CHUNK_EEB00760 =
"Cannot attach value chunk of measurement %s in file %s: expected exactly one buffered "
+ "aligned time chunk, found %d.";
public static final String STORAGE_ENGINE_FAILED_TO_SET_UP = "Storage engine failed to set up.";
public static final String SEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "start sequence memtable timed flush check thread successfully.";
public static final String UNSEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "start unsequence memtable timed flush check thread successfully.";
Expand Down Expand Up @@ -478,6 +526,16 @@ private StorageEngineMessages() {}
public static final String CANNOT_CREATE_TSFILE_FOR_WRITING = "Can not create TsFile {} for writing.";
public static final String CLOSE_TSFILE_IO_WRITER_ERROR = "Close TsFileIOWriter {} error.";
public static final String CLOSE_MODIFICATION_FILE_ERROR = "Close ModificationFile {} error.";
public static final String LOG_PREPARING_LOAD_TSFILE_ARG_SEALING_STAGED_RESOURCES_1FDF1866 =
"Preparing LOAD TsFile {}: sealing staged resources.";
public static final String LOG_COMMITTING_LOAD_TSFILE_ARG_LOADING_PREPARED_RESOURCES_INTO_DATAREGION_EA1D6335 =
"Committing LOAD TsFile {}: loading prepared resources into DataRegion.";
public static final String LOG_RECEIVE_LOAD_TSFILE_NODE_ARG_C36E832B =
"Receive LOAD TsFile node: {}.";
public static final String LOG_RECEIVE_LOAD_TSFILE_NODE_SUCCESS_ARG_27F8ECD6 =
"Receive LOAD TsFile node success: {}.";
public static final String EXCEPTION_TABLE_ARG_ARG_DOES_NOT_EXIST_WHEN_APPLYING_LOAD_CHUNK_DATA_IT_MAY_HAVE_BEEN_DROPPED_AFTER_THE_LOAD_WAS_ANALYZED_DDB35F93 =
"Table '%s.%s' does not exist when applying LOAD chunk data. It may have been dropped after the LOAD was analyzed.";
public static final String TASK_DIR_NOT_EMPTY_SKIP_DELETE = "Task dir {} is not empty, skip deleting.";
public static final String LOAD_CLEANUP_TASK_CANCELED = "Load cleanup task {} is canceled.";
public static final String LOAD_CLEANUP_TASK_STARTS = "Load cleanup task {} starts.";
Expand Down
Loading
Loading