Skip to content
Open
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
Expand Up @@ -895,6 +895,66 @@ public void testExtractorTimeRangeMatch() throws Exception {
}
}

@Test
public void testSourceTimeRangeRespectsHistoryDisable() throws Exception {
final DataNodeWrapper receiverDataNode = receiverEnv.getDataNodeWrapper(0);

final String receiverIp = receiverDataNode.getIp();
final int receiverPort = receiverDataNode.getPort();

TestUtils.executeNonQueries(
senderEnv,
Arrays.asList(
"insert into root.db.history (time, at1) values (2000, 2), (3000, 3)", "flush"),
null);

final Map<String, String> sourceAttributes = new HashMap<>();
final Map<String, String> sinkAttributes = new HashMap<>();

sourceAttributes.put("source.inclusion", "data");
sourceAttributes.put("source.start-time", "2000");
sourceAttributes.put("source.history.enable", "false");
sourceAttributes.put("source.realtime.mode", "stream");
sourceAttributes.put("user", "root");

sinkAttributes.put("sink", "iotdb-thrift-sink");
sinkAttributes.put("sink.batch.enable", "false");
sinkAttributes.put("sink.ip", receiverIp);
sinkAttributes.put("sink.port", Integer.toString(receiverPort));

try (final SyncConfigNodeIServiceClient client =
(SyncConfigNodeIServiceClient) senderEnv.getLeaderConfigNodeConnection()) {
final TSStatus status =
client.createPipe(
new TCreatePipeReq("p1", sinkAttributes).setExtractorAttributes(sourceAttributes));
Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), status.getCode());

TestUtils.assertDataAlwaysOnEnv(
receiverEnv,
"show timeseries root.db.history.**",
"Timeseries,Alias,Database,DataType,Encoding,Compression,Tags,Attributes,Deadband,DeadbandParameters,ViewType,",
Collections.emptySet());

TestUtils.executeNonQueries(
senderEnv,
Collections.singletonList(
"insert into root.db.realtime (time, at1)"
+ " values (1000, 1), (2000, 2), (3000, 3)"),
null);

TestUtils.assertDataEventuallyOnEnv(
receiverEnv,
"select count(at1) from root.db.realtime",
"count(root.db.realtime.at1),",
Collections.singleton("2,"));
TestUtils.assertDataAlwaysOnEnv(
receiverEnv,
"show timeseries root.db.history.**",
"Timeseries,Alias,Database,DataType,Encoding,Compression,Tags,Attributes,Deadband,DeadbandParameters,ViewType,",
Collections.emptySet());
}
}

@Test
public void testSourceStartTimeAndEndTimeWorkingWithOrWithoutPattern() throws Exception {
final DataNodeWrapper receiverDataNode = receiverEnv.getDataNodeWrapper(0);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -651,8 +651,8 @@ public final class DataNodePipeMessages {
"When '{}' ('{}') is set to false, specifying {} and {} is invalid.";
public static final String WHEN_IS_SET_TO_TRUE_SPECIFYING_AND =
"When '{}' ('{}', '{}', '{}') is set to true, specifying {} and {} is invalid.";
public static final String WHEN_OR_IS_SPECIFIED_SPECIFYING_AND_IS =
"When {}, {}, {} or {} is specified, specifying {}, {}, {}, {}, {} and {} is invalid.";
public static final String WHEN_OR_IS_SPECIFIED_SPECIFYING_OR_IS_INVALID =
"When {}, {}, {} or {} is specified, specifying {}, {}, {} or {} is invalid.";

// ===================== SINK =====================

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -616,8 +616,8 @@ public final class DataNodePipeMessages {
"当 '{}'('{}')设置为 false 时,指定 {} 和 {} 无效。";
public static final String WHEN_IS_SET_TO_TRUE_SPECIFYING_AND =
"当 '{}'('{}'、'{}'、'{}')设置为 true 时,指定 {} 和 {} 无效。";
public static final String WHEN_OR_IS_SPECIFIED_SPECIFYING_AND_IS =
"当指定 {}、{}、{} 或 {} 时,指定 {}、{}、{}、{}、{} 和 {} 无效。";
public static final String WHEN_OR_IS_SPECIFIED_SPECIFYING_OR_IS_INVALID =
"当指定 {}、{}、{} 或 {} 时,指定 {}、{}、{} {} 无效。";

// ===================== SINK =====================

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -257,27 +257,23 @@ private void validatePattern(final TreePattern treePattern) {
private void checkInvalidParameters(final PipeParameterValidator validator) {
final PipeParameters parameters = validator.getParameters();

// Enable history and realtime if specifying start-time or end-time
// A global time range takes precedence over a history-specific time range.
if (parameters.hasAnyAttributes(
SOURCE_START_TIME_KEY,
EXTRACTOR_START_TIME_KEY,
SOURCE_END_TIME_KEY,
EXTRACTOR_END_TIME_KEY)
&& parameters.hasAnyAttributes(
EXTRACTOR_HISTORY_ENABLE_KEY,
SOURCE_HISTORY_ENABLE_KEY,
SOURCE_HISTORY_START_TIME_KEY,
EXTRACTOR_HISTORY_START_TIME_KEY,
SOURCE_HISTORY_END_TIME_KEY,
EXTRACTOR_HISTORY_END_TIME_KEY)) {
LOGGER.warn(
DataNodePipeMessages.WHEN_OR_IS_SPECIFIED_SPECIFYING_AND_IS,
DataNodePipeMessages.WHEN_OR_IS_SPECIFIED_SPECIFYING_OR_IS_INVALID,
SOURCE_START_TIME_KEY,
EXTRACTOR_START_TIME_KEY,
SOURCE_END_TIME_KEY,
EXTRACTOR_END_TIME_KEY,
SOURCE_HISTORY_ENABLE_KEY,
EXTRACTOR_HISTORY_ENABLE_KEY,
SOURCE_HISTORY_START_TIME_KEY,
EXTRACTOR_HISTORY_START_TIME_KEY,
SOURCE_HISTORY_END_TIME_KEY,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -243,13 +243,23 @@ public void validate(final PipeParameterValidator validator) {
}
}

// Historical data extraction is enabled in the following cases:
// 1. System restarts the pipe. If the pipe is restarted but historical data extraction is not
// enabled, the pipe will lose some historical data.
// 2. Historical extraction is enabled by the user or by default.
isHistoricalSourceEnabled =
parameters.getBooleanOrDefault(
SystemConstant.RESTART_OR_NEWLY_ADDED_KEY,
SystemConstant.RESTART_OR_NEWLY_ADDED_DEFAULT_VALUE)
|| parameters.getBooleanOrDefault(
Arrays.asList(EXTRACTOR_HISTORY_ENABLE_KEY, SOURCE_HISTORY_ENABLE_KEY),
EXTRACTOR_HISTORY_ENABLE_DEFAULT_VALUE);

if (parameters.hasAnyAttributes(
SOURCE_START_TIME_KEY,
EXTRACTOR_START_TIME_KEY,
SOURCE_END_TIME_KEY,
EXTRACTOR_END_TIME_KEY)) {
isHistoricalSourceEnabled = true;

try {
historicalDataExtractionStartTime =
parameters.hasAnyAttributes(SOURCE_START_TIME_KEY, EXTRACTOR_START_TIME_KEY)
Expand Down Expand Up @@ -284,19 +294,6 @@ public void validate(final PipeParameterValidator validator) {
return;
}

// Historical data extraction is enabled in the following cases:
// 1. System restarts the pipe. If the pipe is restarted but historical data extraction is not
// enabled, the pipe will lose some historical data.
// 2. User may set the EXTRACTOR_HISTORY_START_TIME and EXTRACTOR_HISTORY_END_TIME without
// enabling the historical data extraction, which may affect the realtime data extraction.
isHistoricalSourceEnabled =
parameters.getBooleanOrDefault(
SystemConstant.RESTART_OR_NEWLY_ADDED_KEY,
SystemConstant.RESTART_OR_NEWLY_ADDED_DEFAULT_VALUE)
|| parameters.getBooleanOrDefault(
Arrays.asList(EXTRACTOR_HISTORY_ENABLE_KEY, SOURCE_HISTORY_ENABLE_KEY),
EXTRACTOR_HISTORY_ENABLE_DEFAULT_VALUE);

try {
historicalDataExtractionStartTime =
parameters.hasAnyAttributes(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.iotdb.commons.consensus.index.impl.TimePartitionProgressIndex;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
import org.apache.iotdb.commons.pipe.config.plugin.configuraion.PipeTaskRuntimeConfiguration;
import org.apache.iotdb.commons.pipe.config.plugin.env.PipeTaskSourceRuntimeEnvironment;
import org.apache.iotdb.commons.pipe.datastructure.resource.PersistentResource;
Expand Down Expand Up @@ -193,6 +194,41 @@ public void testHistoricalTsFileQueryPriorityOrderDefaultsToTrue() throws Except
(Boolean) getPrivateField(source, "shouldOrderHistoricalTsFileByQueryPriority"));
}

@Test
public void testGlobalTimeRangeRespectsHistoryEnable() throws Exception {
final Map<String, String> attributes = new HashMap<>();
attributes.put(PipeSourceConstant.SOURCE_START_TIME_KEY, "1000");
attributes.put(PipeSourceConstant.SOURCE_HISTORY_ENABLE_KEY, Boolean.FALSE.toString());

final PipeHistoricalDataRegionTsFileAndDeletionSource realtimeOnlySource =
new PipeHistoricalDataRegionTsFileAndDeletionSource();
realtimeOnlySource.validate(
new PipeParameterValidator(new PipeParameters(new HashMap<>(attributes))));

Assert.assertFalse((Boolean) getPrivateField(realtimeOnlySource, "isHistoricalSourceEnabled"));
Assert.assertEquals(
1000L,
((Long) getPrivateField(realtimeOnlySource, "historicalDataExtractionStartTime"))
.longValue());

final PipeHistoricalDataRegionTsFileAndDeletionSource defaultSource =
new PipeHistoricalDataRegionTsFileAndDeletionSource();
attributes.remove(PipeSourceConstant.SOURCE_HISTORY_ENABLE_KEY);
defaultSource.validate(
new PipeParameterValidator(new PipeParameters(new HashMap<>(attributes))));

Assert.assertTrue((Boolean) getPrivateField(defaultSource, "isHistoricalSourceEnabled"));

final PipeHistoricalDataRegionTsFileAndDeletionSource restartedSource =
new PipeHistoricalDataRegionTsFileAndDeletionSource();
attributes.put(PipeSourceConstant.SOURCE_HISTORY_ENABLE_KEY, Boolean.FALSE.toString());
attributes.put(SystemConstant.RESTART_OR_NEWLY_ADDED_KEY, Boolean.TRUE.toString());
restartedSource.validate(
new PipeParameterValidator(new PipeParameters(new HashMap<>(attributes))));

Assert.assertTrue((Boolean) getPrivateField(restartedSource, "isHistoricalSourceEnabled"));
}

@Test
public void testHistoricalTsFileQueryPriorityOrderMatchesQueryCoverage() throws Exception {
final PipeHistoricalDataRegionTsFileAndDeletionSource source =
Expand Down
Loading