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 @@ -105,6 +105,10 @@ public synchronized String register(
}

for (int connectorIndex = 0; connectorIndex < sinkNum; connectorIndex++) {
final String taskID =
String.format(
"%s_%s_%s", attributeSortedString, environment.getCreationTime(), connectorIndex);
environment.setSinkTaskId(taskID);
final PipeConnector pipeConnector =
isDataRegionSink
? PipeDataNodeAgent.plugin().dataRegion().reflectSink(pipeSinkParameters)
Expand Down Expand Up @@ -135,9 +139,7 @@ public synchronized String register(
// 2. Construct PipeConnectorSubtaskLifeCycle to manage PipeConnectorSubtask's life cycle
final PipeSinkSubtask pipeSinkSubtask =
new PipeSinkSubtask(
String.format(
"%s_%s_%s",
attributeSortedString, environment.getCreationTime(), connectorIndex),
taskID,
environment.getCreationTime(),
attributeSortedString,
connectorIndex,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@
import org.apache.iotdb.db.queryengine.plan.statement.pipe.PipeEnrichedStatement;
import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
import org.apache.iotdb.db.storageengine.load.active.ActiveLoadUtil;
import org.apache.iotdb.db.storageengine.load.converter.PipeTsFileConversionTaskManager;
import org.apache.iotdb.db.storageengine.rescon.disk.FolderManager;
import org.apache.iotdb.db.storageengine.rescon.disk.strategy.DirectoryStrategyType;
import org.apache.iotdb.db.tools.schema.SRStatementGenerator;
Expand Down Expand Up @@ -503,7 +504,7 @@ protected String getSenderPort() {
protected TSStatus loadFileV1(final PipeTransferFileSealReqV1 req, final String fileAbsolutePath)
throws IOException {
return isUsingAsyncLoadTsFileStrategy.get()
? loadTsFileAsync(null, Collections.singletonList(fileAbsolutePath), false)
? loadTsFileAsync(null, Collections.singletonList(fileAbsolutePath), false, null)
: loadTsFileSync(null, fileAbsolutePath, false);
}

Expand All @@ -518,33 +519,126 @@ protected TSStatus loadFileV2(
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);
final String taskId = tsFileSealReq.getConversionTaskId();
final boolean asyncLoadOnTypeMismatch = tsFileSealReq.shouldAsyncLoadOnTypeMismatch();
final TSStatus duplicateStatus =
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(
taskId, asyncLoadOnTypeMismatch);
if (duplicateStatus != null) {
return duplicateStatus;
}
PipeTsFileConversionTaskManager.enter(taskId);
try {
final TSStatus status;
if (isUsingAsyncLoadTsFileStrategy.get()) {
status =
loadTsFileAsync(dataBaseName, fileAbsolutePaths, shouldWaitForSchemaBeforeLoad, taskId);
} else {
PipeTsFileConversionTaskManager.markRunning(taskId);
status =
loadTsFileSync(
dataBaseName,
fileAbsolutePaths.get(req.getFileNames().size() - 1),
shouldWaitForSchemaBeforeLoad);
}

if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
if (!isUsingAsyncLoadTsFileStrategy.get()) {
PipeTsFileConversionTaskManager.markSuccess(taskId);
}
return status;
}

if (shouldTakeOverToAsyncLoad(
status,
isUsingAsyncLoadTsFileStrategy.get(),
shouldConvertDataTypeOnTypeMismatch,
asyncLoadOnTypeMismatch,
PipeTsFileConversionTaskManager.isTypeMismatchDetected(taskId))) {
PipeTsFileConversionTaskManager.clearCurrentContext();
PipeTsFileConversionTaskManager.prepareForActiveLoad(taskId);
try {
final TSStatus takeoverStatus =
loadTsFileAsync(
dataBaseName, fileAbsolutePaths, shouldWaitForSchemaBeforeLoad, taskId);
if (takeoverStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return takeoverStatus;
}
} catch (final Exception ignored) {
// The sender retries the same stable task id after a failed durable handoff.
}
PipeTsFileConversionTaskManager.markRetryable(taskId, status);
return status;
}

PipeTsFileConversionTaskManager.markRetryable(taskId, status);
return status;
} catch (final Exception e) {
final TSStatus status =
new TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode()).setMessage(e.getMessage());
PipeTsFileConversionTaskManager.markRetryable(taskId, status);
throw e;
} finally {
PipeTsFileConversionTaskManager.leave();
}
}

static boolean shouldTakeOverToAsyncLoad(
final TSStatus status,
final boolean usingAsyncLoadStrategy,
final boolean shouldConvertOnTypeMismatch,
final boolean asyncLoadOnTypeMismatch,
final boolean typeMismatchDetected) {
return !usingAsyncLoadStrategy
&& shouldConvertOnTypeMismatch
&& asyncLoadOnTypeMismatch
&& !isLoadTemporarilyUnavailable(status)
&& typeMismatchDetected;
}

private static boolean isLoadTemporarilyUnavailable(final TSStatus status) {
return status != null
&& (status.getCode() == TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()
|| status.getCode()
== TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode());
}

private TSStatus loadTsFileAsync(
final String dataBaseName,
final List<String> absolutePaths,
final boolean shouldWaitForSchemaBeforeLoad)
final boolean shouldWaitForSchemaBeforeLoad,
final String conversionTaskId)
throws IOException {
final Map<String, String> loadAttributes =
buildLoadTsFileAttributesForAsync(
dataBaseName,
shouldConvertDataTypeOnTypeMismatch,
validateTsFile.get(),
shouldMarkAsPipeRequest.get(),
shouldWaitForSchemaBeforeLoad);
shouldWaitForSchemaBeforeLoad,
conversionTaskId);
if (!ActiveLoadUtil.loadFilesToActiveDir(loadAttributes, absolutePaths, true)) {
throw new PipeException("Load active listening pipe dir is not set.");
}
PipeTsFileConversionTaskManager.markReceiverOwned(conversionTaskId);
return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
}

@Override
protected boolean shouldDeleteSealedFilesOnFailure(
final PipeTransferFileSealReqV2 req, final TSStatus loadStatus) {
if (!(req instanceof PipeTransferTsFileSealWithModReq)) {
return true;
}
final String taskId = ((PipeTransferTsFileSealWithModReq) req).getConversionTaskId();
final PipeTsFileConversionTaskManager.Task task = PipeTsFileConversionTaskManager.get(taskId);
return task == null
|| task.isReceiverOwned()
|| (task.getState() != PipeTsFileConversionTaskManager.State.PENDING
&& task.getState() != PipeTsFileConversionTaskManager.State.RUNNING
&& task.getState() != PipeTsFileConversionTaskManager.State.PAUSED);
}

static Map<String, String> buildLoadTsFileAttributesForAsync(
final String dataBaseName,
final boolean shouldConvertDataTypeOnTypeMismatch,
Expand All @@ -555,7 +649,8 @@ static Map<String, String> buildLoadTsFileAttributesForAsync(
shouldConvertDataTypeOnTypeMismatch,
validateTsFile,
shouldMarkAsPipeRequest,
false);
false,
null);
}

static Map<String, String> buildLoadTsFileAttributesForAsync(
Expand All @@ -564,6 +659,22 @@ static Map<String, String> buildLoadTsFileAttributesForAsync(
final boolean validateTsFile,
final boolean shouldMarkAsPipeRequest,
final boolean shouldWaitForSchemaBeforeLoad) {
return buildLoadTsFileAttributesForAsync(
dataBaseName,
shouldConvertDataTypeOnTypeMismatch,
validateTsFile,
shouldMarkAsPipeRequest,
shouldWaitForSchemaBeforeLoad,
null);
}

static Map<String, String> buildLoadTsFileAttributesForAsync(
final String dataBaseName,
final boolean shouldConvertDataTypeOnTypeMismatch,
final boolean validateTsFile,
final boolean shouldMarkAsPipeRequest,
final boolean shouldWaitForSchemaBeforeLoad,
final String conversionTaskId) {
return ActiveLoadPathHelper.buildAttributes(
dataBaseName,
LoadTsFileStatement.getDatabaseLevelByTreeDatabase(dataBaseName),
Expand All @@ -572,7 +683,8 @@ static Map<String, String> buildLoadTsFileAttributesForAsync(
!shouldWaitForSchemaBeforeLoad,
null,
shouldMarkAsPipeRequest,
AuthorityChecker.SUPER_USER);
AuthorityChecker.SUPER_USER,
conversionTaskId);
}

private TSStatus loadTsFileSync(
Expand Down
Loading
Loading