Skip to content
Merged
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 @@ -392,6 +392,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", "count(devices),", Collections.singleton("3,"));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,11 +56,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 @@ -261,13 +270,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(
"PipeDataNodeTaskBuilder failed to parse 'inclusion' and 'exclusion' parameters: {}",
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 @@ -503,32 +503,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 {
if (req instanceof PipeTransferTsFileSealWithModReq) {
final String dataBaseName =
((PipeTransferTsFileSealWithModReq) req).getDatabaseNameByTsFileName();
return isUsingAsyncLoadTsFileStrategy.get()
? loadTsFileAsync(dataBaseName, fileAbsolutePaths)
: loadTsFileSync(dataBaseName, fileAbsolutePaths.get(req.getFileNames().size() - 1));
if (!(req instanceof PipeTransferTsFileSealWithModReq)) {
return loadSchemaSnapShot(req.getParameters(), fileAbsolutePaths);
}
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 (!ActiveLoadUtil.loadFilesToActiveDir(loadAttributes, absolutePaths, true)) {
throw new PipeException("Load active listening pipe dir is not set.");
}
Expand All @@ -540,23 +550,42 @@ 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);
}

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 @@ -565,10 +594,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 @@ -41,6 +41,7 @@ protected PipeRequestType getPlanType() {
}

private 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 getParameters() == null
Expand All @@ -50,23 +51,35 @@ public String getDatabaseNameByTsFileName() {
generateDatabaseNameWithFileNameKey(getFileNames().get(getFileNames().size() - 1)));
}

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

private static String generateDatabaseNameWithFileNameKey(final String fileName) {
return DATABASE_NAME_KEY_PREFIX + fileName;
}

private static Map<String, String> generateDatabaseNameParameter(
final String tsFileName, final String dataBaseName) {
return dataBaseName == null
? new HashMap<>()
: Collections.singletonMap(generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName);
private static Map<String, String> generateParameters(
final String tsFileName,
final String dataBaseName,
final boolean shouldWaitForSchemaBeforeLoad) {
final Map<String, String> parameters = new HashMap<>();
if (dataBaseName != null) {
parameters.put(generateDatabaseNameWithFileNameKey(tsFileName), dataBaseName);
}
if (shouldWaitForSchemaBeforeLoad) {
parameters.put(WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY, Boolean.TRUE.toString());
}
return parameters;
}

/////////////////////////////// Thrift ///////////////////////////////

public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
String modFileName, long modFileLength, String tsFileName, long tsFileLength)
throws IOException {
return toTPipeTransferReq(modFileName, modFileLength, tsFileName, tsFileLength, null);
return toTPipeTransferReq(modFileName, modFileLength, tsFileName, tsFileLength, null, false);
}

public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
Expand All @@ -76,23 +89,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),
generateDatabaseNameParameter(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),
generateDatabaseNameParameter(tsFileName, dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

public static PipeTransferTsFileSealWithModReq fromTPipeTransferReq(TPipeTransferReq req) {
Expand All @@ -105,7 +139,7 @@ public static PipeTransferTsFileSealWithModReq fromTPipeTransferReq(TPipeTransfe
public static byte[] toTPipeTransferBytes(
String modFileName, long modFileLength, String tsFileName, long tsFileLength)
throws IOException {
return toTPipeTransferBytes(modFileName, modFileLength, tsFileName, tsFileLength, null);
return toTPipeTransferBytes(modFileName, modFileLength, tsFileName, tsFileLength, null, false);
}

public static byte[] toTPipeTransferBytes(
Expand All @@ -115,21 +149,42 @@ 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),
generateDatabaseNameParameter(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),
generateDatabaseNameParameter(tsFileName, dataBaseName));
generateParameters(tsFileName, dataBaseName, shouldWaitForSchemaBeforeLoad));
}

/////////////////////////////// Object ///////////////////////////////
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -406,7 +406,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 @@ -420,10 +425,10 @@ private void doTransfer(
transferFilePieces(pipe2WeightMap, tsFile, socket, false);
if (!sendWeighted(
socket,
dataBaseName == null
dataBaseName == null && !shouldWaitForSchemaBeforeLoad
? PipeTransferTsFileSealReq.toTPipeTransferBytes(tsFile.getName(), tsFile.length())
: 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
Loading
Loading