Repository navigation
feat(load): route decoded tsfile pieces through consensus #18698
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
luoluoyuyu
wants to merge
20
commits into
apache:master
Choose a base branch
from
luoluoyuyu:load-tsfile-consensus-ha1
base: master
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+16,212
−2,636
Open
Changes from all commits
Commits
Show all changes
20 commits
Select commit
Hold shift + click to select a range
45946c3
LOAD TsFile: replicate pieces through consensus and clean up staged d…
luoluoyuyu e86844c
LOAD TsFile: commit only after every region agreed to prepare
luoluoyuyu 815576e
LOAD TsFile: retry the terminal commands of the 2PC protocol
luoluoyuyu bc1e5a2
LOAD TsFile: fix the routing, completeness and snapshot findings of t…
luoluoyuyu fa81502
LOAD TsFile: keep staged references valid across roots and nodes
luoluoyuyu 8728f5b
LOAD TsFile: pin the route of a load transaction
luoluoyuyu 8339d1a
LOAD TsFile: answer a repeated COMMIT of a committed task with success
luoluoyuyu d33da4a
LOAD TsFile: cover the replayed commands and the concurrent snapshot
luoluoyuyu d9fad74
LOAD TsFile: snapshot the staging directory only where it is transferred
luoluoyuyu 529a352
LOAD TsFile: follow the route a region has when a command fails
luoluoyuyu 471a9df
Restore the LOAD TsFile consensus and 2PC design docs and sync them
luoluoyuyu ddd3e90
Revert "Restore the LOAD TsFile consensus and 2PC design docs and syn…
luoluoyuyu bcd33df
Merge remote-tracking branch 'iotdb/master' into load-tsfile-consensu…
luoluoyuyu ea38fa3
LOAD TsFile: schedule the cleaner as a periodical job of the DataNode
luoluoyuyu 313d0e2
LOAD TsFile: harden the staged copy, the piece references and the com…
luoluoyuyu 963eac4
Merge remote-tracking branch 'iotdb/master' into load-tsfile-consensu…
luoluoyuyu d0328fb
LOAD TsFile: precalculate the staged chunk offsets and route the load…
luoluoyuyu 732327a
LOAD TsFile: fail a piece of a staged file with the file it cannot re…
luoluoyuyu b0de67b
LOAD TsFile: cover the staged deletions of a LOAD piece
luoluoyuyu 4bedc78
LOAD TsFile: do not invent a failure for a load stopped without one
luoluoyuyu File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
231 changes: 231 additions & 0 deletions
231
integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileClusterIT.java
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,231 @@ | ||
| /* | ||
| * 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 { | ||
|
Comment on lines
+107
to
+108
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Test load concurrently. |
||
| 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() + "\""); | ||
|
|
||
| // queryPointCount also covers the aligned devices, which writtenPointCount includes. | ||
| Assert.assertEquals(writtenPointCount, queryPointCount(statement)); | ||
| } | ||
| } | ||
|
|
||
| @Test | ||
| public void testLoadAfterOneDataNodeDown() throws Exception { | ||
|
Comment on lines
+149
to
+150
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Test load with two failing nodes returns meaningful error message. |
||
| 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++) { | ||
| Thread.sleep(1_000L); | ||
| 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; | ||
| } | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Better to use setUpClass