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 @@ -438,6 +438,11 @@ public void testHistoricalActivationRace() throws Exception {
Assert.assertEquals(
TSStatusCode.SUCCESS_STATUS.getStatusCode(), client.startPipe("testPipe").getCode());

TestUtils.assertDataEventuallyOnEnv(
receiverEnv,
"show paths set device template aligned_template",
"Paths,",
Collections.singleton("root.sg_aligned.device_aligned,"));
TestUtils.assertDataEventuallyOnEnv(
receiverEnv,
"count devices root.sg_aligned.**",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,11 +57,20 @@
import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_REALTIME_ENABLE_DEFAULT_VALUE;
import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_REALTIME_ENABLE_KEY;
import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_REALTIME_ENABLE_KEY;
import static org.apache.iotdb.commons.pipe.datastructure.options.PipeInclusionOptions.areOptionsEnabled;

public class PipeDataNodeTaskBuilder {

private static final Logger LOGGER = LoggerFactory.getLogger(PipeDataNodeTaskBuilder.class);

private static final String[] SCHEMA_OPTIONS_REQUIRED_BEFORE_LOAD = {
"schema.timeseries.ordinary.create",
"schema.timeseries.template.create",
"schema.timeseries.template.alter",
"schema.timeseries.template.set",
"schema.timeseries.template.activate"
};

private final PipeStaticMeta pipeStaticMeta;
private final int regionId;
private final PipeTaskMeta pipeTaskMeta;
Expand Down Expand Up @@ -265,13 +274,46 @@ private static void checkConflict(

private static void injectParameters(
final PipeParameters sourceParameters, final PipeParameters sinkParameters) {
final boolean isSourceExternal =
!BuiltinPipePlugin.BUILTIN_SOURCES.contains(
sourceParameters
.getStringOrDefault(
Arrays.asList(PipeSourceConstant.EXTRACTOR_KEY, PipeSourceConstant.SOURCE_KEY),
BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName())
.toLowerCase());
final String sourcePluginName =
sourceParameters
.getStringOrDefault(
Arrays.asList(PipeSourceConstant.EXTRACTOR_KEY, PipeSourceConstant.SOURCE_KEY),
BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName())
.toLowerCase();
final boolean isIoTDBSource =
BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName().equals(sourcePluginName)
|| BuiltinPipePlugin.IOTDB_SOURCE.getPipePluginName().equals(sourcePluginName);
final boolean shouldMarkAsGeneralWriteRequest =
sinkParameters.getBooleanOrDefault(
Arrays.asList(
PipeSinkConstant.CONNECTOR_MARK_AS_GENERAL_WRITE_REQUEST_KEY,
PipeSinkConstant.SINK_MARK_AS_GENERAL_WRITE_REQUEST_KEY),
PipeSinkConstant.CONNECTOR_MARK_AS_GENERAL_WRITE_REQUEST_DEFAULT_VALUE);
final boolean shouldMarkAsPipeRequest =
!shouldMarkAsGeneralWriteRequest
&& sinkParameters.getBooleanOrDefault(
Arrays.asList(
PipeSinkConstant.CONNECTOR_MARK_AS_PIPE_REQUEST_KEY,
PipeSinkConstant.SINK_MARK_AS_PIPE_REQUEST_KEY),
PipeSinkConstant.CONNECTOR_MARK_AS_PIPE_REQUEST_DEFAULT_VALUE);

boolean shouldWaitForSchemaBeforeLoad = false;
try {
shouldWaitForSchemaBeforeLoad =
isIoTDBSource
&& shouldMarkAsPipeRequest
&& areOptionsEnabled(sourceParameters, SCHEMA_OPTIONS_REQUIRED_BEFORE_LOAD);
} catch (final IllegalPathException e) {
LOGGER.warn(
DataNodePipeMessages.PIPEDATANODETASKBUILDER_FAILED_TO_PARSE_INCLUSION_AND_EXCLUSION,
e.getMessage(),
e);
}
sinkParameters.addAttribute(
SystemConstant.SINK_WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY,
Boolean.toString(shouldWaitForSchemaBeforeLoad));

final boolean isSourceExternal = !BuiltinPipePlugin.BUILTIN_SOURCES.contains(sourcePluginName);

final String sinkPluginName =
sinkParameters
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -610,34 +610,42 @@ protected String getSenderPort() {
protected TSStatus loadFileV1(final PipeTransferFileSealReqV1 req, final String fileAbsolutePath)
throws IOException {
return isUsingAsyncLoadTsFileStrategy.get()
? loadTsFileAsync(null, Collections.singletonList(fileAbsolutePath))
: loadTsFileSync(null, fileAbsolutePath);
? loadTsFileAsync(null, Collections.singletonList(fileAbsolutePath), false)
: loadTsFileSync(null, fileAbsolutePath, false);
}

@Override
protected TSStatus loadFileV2(
final PipeTransferFileSealReqV2 req, final List<String> fileAbsolutePaths)
throws IOException, IllegalPathException {
return req instanceof PipeTransferTsFileSealWithModReq
// TsFile's absolute path will be the second element
? (isUsingAsyncLoadTsFileStrategy.get()
? loadTsFileAsync(
((PipeTransferTsFileSealWithModReq) req).getDatabaseNameByTsFileName(),
fileAbsolutePaths)
: loadTsFileSync(
((PipeTransferTsFileSealWithModReq) req).getDatabaseNameByTsFileName(),
fileAbsolutePaths.get(req.getFileNames().size() - 1)))
: loadSchemaSnapShot(req.getParameters(), fileAbsolutePaths);
if (!(req instanceof PipeTransferTsFileSealWithModReq)) {
return loadSchemaSnapShot(req.getParameters(), fileAbsolutePaths);
}

final PipeTransferTsFileSealWithModReq tsFileSealReq = (PipeTransferTsFileSealWithModReq) req;
final String databaseName = tsFileSealReq.getDatabaseNameByTsFileName();
final boolean shouldWaitForSchemaBeforeLoad = tsFileSealReq.shouldWaitForSchemaBeforeLoad();
// TsFile's absolute path will be the second element when the request contains a mod file.
return isUsingAsyncLoadTsFileStrategy.get()
? loadTsFileAsync(databaseName, fileAbsolutePaths, shouldWaitForSchemaBeforeLoad)
: loadTsFileSync(
databaseName,
fileAbsolutePaths.get(req.getFileNames().size() - 1),
shouldWaitForSchemaBeforeLoad);
}

private TSStatus loadTsFileAsync(final String dataBaseName, final List<String> absolutePaths)
private TSStatus loadTsFileAsync(
final String dataBaseName,
final List<String> absolutePaths,
final boolean shouldWaitForSchemaBeforeLoad)
throws IOException {
final Map<String, String> loadAttributes =
buildLoadTsFileAttributesForAsync(
dataBaseName,
shouldConvertDataTypeOnTypeMismatch,
validateTsFile.get(),
shouldMarkAsPipeRequest.get());
shouldMarkAsPipeRequest.get(),
shouldWaitForSchemaBeforeLoad);

if (!LoadUtil.loadFilesToActiveDir(loadAttributes, absolutePaths, true)) {
throw new PipeException(DataNodePipeMessages.LOAD_ACTIVE_LISTENING_PIPE_DIR_IS_NOT);
Expand All @@ -650,24 +658,43 @@ static Map<String, String> buildLoadTsFileAttributesForAsync(
final boolean shouldConvertDataTypeOnTypeMismatch,
final boolean validateTsFile,
final boolean shouldMarkAsPipeRequest) {
return buildLoadTsFileAttributesForAsync(
dataBaseName,
shouldConvertDataTypeOnTypeMismatch,
validateTsFile,
shouldMarkAsPipeRequest,
false);
}

static Map<String, String> buildLoadTsFileAttributesForAsync(
final String dataBaseName,
final boolean shouldConvertDataTypeOnTypeMismatch,
final boolean validateTsFile,
final boolean shouldMarkAsPipeRequest,
final boolean shouldWaitForSchemaBeforeLoad) {
return ActiveLoadPathHelper.buildAttributes(
dataBaseName,
LoadTsFileStatement.getDatabaseLevelByTreeDatabase(dataBaseName),
shouldConvertDataTypeOnTypeMismatch,
validateTsFile || shouldConvertDataTypeOnTypeMismatch,
validateTsFile || shouldConvertDataTypeOnTypeMismatch || shouldWaitForSchemaBeforeLoad,
!shouldWaitForSchemaBeforeLoad,
null,
shouldMarkAsPipeRequest,
AuthorityChecker.SUPER_USER);
}

private TSStatus loadTsFileSync(final String dataBaseName, final String fileAbsolutePath)
private TSStatus loadTsFileSync(
final String dataBaseName,
final String fileAbsolutePath,
final boolean shouldWaitForSchemaBeforeLoad)
throws FileNotFoundException {
return executeStatementAndClassifyExceptions(
buildLoadTsFileStatementForSync(
dataBaseName,
fileAbsolutePath,
validateTsFile.get(),
shouldConvertDataTypeOnTypeMismatch));
shouldConvertDataTypeOnTypeMismatch,
shouldWaitForSchemaBeforeLoad));
}

static LoadTsFileStatement buildLoadTsFileStatementForSync(
Expand All @@ -676,10 +703,23 @@ static LoadTsFileStatement buildLoadTsFileStatementForSync(
final boolean validateTsFile,
final boolean shouldConvertDataTypeOnTypeMismatch)
throws FileNotFoundException {
return buildLoadTsFileStatementForSync(
dataBaseName, fileAbsolutePath, validateTsFile, shouldConvertDataTypeOnTypeMismatch, false);
}

static LoadTsFileStatement buildLoadTsFileStatementForSync(
final String dataBaseName,
final String fileAbsolutePath,
final boolean validateTsFile,
final boolean shouldConvertDataTypeOnTypeMismatch,
final boolean shouldWaitForSchemaBeforeLoad)
throws FileNotFoundException {
final LoadTsFileStatement statement = LoadTsFileStatement.createUnchecked(fileAbsolutePath);
statement.setDeleteAfterLoad(true);
statement.setConvertOnTypeMismatch(shouldConvertDataTypeOnTypeMismatch);
statement.setVerifySchema(validateTsFile || shouldConvertDataTypeOnTypeMismatch);
statement.setVerifySchema(
validateTsFile || shouldConvertDataTypeOnTypeMismatch || shouldWaitForSchemaBeforeLoad);
statement.setAutoCreateSchema(!shouldWaitForSchemaBeforeLoad);
statement.setAutoCreateDatabase(
IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled());
statement.setDatabase(dataBaseName);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,19 @@ protected PipeRequestType getPlanType() {
}

protected static final String DATABASE_NAME_KEY_PREFIX = "DATABASE_NAME_";
private static final String WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY = "WAIT_FOR_SCHEMA_BEFORE_LOAD";

public String getDatabaseNameByTsFileName() {
return parameters == null
? null
: parameters.get(generateDatabaseNameWithFileNameKey(fileNames.get(fileNames.size() - 1)));
}

public boolean shouldWaitForSchemaBeforeLoad() {
return parameters != null
&& Boolean.parseBoolean(parameters.get(WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY));
}

protected static String generateDatabaseNameWithFileNameKey(final String fileName) {
return DATABASE_NAME_KEY_PREFIX + fileName;
}
Expand All @@ -69,25 +75,44 @@ public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
final long tsFileLength,
final String dataBaseName)
throws IOException {
return toTPipeTransferReq(
modFileName, modFileLength, tsFileName, tsFileLength, dataBaseName, false);
}

public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
final String modFileName,
final long modFileLength,
final String tsFileName,
final long tsFileLength,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad)
throws IOException {
return (PipeTransferTsFileSealWithModReq)
new PipeTransferTsFileSealWithModReq()
.convertToTPipeTransferReq(
Arrays.asList(modFileName, tsFileName),
Arrays.asList(modFileLength, tsFileLength),
Collections.singletonMap(
generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
final String tsFileName, final long tsFileLength, final String dataBaseName)
throws IOException {
return toTPipeTransferReq(tsFileName, tsFileLength, dataBaseName, false);
}

public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
final String tsFileName,
final long tsFileLength,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad)
throws IOException {
return (PipeTransferTsFileSealWithModReq)
new PipeTransferTsFileSealWithModReq()
.convertToTPipeTransferReq(
Collections.singletonList(tsFileName),
Collections.singletonList(tsFileLength),
Collections.singletonMap(
generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

public static PipeTransferTsFileSealWithModReq fromTPipeTransferReq(final TPipeTransferReq req) {
Expand Down Expand Up @@ -117,23 +142,54 @@ public static byte[] toTPipeTransferBytes(
final long tsFileLength,
final String dataBaseName)
throws IOException {
return toTPipeTransferBytes(
modFileName, modFileLength, tsFileName, tsFileLength, dataBaseName, false);
}

public static byte[] toTPipeTransferBytes(
final String modFileName,
final long modFileLength,
final String tsFileName,
final long tsFileLength,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad)
throws IOException {
return new PipeTransferTsFileSealWithModReq()
.convertToTPipeTransferSnapshotSealBytes(
Arrays.asList(modFileName, tsFileName),
Arrays.asList(modFileLength, tsFileLength),
Collections.singletonMap(
generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

public static byte[] toTPipeTransferBytes(
final String tsFileName, final long tsFileLength, final String dataBaseName)
throws IOException {
return toTPipeTransferBytes(tsFileName, tsFileLength, dataBaseName, false);
}

public static byte[] toTPipeTransferBytes(
final String tsFileName,
final long tsFileLength,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad)
throws IOException {
return new PipeTransferTsFileSealWithModReq()
.convertToTPipeTransferSnapshotSealBytes(
Collections.singletonList(tsFileName),
Collections.singletonList(tsFileLength),
Collections.singletonMap(
generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

private static HashMap<String, String> generateParameters(
final String tsFileName,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad) {
final HashMap<String, String> parameters = new HashMap<>();
parameters.put(generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName);
if (shouldWaitForSchemaBeforeLoad) {
parameters.put(WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY, Boolean.TRUE.toString());
}
return parameters;
}

/////////////////////////////// Object ///////////////////////////////
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -475,7 +475,12 @@ private void doTransfer(
if (!sendWeighted(
socket,
PipeTransferTsFileSealWithModReq.toTPipeTransferBytes(
modFile.getName(), modFile.length(), tsFile.getName(), tsFile.length(), dataBaseName),
modFile.getName(),
modFile.length(),
tsFile.getName(),
tsFile.length(),
dataBaseName,
shouldWaitForSchemaBeforeLoad),
pipe2WeightMap)) {
receiverStatusHandler.handle(
new TSStatus(TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode())
Expand All @@ -490,7 +495,7 @@ private void doTransfer(
if (!sendWeighted(
socket,
PipeTransferTsFileSealWithModReq.toTPipeTransferBytes(
tsFile.getName(), tsFile.length(), dataBaseName),
tsFile.getName(), tsFile.length(), dataBaseName, shouldWaitForSchemaBeforeLoad),
pipe2WeightMap)) {
receiverStatusHandler.handle(
new TSStatus(TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -206,9 +206,13 @@ public void transfer(
modFile.length(),
tsFile.getName(),
tsFile.length(),
dataBaseName)
dataBaseName,
sink.shouldWaitForSchemaBeforeLoad())
: PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
tsFile.getName(), tsFile.length(), dataBaseName);
tsFile.getName(),
tsFile.length(),
dataBaseName,
sink.shouldWaitForSchemaBeforeLoad());
final TPipeTransferReq req = sink.compressIfNeeded(uncompressedReq);

pipeName2WeightMap.forEach(
Expand Down
Loading
Loading