This is an automated email from the ASF dual-hosted git repository.

Caideyipi pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new 6ac48ee161e [To dev/1.3] [Pipe] Support receiver-side async takeover 
for TsFile conversion (1.3) (#18489)
6ac48ee161e is described below

commit 6ac48ee161eb7ee2d445246deee7f1086b837466
Author: Caideyipi <[email protected]>
AuthorDate: Thu Aug 20 09:50:21 2026 +0800

    [To dev/1.3] [Pipe] Support receiver-side async takeover for TsFile 
conversion (1.3) (#18489)
    
    * [Pipe] Support receiver-side async takeover for TsFile conversion (#18421)
    
    * [Pipe] Support async TsFile conversion takeover
    
    * [Pipe] Retry receiver-side TsFile load failures
    
    * [Pipe] Address async conversion review comments
    
    * test(load): fix active load scanner mock verification
    
    (cherry picked from commit 4eff4c1654c74cb5f4cc85a76fa95c4af00987d1)
    
    * [To dev/1.3] Fix cherry-pick test regressions
    
    * [Pipe] Fix conversion task ID collision after leader switch (#18494)
---
 .../task/subtask/sink/PipeSinkSubtaskManager.java  |   8 +-
 .../protocol/thrift/IoTDBDataNodeReceiver.java     | 136 +++++-
 .../request/PipeTransferTsFileSealWithModReq.java  | 198 ++++++++
 .../protocol/airgap/IoTDBDataRegionAirGapSink.java |  33 +-
 .../thrift/async/IoTDBDataRegionAsyncSink.java     |   6 +-
 .../async/handler/PipeTransferTsFileHandler.java   |  75 ++-
 .../thrift/sync/IoTDBDataRegionSyncSink.java       |  48 +-
 .../plan/analyze/load/LoadTsFileAnalyzer.java      |   4 +
 .../load/active/ActiveLoadDirScanner.java          |  20 +-
 .../load/active/ActiveLoadPathHelper.java          |  76 +++
 .../load/active/ActiveLoadPendingQueue.java        |  22 +-
 .../load/active/ActiveLoadTsFileLoader.java        |  57 ++-
 .../storageengine/load/active/ActiveLoadUtil.java  | 112 ++++-
 ...eeStatementDataTypeConvertExecutionVisitor.java | 204 +++++---
 .../converter/LoadTreeTsFileTabletIterator.java    |  14 +-
 .../converter/LoadTsFileDataTypeConverter.java     |   2 +-
 .../converter/PipeTsFileConversionTaskManager.java | 541 +++++++++++++++++++++
 .../pipe/sink/PipeDataNodeThriftRequestTest.java   |  30 ++
 .../airgap/IoTDBDataRegionAirGapSinkTest.java      |   3 +-
 .../load/active/ActiveLoadTsFileLoaderTest.java    |  13 +
 .../load/active/ActiveLoadUtilTest.java            |  99 +++-
 .../PipeTsFileConversionTaskManagerTest.java       | 381 +++++++++++++++
 .../pipe/config/constant/PipeSinkConstant.java     |   7 +
 .../plugin/env/PipeTaskSinkRuntimeEnvironment.java |   9 +
 .../commons/pipe/receiver/IoTDBFileReceiver.java   |  16 +-
 .../commons/pipe/sink/protocol/IoTDBSink.java      |  25 +
 26 files changed, 1990 insertions(+), 149 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
index 6307dbba648..072b31e8f8f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
@@ -105,6 +105,10 @@ public class PipeSinkSubtaskManager {
       }
 
       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)
@@ -135,9 +139,7 @@ public class PipeSinkSubtaskManager {
         // 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,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
index 9439d81f6bc..3fbda2083dd 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
@@ -90,6 +90,7 @@ import 
org.apache.iotdb.db.queryengine.plan.statement.metadata.DatabaseSchemaSta
 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;
@@ -519,7 +520,7 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
   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);
   }
 
@@ -534,19 +535,95 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
     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(
@@ -554,13 +631,30 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
             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,
@@ -571,7 +665,8 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
         shouldConvertDataTypeOnTypeMismatch,
         validateTsFile,
         shouldMarkAsPipeRequest,
-        false);
+        false,
+        null);
   }
 
   static Map<String, String> buildLoadTsFileAttributesForAsync(
@@ -580,6 +675,22 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
       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),
@@ -588,7 +699,8 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
         !shouldWaitForSchemaBeforeLoad,
         null,
         shouldMarkAsPipeRequest,
-        AuthorityChecker.SUPER_USER);
+        AuthorityChecker.SUPER_USER,
+        conversionTaskId);
   }
 
   private TSStatus loadTsFileSync(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
index 63b3dbfd296..c9da41dbd1c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTsFileSealWithModReq.java
@@ -19,15 +19,21 @@
 
 package org.apache.iotdb.db.pipe.sink.payload.evolvable.request;
 
+import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.request.PipeRequestType;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.request.PipeTransferFileSealReqV2;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 
 import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
+import java.util.Comparator;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
+import java.util.UUID;
 
 public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2 {
 
@@ -42,6 +48,9 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
 
   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 static final String CONVERSION_TASK_ID_KEY = "CONVERSION_TASK_ID";
+  public static final String ASYNC_LOAD_ON_TYPE_MISMATCH_KEY = 
"ASYNC_LOAD_ON_TYPE_MISMATCH";
+  private static final String UNSUPPORTED_PROGRESS_INDEX = 
"unsupported-progress-index";
 
   public String getDatabaseNameByTsFileName() {
     return getParameters() == null
@@ -56,6 +65,94 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
         && 
Boolean.parseBoolean(getParameters().get(WAIT_FOR_SCHEMA_BEFORE_LOAD_KEY));
   }
 
+  public String getConversionTaskId() {
+    return getParameters() == null ? null : 
getParameters().get(CONVERSION_TASK_ID_KEY);
+  }
+
+  public boolean shouldAsyncLoadOnTypeMismatch() {
+    if (getParameters() == null) {
+      return true;
+    }
+    final String value = getParameters().get(ASYNC_LOAD_ON_TYPE_MISMATCH_KEY);
+    return value == null || Boolean.parseBoolean(value);
+  }
+
+  public PipeTransferTsFileSealWithModReq setConversionTaskInfo(
+      final String conversionTaskId, final boolean 
shouldAsyncLoadOnTypeMismatch)
+      throws IOException {
+    final Map<String, String> parameters =
+        getParameters() == null ? new HashMap<>() : new 
HashMap<>(getParameters());
+    if (conversionTaskId != null) {
+      parameters.put(CONVERSION_TASK_ID_KEY, conversionTaskId);
+    }
+    parameters.put(
+        ASYNC_LOAD_ON_TYPE_MISMATCH_KEY, 
Boolean.toString(shouldAsyncLoadOnTypeMismatch));
+    return (PipeTransferTsFileSealWithModReq)
+        convertToTPipeTransferReq(getFileNames(), getFileLengths(), 
parameters);
+  }
+
+  public static String generateConversionTaskId(
+      final String sinkTaskId,
+      final Iterable<? extends EnrichedEvent> events,
+      final String databaseName,
+      final int outputIndex) {
+    return generateConversionTaskId(sinkTaskId, events, databaseName, 
outputIndex, false);
+  }
+
+  public static String generateConversionTaskId(
+      final String sinkTaskId,
+      final Iterable<? extends EnrichedEvent> events,
+      final String databaseName,
+      final int outputIndex,
+      final boolean hasModFile) {
+    final StringBuilder stableKey = new StringBuilder();
+    appendStablePart(stableKey, sinkTaskId);
+    appendStablePart(stableKey, databaseName);
+    appendStablePart(stableKey, Integer.toString(outputIndex));
+    appendStablePart(stableKey, Boolean.toString(hasModFile));
+
+    final List<String> eventIdentities = new ArrayList<>();
+    if (events != null) {
+      for (final EnrichedEvent event : events) {
+        if (event == null) {
+          continue;
+        }
+        final StringBuilder eventIdentity = new StringBuilder();
+        appendStablePart(eventIdentity, event.getClass().getName());
+        appendStablePart(eventIdentity, event.getPipeName());
+        appendStablePart(eventIdentity, 
Long.toString(event.getCreationTime()));
+        appendStablePart(eventIdentity, Integer.toString(event.getRegionId()));
+        if (event.getCommitterKey() != null) {
+          appendStablePart(eventIdentity, 
event.getCommitterKey().getPipeName());
+          appendStablePart(eventIdentity, 
Long.toString(event.getCommitterKey().getCreationTime()));
+          appendStablePart(eventIdentity, 
Integer.toString(event.getCommitterKey().getRegionId()));
+        }
+        final List<Long> commitIds = new ArrayList<>();
+        if (event.getCommitIds() != null) {
+          commitIds.addAll(event.getCommitIds());
+        }
+        commitIds.sort(Comparator.naturalOrder());
+        commitIds.forEach(id -> appendStablePart(eventIdentity, 
Long.toString(id)));
+        // Commit ids are local to a DataNode and may collide after a leader 
change.
+        try {
+          appendStablePart(eventIdentity, 
String.valueOf(event.getProgressIndex()));
+        } catch (final UnsupportedOperationException e) {
+          appendStablePart(eventIdentity, UNSUPPORTED_PROGRESS_INDEX);
+        }
+        eventIdentities.add(eventIdentity.toString());
+      }
+    }
+    eventIdentities.sort(Comparator.naturalOrder());
+    appendStablePart(stableKey, Integer.toString(eventIdentities.size()));
+    eventIdentities.forEach(identity -> appendStablePart(stableKey, identity));
+    return 
UUID.nameUUIDFromBytes(stableKey.toString().getBytes(StandardCharsets.UTF_8)).toString();
+  }
+
+  private static void appendStablePart(final StringBuilder builder, final 
String value) {
+    final String normalizedValue = value == null ? "" : value;
+    
builder.append(normalizedValue.length()).append(':').append(normalizedValue).append('\0');
+  }
+
   private static String generateDatabaseNameWithFileNameKey(final String 
fileName) {
     return DATABASE_NAME_KEY_PREFIX + fileName;
   }
@@ -74,6 +171,21 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
     return parameters;
   }
 
+  private static Map<String, String> generateParameters(
+      final String tsFileName,
+      final String dataBaseName,
+      final boolean shouldWaitForSchemaBeforeLoad,
+      final String conversionTaskId,
+      final boolean asyncLoadOnTypeMismatch) {
+    final Map<String, String> parameters =
+        generateParameters(tsFileName, dataBaseName, 
shouldWaitForSchemaBeforeLoad);
+    if (conversionTaskId != null) {
+      parameters.put(CONVERSION_TASK_ID_KEY, conversionTaskId);
+    }
+    parameters.put(ASYNC_LOAD_ON_TYPE_MISMATCH_KEY, 
Boolean.toString(asyncLoadOnTypeMismatch));
+    return parameters;
+  }
+
   /////////////////////////////// Thrift ///////////////////////////////
 
   public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
@@ -109,6 +221,29 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
                 generateParameters(tsFileName, dataBaseName, 
shouldWaitForSchemaBeforeLoad));
   }
 
+  public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
+      final String modFileName,
+      final long modFileLength,
+      final String tsFileName,
+      final long tsFileLength,
+      final String dataBaseName,
+      final boolean shouldWaitForSchemaBeforeLoad,
+      final String conversionTaskId,
+      final boolean asyncLoadOnTypeMismatch)
+      throws IOException {
+    return (PipeTransferTsFileSealWithModReq)
+        new PipeTransferTsFileSealWithModReq()
+            .convertToTPipeTransferReq(
+                Arrays.asList(modFileName, tsFileName),
+                Arrays.asList(modFileLength, tsFileLength),
+                generateParameters(
+                    tsFileName,
+                    dataBaseName,
+                    shouldWaitForSchemaBeforeLoad,
+                    conversionTaskId,
+                    asyncLoadOnTypeMismatch));
+  }
+
   public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
       final String tsFileName, final long tsFileLength, final String 
dataBaseName)
       throws IOException {
@@ -129,6 +264,27 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
                 generateParameters(tsFileName, dataBaseName, 
shouldWaitForSchemaBeforeLoad));
   }
 
+  public static PipeTransferTsFileSealWithModReq toTPipeTransferReq(
+      final String tsFileName,
+      final long tsFileLength,
+      final String dataBaseName,
+      final boolean shouldWaitForSchemaBeforeLoad,
+      final String conversionTaskId,
+      final boolean asyncLoadOnTypeMismatch)
+      throws IOException {
+    return (PipeTransferTsFileSealWithModReq)
+        new PipeTransferTsFileSealWithModReq()
+            .convertToTPipeTransferReq(
+                Collections.singletonList(tsFileName),
+                Collections.singletonList(tsFileLength),
+                generateParameters(
+                    tsFileName,
+                    dataBaseName,
+                    shouldWaitForSchemaBeforeLoad,
+                    conversionTaskId,
+                    asyncLoadOnTypeMismatch));
+  }
+
   public static PipeTransferTsFileSealWithModReq 
fromTPipeTransferReq(TPipeTransferReq req) {
     return (PipeTransferTsFileSealWithModReq)
         new 
PipeTransferTsFileSealWithModReq().translateFromTPipeTransferReq(req);
@@ -168,6 +324,28 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
             generateParameters(tsFileName, dataBaseName, 
shouldWaitForSchemaBeforeLoad));
   }
 
+  public static byte[] toTPipeTransferBytes(
+      final String modFileName,
+      final long modFileLength,
+      final String tsFileName,
+      final long tsFileLength,
+      final String dataBaseName,
+      final boolean shouldWaitForSchemaBeforeLoad,
+      final String conversionTaskId,
+      final boolean asyncLoadOnTypeMismatch)
+      throws IOException {
+    return new PipeTransferTsFileSealWithModReq()
+        .convertToTPipeTransferSnapshotSealBytes(
+            Arrays.asList(modFileName, tsFileName),
+            Arrays.asList(modFileLength, tsFileLength),
+            generateParameters(
+                tsFileName,
+                dataBaseName,
+                shouldWaitForSchemaBeforeLoad,
+                conversionTaskId,
+                asyncLoadOnTypeMismatch));
+  }
+
   public static byte[] toTPipeTransferBytes(
       final String tsFileName, final long tsFileLength, final String 
dataBaseName)
       throws IOException {
@@ -187,6 +365,26 @@ public class PipeTransferTsFileSealWithModReq extends 
PipeTransferFileSealReqV2
             generateParameters(tsFileName, dataBaseName, 
shouldWaitForSchemaBeforeLoad));
   }
 
+  public static byte[] toTPipeTransferBytes(
+      final String tsFileName,
+      final long tsFileLength,
+      final String dataBaseName,
+      final boolean shouldWaitForSchemaBeforeLoad,
+      final String conversionTaskId,
+      final boolean asyncLoadOnTypeMismatch)
+      throws IOException {
+    return new PipeTransferTsFileSealWithModReq()
+        .convertToTPipeTransferSnapshotSealBytes(
+            Collections.singletonList(tsFileName),
+            Collections.singletonList(tsFileLength),
+            generateParameters(
+                tsFileName,
+                dataBaseName,
+                shouldWaitForSchemaBeforeLoad,
+                conversionTaskId,
+                asyncLoadOnTypeMismatch));
+  }
+
   /////////////////////////////// Object ///////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSink.java
index 27564aa1505..535f2b2ed41 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSink.java
@@ -255,10 +255,13 @@ public class IoTDBDataRegionAirGapSink extends 
IoTDBDataNodeAirGapSink {
       throws IOException, WriteProcessException {
     final List<File> sealedFiles = batchToTransfer.sealTsFiles();
     final Map<Pair<String, Long>, Double> pipe2WeightMap = 
batchToTransfer.deepCopyPipe2WeightMap();
+    final List<EnrichedEvent> events = batchToTransfer.deepCopyEvents();
 
     try {
-      for (final File tsFile : sealedFiles) {
-        doTransfer(pipe2WeightMap, socket, tsFile, null, null, 
tsFile.getName());
+      for (int outputIndex = 0; outputIndex < sealedFiles.size(); 
outputIndex++) {
+        final File tsFile = sealedFiles.get(outputIndex);
+        doTransfer(
+            pipe2WeightMap, socket, tsFile, null, null, tsFile.getName(), 
events, outputIndex);
       }
     } finally {
       for (final File tsFile : sealedFiles) {
@@ -387,7 +390,9 @@ public class IoTDBDataRegionAirGapSink extends 
IoTDBDataNodeAirGapSink {
             ? pipeTsFileInsertionEvent.getModFile()
             : null,
         pipeTsFileInsertionEvent.getDatabaseName(),
-        pipeTsFileInsertionEvent.toString());
+        pipeTsFileInsertionEvent.toString(),
+        Collections.singletonList(pipeTsFileInsertionEvent),
+        0);
   }
 
   private void doTransfer(
@@ -396,9 +401,16 @@ public class IoTDBDataRegionAirGapSink extends 
IoTDBDataNodeAirGapSink {
       final File tsFile,
       final File modFile,
       final String dataBaseName,
-      final String receiverStatusContext)
+      final String receiverStatusContext,
+      final Iterable<? extends EnrichedEvent> events,
+      final int outputIndex)
       throws PipeException, IOException {
     final String errorMessage = String.format("Seal file %s error. Socket 
%s.", tsFile, socket);
+    final String conversionTaskId =
+        shouldAsyncLoadTsFileOnTypeMismatch
+            ? PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+                getSinkTaskId(), events, dataBaseName, outputIndex, 
Objects.nonNull(modFile))
+            : null;
 
     if (Objects.nonNull(modFile)) {
       transferFilePieces(pipe2WeightMap, modFile, socket, true);
@@ -411,7 +423,9 @@ public class IoTDBDataRegionAirGapSink extends 
IoTDBDataNodeAirGapSink {
               tsFile.getName(),
               tsFile.length(),
               dataBaseName,
-              shouldWaitForSchemaBeforeLoad),
+              shouldWaitForSchemaBeforeLoad,
+              conversionTaskId,
+              shouldAsyncLoadTsFileOnTypeMismatch),
           pipe2WeightMap)) {
         receiverStatusHandler.handle(
             new 
TSStatus(TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode())
@@ -425,10 +439,15 @@ public class IoTDBDataRegionAirGapSink extends 
IoTDBDataNodeAirGapSink {
       transferFilePieces(pipe2WeightMap, tsFile, socket, false);
       if (!sendWeighted(
           socket,
-          dataBaseName == null && !shouldWaitForSchemaBeforeLoad
+          conversionTaskId == null && dataBaseName == null && 
!shouldWaitForSchemaBeforeLoad
               ? 
PipeTransferTsFileSealReq.toTPipeTransferBytes(tsFile.getName(), 
tsFile.length())
               : PipeTransferTsFileSealWithModReq.toTPipeTransferBytes(
-                  tsFile.getName(), tsFile.length(), dataBaseName, 
shouldWaitForSchemaBeforeLoad),
+                  tsFile.getName(),
+                  tsFile.length(),
+                  dataBaseName,
+                  shouldWaitForSchemaBeforeLoad,
+                  conversionTaskId,
+                  shouldAsyncLoadTsFileOnTypeMismatch),
           pipe2WeightMap)) {
         receiverStatusHandler.handle(
             new 
TSStatus(TSStatusCode.PIPE_RECEIVER_USER_CONFLICT_EXCEPTION.getStatusCode())
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
index 52d22ac9774..cd780d963d0 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/IoTDBDataRegionAsyncSink.java
@@ -260,7 +260,8 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
 
       int transferredFileCount = 0;
       try {
-        for (final File sealedFile : sealedFiles) {
+        for (int outputIndex = 0; outputIndex < sealedFiles.size(); 
outputIndex++) {
+          final File sealedFile = sealedFiles.get(outputIndex);
           transfer(
               new PipeTransferTsFileHandler(
                   this,
@@ -271,7 +272,8 @@ public class IoTDBDataRegionAsyncSink extends IoTDBSink {
                   sealedFile,
                   null,
                   false,
-                  null));
+                  null,
+                  outputIndex));
           transferredFileCount++;
         }
       } catch (final Exception e) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java
index d5a10adf75e..8e13e3171c3 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/async/handler/PipeTransferTsFileHandler.java
@@ -80,6 +80,7 @@ public class PipeTransferTsFileHandler extends 
PipeTransferTrackableHandler {
 
   private final boolean transferMod;
   private final String dataBaseName;
+  private final String conversionTaskId;
 
   private final int readFileBufferSize;
   private PipeTsFileMemoryBlock memoryBlock;
@@ -103,6 +104,31 @@ public class PipeTransferTsFileHandler extends 
PipeTransferTrackableHandler {
       final boolean transferMod,
       final String dataBaseName)
       throws InterruptedException {
+    this(
+        connector,
+        pipeName2WeightMap,
+        events,
+        eventsReferenceCount,
+        eventsHadBeenAddedToRetryQueue,
+        tsFile,
+        modFile,
+        transferMod,
+        dataBaseName,
+        0);
+  }
+
+  public PipeTransferTsFileHandler(
+      final IoTDBDataRegionAsyncSink connector,
+      final Map<Pair<String, Long>, Double> pipeName2WeightMap,
+      final List<EnrichedEvent> events,
+      final AtomicInteger eventsReferenceCount,
+      final AtomicBoolean eventsHadBeenAddedToRetryQueue,
+      final File tsFile,
+      final File modFile,
+      final boolean transferMod,
+      final String dataBaseName,
+      final int outputIndex)
+      throws InterruptedException {
     super(connector);
 
     this.pipeName2WeightMap = pipeName2WeightMap;
@@ -115,6 +141,11 @@ public class PipeTransferTsFileHandler extends 
PipeTransferTrackableHandler {
     this.modFile = modFile;
     this.transferMod = transferMod;
     this.dataBaseName = dataBaseName;
+    conversionTaskId =
+        connector.shouldAsyncLoadTsFileOnTypeMismatch()
+            ? PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+                connector.getSinkTaskId(), events, dataBaseName, outputIndex, 
transferMod)
+            : null;
     currentFile = transferMod ? modFile : tsFile;
 
     // NOTE: Waiting for resource enough for slicing here may cause deadlock!
@@ -189,23 +220,33 @@ public class PipeTransferTsFileHandler extends 
PipeTransferTrackableHandler {
       } else if (currentFile == tsFile) {
         isSealSignalSent.set(true);
 
-        final TPipeTransferReq uncompressedReq =
-            transferMod
-                ? PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
-                    modFile.getName(),
-                    modFile.length(),
-                    tsFile.getName(),
-                    tsFile.length(),
-                    dataBaseName,
-                    sink.shouldWaitForSchemaBeforeLoad())
-                : dataBaseName == null && !sink.shouldWaitForSchemaBeforeLoad()
-                    ? PipeTransferTsFileSealReq.toTPipeTransferReq(
-                        tsFile.getName(), tsFile.length())
-                    : PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
-                        tsFile.getName(),
-                        tsFile.length(),
-                        dataBaseName,
-                        sink.shouldWaitForSchemaBeforeLoad());
+        final TPipeTransferReq uncompressedReq;
+        if (transferMod) {
+          uncompressedReq =
+              PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
+                      modFile.getName(),
+                      modFile.length(),
+                      tsFile.getName(),
+                      tsFile.length(),
+                      dataBaseName,
+                      sink.shouldWaitForSchemaBeforeLoad())
+                  .setConversionTaskInfo(
+                      conversionTaskId, 
sink.shouldAsyncLoadTsFileOnTypeMismatch());
+        } else if (conversionTaskId != null
+            || dataBaseName != null
+            || sink.shouldWaitForSchemaBeforeLoad()) {
+          uncompressedReq =
+              PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
+                      tsFile.getName(),
+                      tsFile.length(),
+                      dataBaseName,
+                      sink.shouldWaitForSchemaBeforeLoad())
+                  .setConversionTaskInfo(
+                      conversionTaskId, 
sink.shouldAsyncLoadTsFileOnTypeMismatch());
+        } else {
+          uncompressedReq =
+              PipeTransferTsFileSealReq.toTPipeTransferReq(tsFile.getName(), 
tsFile.length());
+        }
         final TPipeTransferReq req = sink.compressIfNeeded(uncompressedReq);
 
         pipeName2WeightMap.forEach(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
index 6b388a6b2be..e5e1aca2245 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/thrift/sync/IoTDBDataRegionSyncSink.java
@@ -278,10 +278,12 @@ public class IoTDBDataRegionSyncSink extends 
IoTDBDataNodeSyncSink {
       throws IOException, WriteProcessException {
     final List<File> sealedFiles = batchToTransfer.sealTsFiles();
     final Map<Pair<String, Long>, Double> pipe2WeightMap = 
batchToTransfer.deepCopyPipe2WeightMap();
+    final List<EnrichedEvent> events = batchToTransfer.deepCopyEvents();
 
     try {
-      for (final File tsFile : sealedFiles) {
-        doTransfer(pipe2WeightMap, tsFile, null, null);
+      for (int outputIndex = 0; outputIndex < sealedFiles.size(); 
outputIndex++) {
+        final File tsFile = sealedFiles.get(outputIndex);
+        doTransfer(pipe2WeightMap, tsFile, null, null, events, outputIndex);
       }
     } finally {
       for (final File tsFile : sealedFiles) {
@@ -439,7 +441,9 @@ public class IoTDBDataRegionSyncSink extends 
IoTDBDataNodeSyncSink {
               1.0),
           pipeTsFileInsertionEvent.getTsFile(),
           pipeTsFileInsertionEvent.isWithMod() ? 
pipeTsFileInsertionEvent.getModFile() : null,
-          pipeTsFileInsertionEvent.getDatabaseName());
+          pipeTsFileInsertionEvent.getDatabaseName(),
+          Collections.singletonList(pipeTsFileInsertionEvent),
+          0);
     } finally {
       pipeTsFileInsertionEvent.decreaseReferenceCount(
           IoTDBDataRegionSyncSink.class.getName(), false);
@@ -450,11 +454,22 @@ public class IoTDBDataRegionSyncSink extends 
IoTDBDataNodeSyncSink {
       final Map<Pair<String, Long>, Double> pipeName2WeightMap,
       final File tsFile,
       final File modFile,
-      final String dataBaseName)
+      final String dataBaseName,
+      final Iterable<? extends EnrichedEvent> events,
+      final int outputIndex)
       throws PipeException, IOException {
 
     final Pair<IoTDBSyncClient, Boolean> clientAndStatus = 
clientManager.getClient();
     final TPipeTransferResp resp;
+    final String conversionTaskId =
+        shouldAsyncLoadTsFileOnTypeMismatch
+            ? PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+                sinkTaskId,
+                events,
+                dataBaseName,
+                outputIndex,
+                Objects.nonNull(modFile) && 
clientManager.supportModsIfIsDataNodeReceiver())
+            : null;
 
     // 1. Transfer tsFile, and mod file if exists and receiver's version >= 2
     if (Objects.nonNull(modFile) && 
clientManager.supportModsIfIsDataNodeReceiver()) {
@@ -466,12 +481,13 @@ public class IoTDBDataRegionSyncSink extends 
IoTDBDataNodeSyncSink {
         final TPipeTransferReq req =
             compressIfNeeded(
                 PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
-                    modFile.getName(),
-                    modFile.length(),
-                    tsFile.getName(),
-                    tsFile.length(),
-                    dataBaseName,
-                    shouldWaitForSchemaBeforeLoad));
+                        modFile.getName(),
+                        modFile.length(),
+                        tsFile.getName(),
+                        tsFile.length(),
+                        dataBaseName,
+                        shouldWaitForSchemaBeforeLoad)
+                    .setConversionTaskInfo(conversionTaskId, 
shouldAsyncLoadTsFileOnTypeMismatch));
 
         pipeName2WeightMap.forEach(
             (pipePair, weight) ->
@@ -496,14 +512,16 @@ public class IoTDBDataRegionSyncSink extends 
IoTDBDataNodeSyncSink {
       try {
         final TPipeTransferReq req =
             compressIfNeeded(
-                dataBaseName == null && !shouldWaitForSchemaBeforeLoad
+                conversionTaskId == null && dataBaseName == null && 
!shouldWaitForSchemaBeforeLoad
                     ? PipeTransferTsFileSealReq.toTPipeTransferReq(
                         tsFile.getName(), tsFile.length())
                     : PipeTransferTsFileSealWithModReq.toTPipeTransferReq(
-                        tsFile.getName(),
-                        tsFile.length(),
-                        dataBaseName,
-                        shouldWaitForSchemaBeforeLoad));
+                            tsFile.getName(),
+                            tsFile.length(),
+                            dataBaseName,
+                            shouldWaitForSchemaBeforeLoad)
+                        .setConversionTaskInfo(
+                            conversionTaskId, 
shouldAsyncLoadTsFileOnTypeMismatch));
 
         pipeName2WeightMap.forEach(
             (pipePair, weight) ->
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
index a4a23478612..1715fe4bb27 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
@@ -69,6 +69,7 @@ import 
org.apache.iotdb.db.storageengine.dataregion.utils.TsFileResourceUtils;
 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.LoadTsFileDataTypeConverter;
+import 
org.apache.iotdb.db.storageengine.load.converter.PipeTsFileConversionTaskManager;
 import org.apache.iotdb.db.storageengine.load.memory.LoadTsFileMemoryBlock;
 import org.apache.iotdb.db.storageengine.load.memory.LoadTsFileMemoryManager;
 import org.apache.iotdb.db.storageengine.load.metrics.LoadTsFileCostMetricsSet;
@@ -525,6 +526,9 @@ public class LoadTsFileAnalyzer implements AutoCloseable {
 
   private Analysis executeTabletConversionOnException(
       final Analysis analysis, final LoadAnalyzeException e) {
+    if (e instanceof LoadAnalyzeTypeMismatchException) {
+      PipeTsFileConversionTaskManager.markTypeMismatchDetected();
+    }
     if (setTemporaryUnavailableStatusIfNecessary(analysis, e)) {
       return analysis;
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadDirScanner.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadDirScanner.java
index a5c419c4d85..c07c0d8e61a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadDirScanner.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadDirScanner.java
@@ -108,6 +108,7 @@ public class ActiveLoadDirScanner extends 
ActiveLoadScheduledExecutorService {
           FileUtils.streamFiles(listeningDirFile, true, (String[]) null)) {
         try {
           fileStream
+              .filter(file -> 
!ActiveLoadPathHelper.isTransferStagingFile(file, listeningDirFile))
               .map(file -> new File(getTsFilePath(file.getAbsolutePath())))
               .distinct()
               .filter(file -> 
!activeLoadTsFileLoader.isFilePendingOrLoading(file))
@@ -116,10 +117,21 @@ public class ActiveLoadDirScanner extends 
ActiveLoadScheduledExecutorService {
               .limit(currentAllowedPendingSize)
               .forEach(
                   tsFile -> {
-                    activeLoadTsFileLoader.tryTriggerTsFileLoad(
-                        tsFile.getAbsolutePath(),
-                        listeningDirFile.getAbsolutePath(),
-                        isGeneratedByPipe);
+                    final String conversionTaskId =
+                        ActiveLoadPathHelper.parseAttributes(tsFile, 
listeningDirFile)
+                            
.get(ActiveLoadPathHelper.PIPE_CONVERSION_TASK_ID_KEY);
+                    if (conversionTaskId == null) {
+                      activeLoadTsFileLoader.tryTriggerTsFileLoad(
+                          tsFile.getAbsolutePath(),
+                          listeningDirFile.getAbsolutePath(),
+                          isGeneratedByPipe);
+                    } else {
+                      activeLoadTsFileLoader.tryTriggerTsFileLoad(
+                          tsFile.getAbsolutePath(),
+                          listeningDirFile.getAbsolutePath(),
+                          isGeneratedByPipe,
+                          conversionTaskId);
+                    }
                   });
         } catch (UncheckedIOException e) {
           LOGGER.debug("The file has been deleted. Ignore this exception.");
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPathHelper.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPathHelper.java
index 7c131a2e912..1704f58f1f2 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPathHelper.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPathHelper.java
@@ -45,6 +45,8 @@ public final class ActiveLoadPathHelper {
 
   private static final String SEGMENT_SEPARATOR = "-";
   public static final String USER_KEY = "user";
+  public static final String PIPE_CONVERSION_TASK_ID_KEY = 
"pipe-conversion-task-id";
+  private static final String TRANSFER_STAGING_DIRECTORY_PREFIX = 
".iotdb-load-staging-";
   // Keep a version in the user path segment so future encryption algorithms 
can be added safely.
   private static final String USER_VALUE_MASK_PREFIX = "v1-";
   private static final BaseEncoding USER_VALUE_ENCODING = 
BaseEncoding.base32().omitPadding();
@@ -53,6 +55,7 @@ public final class ActiveLoadPathHelper {
       Collections.unmodifiableList(
           Arrays.asList(
               USER_KEY,
+              PIPE_CONVERSION_TASK_ID_KEY,
               LoadTsFileConfigurator.DATABASE_NAME_KEY,
               LoadTsFileConfigurator.DATABASE_LEVEL_KEY,
               LoadTsFileConfigurator.CONVERT_ON_TYPE_MISMATCH_KEY,
@@ -133,6 +136,32 @@ public final class ActiveLoadPathHelper {
     return attributes;
   }
 
+  public static Map<String, String> buildAttributes(
+      final String databaseName,
+      final Integer databaseLevel,
+      final Boolean convertOnTypeMismatch,
+      final Boolean verify,
+      final Boolean autoCreateSchema,
+      final Long tabletConversionThresholdBytes,
+      final Boolean pipeGenerated,
+      final String userName,
+      final String conversionTaskId) {
+    final Map<String, String> attributes =
+        buildAttributes(
+            databaseName,
+            databaseLevel,
+            convertOnTypeMismatch,
+            verify,
+            autoCreateSchema,
+            tabletConversionThresholdBytes,
+            pipeGenerated,
+            userName);
+    if (conversionTaskId != null && !conversionTaskId.isEmpty()) {
+      attributes.put(PIPE_CONVERSION_TASK_ID_KEY, conversionTaskId);
+    }
+    return attributes;
+  }
+
   public static File resolveTargetDir(final File baseDir, final Map<String, 
String> attributes) {
     File current = baseDir;
     for (final String key : KEY_ORDER) {
@@ -145,6 +174,48 @@ public final class ActiveLoadPathHelper {
     return current;
   }
 
+  public static File resolvePipeTransferTargetDir(
+      final File baseDir, final Map<String, String> attributes) {
+    File current = baseDir;
+    for (final String key : KEY_ORDER) {
+      if (PIPE_CONVERSION_TASK_ID_KEY.equals(key)) {
+        continue;
+      }
+      final String value = attributes.get(key);
+      if (value == null) {
+        continue;
+      }
+      current = new File(current, formatSegment(key, value));
+    }
+    return current;
+  }
+
+  public static String formatPipeTaskTransferDirectoryName(final String 
conversionTaskId) {
+    return formatSegment(PIPE_CONVERSION_TASK_ID_KEY, conversionTaskId);
+  }
+
+  public static String formatTransferStagingDirectoryName(final String 
uniqueSuffix) {
+    return TRANSFER_STAGING_DIRECTORY_PREFIX + uniqueSuffix;
+  }
+
+  public static boolean isTransferStagingFile(final File file, final File 
pendingDir) {
+    if (file == null) {
+      return false;
+    }
+    final File normalizedPendingDir = pendingDir == null ? null : 
pendingDir.getAbsoluteFile();
+    File current = file.getAbsoluteFile();
+    while (current != null) {
+      if (normalizedPendingDir != null && 
current.equals(normalizedPendingDir)) {
+        return false;
+      }
+      if (current.getName().startsWith(TRANSFER_STAGING_DIRECTORY_PREFIX)) {
+        return true;
+      }
+      current = current.getParentFile();
+    }
+    return false;
+  }
+
   public static Map<String, String> parseAttributes(final File file, final 
File pendingDir) {
     if (file == null) {
       return Collections.emptyMap();
@@ -298,6 +369,11 @@ public final class ActiveLoadPathHelper {
       case LoadTsFileConfigurator.AUTO_CREATE_SCHEMA_KEY:
         LoadTsFileConfigurator.validateAutoCreateSchemaParam(value);
         break;
+      case PIPE_CONVERSION_TASK_ID_KEY:
+        if (value == null || value.isEmpty()) {
+          throw new SemanticException("Pipe conversion task id must not be 
empty.");
+        }
+        break;
       case USER_KEY:
         if (value == null || value.isEmpty()) {
           throw new SemanticException("User name must not be empty.");
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPendingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPendingQueue.java
index f50af8d8d60..3275ba8368a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPendingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadPendingQueue.java
@@ -35,8 +35,17 @@ public class ActiveLoadPendingQueue {
 
   public synchronized boolean enqueue(
       final String file, final String pendingDir, final boolean 
isGeneratedByPipe) {
+    return enqueue(file, pendingDir, isGeneratedByPipe, null);
+  }
+
+  public synchronized boolean enqueue(
+      final String file,
+      final String pendingDir,
+      final boolean isGeneratedByPipe,
+      final String conversionTaskId) {
     if (!loadingFileSet.contains(file) && pendingFileSet.add(file)) {
-      pendingFileQueue.offer(new ActiveLoadEntry(file, pendingDir, 
isGeneratedByPipe));
+      pendingFileQueue.offer(
+          new ActiveLoadEntry(file, pendingDir, isGeneratedByPipe, 
conversionTaskId));
 
       
ActiveLoadingFilesNumberMetricsSet.getInstance().increaseQueuingFileCounter(1);
       return true;
@@ -92,11 +101,18 @@ public class ActiveLoadPendingQueue {
     private final String file;
     private final String pendingDir;
     private final boolean isGeneratedByPipe;
+    private final String conversionTaskId;
 
     public ActiveLoadEntry(String file, String pendingDir, boolean 
isGeneratedByPipe) {
+      this(file, pendingDir, isGeneratedByPipe, null);
+    }
+
+    public ActiveLoadEntry(
+        String file, String pendingDir, boolean isGeneratedByPipe, String 
conversionTaskId) {
       this.file = file;
       this.pendingDir = pendingDir;
       this.isGeneratedByPipe = isGeneratedByPipe;
+      this.conversionTaskId = conversionTaskId;
     }
 
     public String getFile() {
@@ -110,5 +126,9 @@ public class ActiveLoadPendingQueue {
     public boolean isGeneratedByPipe() {
       return isGeneratedByPipe;
     }
+
+    public String getConversionTaskId() {
+      return conversionTaskId;
+    }
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java
index 7d661a18109..fa26572e811 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java
@@ -38,6 +38,7 @@ import 
org.apache.iotdb.db.queryengine.plan.analyze.schema.ClusterSchemaFetcher;
 import org.apache.iotdb.db.queryengine.plan.statement.Statement;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
 import 
org.apache.iotdb.db.queryengine.plan.statement.pipe.PipeEnrichedStatement;
+import 
org.apache.iotdb.db.storageengine.load.converter.PipeTsFileConversionTaskManager;
 import 
org.apache.iotdb.db.storageengine.load.metrics.ActiveLoadingFilesNumberMetricsSet;
 import 
org.apache.iotdb.db.storageengine.load.metrics.ActiveLoadingFilesSizeMetricsSet;
 import org.apache.iotdb.rpc.TSStatusCode;
@@ -85,11 +86,16 @@ public class ActiveLoadTsFileLoader {
 
   public void tryTriggerTsFileLoad(
       String absolutePath, String pendingDir, boolean isGeneratedByPipe) {
+    tryTriggerTsFileLoad(absolutePath, pendingDir, isGeneratedByPipe, null);
+  }
+
+  public void tryTriggerTsFileLoad(
+      String absolutePath, String pendingDir, boolean isGeneratedByPipe, 
String conversionTaskId) {
     if (CommonDescriptor.getInstance().getConfig().isReadOnly()) {
       return;
     }
 
-    if (pendingQueue.enqueue(absolutePath, pendingDir, isGeneratedByPipe)) {
+    if (pendingQueue.enqueue(absolutePath, pendingDir, isGeneratedByPipe, 
conversionTaskId)) {
       initFailDirIfNecessary();
       adjustExecutorIfNecessary();
     }
@@ -262,13 +268,31 @@ public class ActiveLoadTsFileLoader {
             ? ActiveLoadPathHelper.findPendingDirectory(tsFile)
             : new File(entry.getPendingDir());
     final Map<String, String> attributes = 
ActiveLoadPathHelper.parseAttributes(tsFile, pendingDir);
-    ActiveLoadPathHelper.applyAttributesToStatement(attributes, statement, 
isVerify);
-    final String userName =
-        attributes.getOrDefault(ActiveLoadPathHelper.USER_KEY, 
AuthorityChecker.SUPER_USER);
-    session.setUsername(userName);
-
-    return executeStatement(
-        entry.isGeneratedByPipe() ? new PipeEnrichedStatement(statement) : 
statement, session);
+    final String conversionTaskId = entry.getConversionTaskId();
+    PipeTsFileConversionTaskManager.registerIfAbsent(conversionTaskId);
+    PipeTsFileConversionTaskManager.markReceiverOwned(conversionTaskId);
+    PipeTsFileConversionTaskManager.markRunning(conversionTaskId);
+    PipeTsFileConversionTaskManager.enter(conversionTaskId);
+    try {
+      ActiveLoadPathHelper.applyAttributesToStatement(attributes, statement, 
isVerify);
+      final String userName =
+          attributes.getOrDefault(ActiveLoadPathHelper.USER_KEY, 
AuthorityChecker.SUPER_USER);
+      session.setUsername(userName);
+
+      final TSStatus result =
+          executeStatement(
+              entry.isGeneratedByPipe() ? new PipeEnrichedStatement(statement) 
: statement,
+              session);
+      if (result.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+          || result.getCode() == 
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()) {
+        PipeTsFileConversionTaskManager.markSuccess(conversionTaskId);
+      } else {
+        PipeTsFileConversionTaskManager.markPaused(conversionTaskId, result);
+      }
+      return result;
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
   }
 
   private TSStatus executeStatement(final Statement statement, final 
IClientSession session) {
@@ -293,7 +317,10 @@ public class ActiveLoadTsFileLoader {
 
   private void handleLoadFailure(
       final ActiveLoadPendingQueue.ActiveLoadEntry entry, final TSStatus 
status) {
-    if (!ActiveLoadFailedMessageHandler.isStatusShouldRetry(entry, status)) {
+    if (ActiveLoadFailedMessageHandler.isStatusShouldRetry(entry, status)) {
+      PipeTsFileConversionTaskManager.markPaused(entry.getConversionTaskId(), 
status);
+    } else {
+      PipeTsFileConversionTaskManager.markFailed(entry.getConversionTaskId(), 
status);
       LOGGER.warn(
           "Failed to auto load tsfile {} (isGeneratedByPipe = {}), status: {}. 
File will be moved to fail directory.",
           entry.getFile(),
@@ -304,6 +331,8 @@ public class ActiveLoadTsFileLoader {
   }
 
   private void handleFileNotFoundException(final 
ActiveLoadPendingQueue.ActiveLoadEntry entry) {
+    PipeTsFileConversionTaskManager.markFailed(
+        entry.getConversionTaskId(), new 
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode()));
     LOGGER.warn(
         "Failed to auto load tsfile {} (isGeneratedByPipe = {}) due to file 
not found, will skip this file.",
         entry.getFile(),
@@ -313,7 +342,15 @@ public class ActiveLoadTsFileLoader {
 
   private void handleOtherException(
       final ActiveLoadPendingQueue.ActiveLoadEntry entry, final Exception e) {
-    if (!ActiveLoadFailedMessageHandler.isExceptionMessageShouldRetry(entry, 
e.getMessage())) {
+    if (ActiveLoadFailedMessageHandler.isExceptionMessageShouldRetry(entry, 
e.getMessage())) {
+      PipeTsFileConversionTaskManager.markPaused(
+          entry.getConversionTaskId(),
+          new 
TSStatus(TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+              .setMessage(e.getMessage()));
+    } else {
+      PipeTsFileConversionTaskManager.markFailed(
+          entry.getConversionTaskId(),
+          new 
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode()).setMessage(e.getMessage()));
       LOGGER.warn(
           "Failed to auto load tsfile {} (isGeneratedByPipe = {}) because of 
an unexpected exception. File will be moved to fail directory.",
           entry.getFile(),
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtil.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtil.java
index 93f8f04b481..e8807711a07 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtil.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtil.java
@@ -36,6 +36,7 @@ import org.slf4j.LoggerFactory;
 
 import java.io.File;
 import java.io.IOException;
+import java.nio.file.AtomicMoveNotSupportedException;
 import java.nio.file.Files;
 import java.nio.file.StandardCopyOption;
 import java.util.ArrayList;
@@ -96,7 +97,8 @@ public class ActiveLoadUtil {
       return false;
     }
     final Map<String, String> attributes = 
appendCurrentUserIfAbsent(loadAttributes);
-    final File targetDir = 
ActiveLoadPathHelper.resolveTargetDir(targetFilePath, attributes);
+    final File targetDir =
+        ActiveLoadPathHelper.resolvePipeTransferTargetDir(targetFilePath, 
attributes);
 
     transferFilesToActiveDir(
         targetDir,
@@ -149,53 +151,147 @@ public class ActiveLoadUtil {
       return false;
     }
     final Map<String, String> attributes = 
appendCurrentUserIfAbsent(loadAttributes);
-    final File targetDir = 
ActiveLoadPathHelper.resolveTargetDir(targetFilePath, attributes);
+    final File targetDir =
+        ActiveLoadPathHelper.resolvePipeTransferTargetDir(targetFilePath, 
attributes);
 
     final List<File> sourceFiles = new ArrayList<>(files.size());
     for (final String file : files) {
       sourceFiles.add(new File(file));
     }
     sourceFiles.sort(Comparator.comparing(ActiveLoadUtil::isTsFile));
-    transferFilesToActiveDir(targetDir, sourceFiles, isDeleteAfterLoad);
+    transferFilesToActiveDir(
+        targetDir,
+        sourceFiles,
+        isDeleteAfterLoad,
+        attributes.get(ActiveLoadPathHelper.PIPE_CONVERSION_TASK_ID_KEY));
     return true;
   }
 
   static void transferFilesToActiveDir(
       final File targetDir, final List<File> sourceFiles, final boolean 
isDeleteAfterLoad)
       throws IOException {
+    transferFilesToActiveDir(targetDir, sourceFiles, isDeleteAfterLoad, null);
+  }
+
+  static void transferFilesToActiveDir(
+      final File targetDir,
+      final List<File> sourceFiles,
+      final boolean isDeleteAfterLoad,
+      final String conversionTaskId)
+      throws IOException {
     final List<File> existingSourceFiles = new ArrayList<>(sourceFiles.size());
     for (final File sourceFile : sourceFiles) {
       if (sourceFile.exists()) {
         existingSourceFiles.add(sourceFile);
       }
     }
+    final File transferDir =
+        new File(
+            targetDir,
+            conversionTaskId == null
+                ? UUID.randomUUID().toString()
+                : 
ActiveLoadPathHelper.formatPipeTaskTransferDirectoryName(conversionTaskId));
+
+    if (conversionTaskId != null && transferDir.exists()) {
+      if (!isExistingTaskComplete(transferDir, sourceFiles)) {
+        throw new IOException("Failed to load TsFile to active directory.");
+      }
+      if (isDeleteAfterLoad) {
+        deleteSourceFiles(existingSourceFiles);
+      }
+      return;
+    }
+
     if (existingSourceFiles.isEmpty()) {
+      if (conversionTaskId != null) {
+        throw new IOException("Failed to load TsFile to active directory.");
+      }
       return;
     }
 
-    final File transferDir = new File(targetDir, UUID.randomUUID().toString());
+    final File stagingDir =
+        new File(
+            targetDir,
+            
ActiveLoadPathHelper.formatTransferStagingDirectoryName(UUID.randomUUID().toString()));
     try {
-      Files.createDirectories(transferDir.toPath());
+      Files.createDirectories(stagingDir.toPath());
       for (final File sourceFile : existingSourceFiles) {
-        final File targetFile = new File(transferDir, sourceFile.getName());
+        final File targetFile = new File(stagingDir, sourceFile.getName());
         RetryUtils.retryOnException(
             () -> {
               transferFile(sourceFile, targetFile, isDeleteAfterLoad);
               return null;
             });
       }
+      try {
+        publishTransferDirectory(stagingDir, transferDir);
+      } catch (final IOException e) {
+        if (conversionTaskId == null || !isExistingTaskComplete(transferDir, 
sourceFiles)) {
+          throw e;
+        }
+      }
     } catch (final IOException | RuntimeException e) {
-      if (transferDir.exists()) {
-        FileUtils.deleteFileOrDirectoryWithRetry(transferDir);
+      if (stagingDir.exists()) {
+        FileUtils.deleteFileOrDirectoryWithRetry(stagingDir);
       }
       throw e;
     }
 
+    if (stagingDir.exists()) {
+      FileUtils.deleteFileOrDirectoryWithRetry(stagingDir);
+    }
     if (isDeleteAfterLoad) {
       deleteSourceFiles(existingSourceFiles);
     }
   }
 
+  private static boolean isExistingTaskComplete(
+      final File transferDir, final List<File> sourceFiles) {
+    if (!transferDir.isDirectory()) {
+      return false;
+    }
+    final File[] targetFiles = transferDir.listFiles(File::isFile);
+    if (targetFiles == null || targetFiles.length == 0) {
+      return false;
+    }
+
+    final List<File> existingSourceFiles =
+        
sourceFiles.stream().filter(File::isFile).collect(java.util.stream.Collectors.toList());
+    if (existingSourceFiles.isEmpty()) {
+      return Arrays.stream(targetFiles).anyMatch(ActiveLoadUtil::isTsFile);
+    }
+    if (targetFiles.length != existingSourceFiles.size()) {
+      return false;
+    }
+
+    final boolean[] matched = new boolean[targetFiles.length];
+    for (final File sourceFile : existingSourceFiles) {
+      boolean found = false;
+      for (int i = 0; i < targetFiles.length; i++) {
+        if (!matched[i]
+            && isTsFile(sourceFile) == isTsFile(targetFiles[i])
+            && sourceFile.length() == targetFiles[i].length()) {
+          matched[i] = true;
+          found = true;
+          break;
+        }
+      }
+      if (!found) {
+        return false;
+      }
+    }
+    return true;
+  }
+
+  private static void publishTransferDirectory(final File stagingDir, final 
File transferDir)
+      throws IOException {
+    try {
+      Files.move(stagingDir.toPath(), transferDir.toPath(), 
StandardCopyOption.ATOMIC_MOVE);
+    } catch (final AtomicMoveNotSupportedException e) {
+      Files.move(stagingDir.toPath(), transferDir.toPath());
+    }
+  }
+
   private static void transferFile(
       final File sourceFile, final File targetFile, final boolean useHardLink) 
throws IOException {
     Exception linkException = null;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
index a1da7095246..03e13f01dca 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
@@ -80,87 +80,93 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
 
     LOGGER.info("Start data type conversion for LoadTsFileStatement: {}", 
loadTsFileStatement);
 
-    final LoadTsFileMemoryBlock block =
-        LoadTsFileMemoryManager.getInstance()
-            .allocateMemoryBlock(TABLET_BATCH_MEMORY_SIZE_IN_BYTES);
-    final List<PipeTransferTabletRawReq> tabletRawReqs = new ArrayList<>();
-    final List<Long> tabletRawReqSizes = new ArrayList<>();
+    final boolean isManagedTask = 
PipeTsFileConversionTaskManager.getCurrentTaskId() != null;
+    final TreeConversionContext conversionContext =
+        isManagedTask
+            ? 
PipeTsFileConversionTaskManager.getOrCreateCurrentContext(TreeConversionContext::new)
+            : new TreeConversionContext();
+    boolean shouldReleaseContext = !isManagedTask;
 
     try {
-      for (final File file : loadTsFileStatement.getTsFiles()) {
-        try (final LoadTreeTsFileTabletIterator tabletIterator =
-            new LoadTreeTsFileTabletIterator(file, true)) {
-          for (final Pair<Tablet, Boolean> tabletWithIsAligned : 
tabletIterator) {
-            final PipeTransferTabletRawReq tabletRawReq =
-                PipeTransferTabletRawReq.toTPipeTransferRawReq(
-                    tabletWithIsAligned.getLeft(), 
tabletWithIsAligned.getRight());
-            final long curMemory = 
calculateTabletSizeInBytes(tabletWithIsAligned.getLeft()) + 1;
-            if (block.hasEnoughMemory(curMemory)) {
-              tabletRawReqs.add(tabletRawReq);
-              tabletRawReqSizes.add(curMemory);
-              block.addMemoryUsage(curMemory);
-              continue;
-            }
-
-            final TSStatus result =
-                executeInsertMultiTabletsWithRetry(
-                    tabletRawReqs, 
loadTsFileStatement.isConvertOnTypeMismatch());
-
-            for (final long memoryCost : tabletRawReqSizes) {
-              block.reduceMemoryUsage(memoryCost);
-            }
-            tabletRawReqs.clear();
-            tabletRawReqSizes.clear();
-
-            if (!handleTSStatus(result, loadTsFileStatement)) {
-              return Optional.of(result);
-            }
-
-            tabletRawReqs.add(tabletRawReq);
-            tabletRawReqSizes.add(curMemory);
-            block.addMemoryUsage(curMemory);
-          }
-        } catch (final Exception e) {
-          LOGGER.warn(
-              "Failed to convert data type for LoadTsFileStatement: {}.", 
loadTsFileStatement, e);
-          return Optional.of(
-              loadTsFileStatement.accept(
-                  LoadTsFileDataTypeConverter.STATEMENT_EXCEPTION_VISITOR, e));
+      final List<File> files = loadTsFileStatement.getTsFiles();
+      while (conversionContext.fileIndex < files.size()) {
+        if (conversionContext.tabletIterator == null) {
+          conversionContext.tabletIterator =
+              new 
LoadTreeTsFileTabletIterator(files.get(conversionContext.fileIndex), true);
         }
-      }
 
-      if (!tabletRawReqs.isEmpty()) {
-        try {
+        if (conversionContext.deferredTabletRawReq != null) {
           final TSStatus result =
-              executeInsertMultiTabletsWithRetry(
-                  tabletRawReqs, 
loadTsFileStatement.isConvertOnTypeMismatch());
+              flushPendingTablets(conversionContext, 
loadTsFileStatement.isConvertOnTypeMismatch());
+          if (!handleTSStatus(result, loadTsFileStatement)) {
+            shouldReleaseContext = !isManagedTask || 
!isTemporaryUnavailable(result);
+            return Optional.of(result);
+          }
+          conversionContext.addDeferredTablet();
+        }
 
-          for (final long memoryCost : tabletRawReqSizes) {
-            block.reduceMemoryUsage(memoryCost);
+        while (conversionContext.deferredTabletWithIsAligned != null
+            || conversionContext.tabletIterator.hasNext()) {
+          if (conversionContext.deferredTabletWithIsAligned == null) {
+            conversionContext.deferredTabletWithIsAligned = 
conversionContext.tabletIterator.next();
+          }
+          final Pair<Tablet, Boolean> tabletWithIsAligned =
+              conversionContext.deferredTabletWithIsAligned;
+          final PipeTransferTabletRawReq tabletRawReq =
+              PipeTransferTabletRawReq.toTPipeTransferRawReq(
+                  tabletWithIsAligned.getLeft(), 
tabletWithIsAligned.getRight());
+          final long currentMemory = 
calculateTabletSizeInBytes(tabletWithIsAligned.getLeft()) + 1;
+          if (conversionContext.block.hasEnoughMemory(currentMemory)) {
+            conversionContext.addTablet(tabletRawReq, currentMemory);
+            conversionContext.deferredTabletWithIsAligned = null;
+            continue;
           }
-          tabletRawReqs.clear();
-          tabletRawReqSizes.clear();
 
+          final TSStatus result =
+              flushPendingTablets(conversionContext, 
loadTsFileStatement.isConvertOnTypeMismatch());
           if (!handleTSStatus(result, loadTsFileStatement)) {
+            conversionContext.deferredTabletRawReq = tabletRawReq;
+            conversionContext.deferredTabletRawReqSize = currentMemory;
+            conversionContext.deferredTabletWithIsAligned = null;
+            shouldReleaseContext = !isManagedTask || 
!isTemporaryUnavailable(result);
             return Optional.of(result);
           }
-        } catch (final Exception e) {
-          LOGGER.warn(
-              "Failed to convert data type for LoadTsFileStatement: {}.", 
loadTsFileStatement, e);
-          return Optional.of(
-              loadTsFileStatement.accept(
-                  LoadTsFileDataTypeConverter.STATEMENT_EXCEPTION_VISITOR, e));
+          conversionContext.addTablet(tabletRawReq, currentMemory);
+          conversionContext.deferredTabletWithIsAligned = null;
         }
+
+        conversionContext.tabletIterator.close();
+        conversionContext.tabletIterator = null;
+        conversionContext.fileIndex++;
       }
+
+      if (!conversionContext.tabletRawReqs.isEmpty()) {
+        final TSStatus result =
+            flushPendingTablets(conversionContext, 
loadTsFileStatement.isConvertOnTypeMismatch());
+        if (!handleTSStatus(result, loadTsFileStatement)) {
+          shouldReleaseContext = !isManagedTask || 
!isTemporaryUnavailable(result);
+          return Optional.of(result);
+        }
+      }
+    } catch (final Exception e) {
+      LOGGER.warn(
+          "Failed to convert data type for LoadTsFileStatement: {}.", 
loadTsFileStatement, e);
+      final TSStatus status =
+          
loadTsFileStatement.accept(LoadTsFileDataTypeConverter.STATEMENT_EXCEPTION_VISITOR,
 e);
+      shouldReleaseContext =
+          !isManagedTask || 
!LoadTsFileDataTypeConverter.isMemoryPressureException(e);
+      return Optional.of(status);
     } finally {
-      for (final long memoryCost : tabletRawReqSizes) {
-        block.reduceMemoryUsage(memoryCost);
+      if (shouldReleaseContext) {
+        if (isManagedTask) {
+          PipeTsFileConversionTaskManager.clearCurrentContext();
+        } else {
+          conversionContext.close();
+        }
       }
-      tabletRawReqs.clear();
-      tabletRawReqSizes.clear();
-      block.close();
     }
 
+    shouldReleaseContext = true;
     if (loadTsFileStatement.isDeleteAfterLoad()) {
       loadTsFileStatement
           .getTsFiles()
@@ -179,6 +185,72 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
     return Optional.of(new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
   }
 
+  private TSStatus flushPendingTablets(
+      final TreeConversionContext context, final boolean 
isConvertedOnTypeMismatch) {
+    if (context.tabletRawReqs.isEmpty()) {
+      return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+    }
+    final TSStatus result =
+        executeInsertMultiTabletsWithRetry(context.tabletRawReqs, 
isConvertedOnTypeMismatch);
+    if (handleTSStatus(result, context)) {
+      context.clearPendingTablets();
+    }
+    return result;
+  }
+
+  private static boolean isTemporaryUnavailable(final TSStatus status) {
+    return status != null
+        && (status.getCode() == 
TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()
+            || status.getCode()
+                == 
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode());
+  }
+
+  private static final class TreeConversionContext implements AutoCloseable {
+    private final LoadTsFileMemoryBlock block =
+        LoadTsFileMemoryManager.getInstance()
+            .allocateMemoryBlock(TABLET_BATCH_MEMORY_SIZE_IN_BYTES);
+    private final List<PipeTransferTabletRawReq> tabletRawReqs = new 
ArrayList<>();
+    private final List<Long> tabletRawReqSizes = new ArrayList<>();
+    private int fileIndex;
+    private LoadTreeTsFileTabletIterator tabletIterator;
+    private Pair<Tablet, Boolean> deferredTabletWithIsAligned;
+    private PipeTransferTabletRawReq deferredTabletRawReq;
+    private long deferredTabletRawReqSize;
+
+    private void addTablet(final PipeTransferTabletRawReq request, final long 
size) {
+      tabletRawReqs.add(request);
+      tabletRawReqSizes.add(size);
+      block.addMemoryUsage(size);
+    }
+
+    private void addDeferredTablet() {
+      addTablet(deferredTabletRawReq, deferredTabletRawReqSize);
+      deferredTabletRawReq = null;
+      deferredTabletRawReqSize = 0;
+    }
+
+    private void clearPendingTablets() {
+      for (final long memoryCost : tabletRawReqSizes) {
+        block.reduceMemoryUsage(memoryCost);
+      }
+      tabletRawReqs.clear();
+      tabletRawReqSizes.clear();
+    }
+
+    @Override
+    public void close() {
+      if (tabletIterator != null) {
+        tabletIterator.close();
+        tabletIterator = null;
+      }
+      clearPendingTablets();
+      deferredTabletWithIsAligned = null;
+      deferredTabletRawReq = null;
+      deferredTabletRawReqSize = 0;
+      block.close();
+    }
+  }
+
   private TSStatus executeInsertMultiTabletsWithRetry(
       final List<PipeTransferTabletRawReq> tabletRawReqs, boolean 
isConvertOnTypeMismatch) {
     final InsertMultiTabletsStatement batchStatement = new 
InsertMultiTabletsStatement();
@@ -220,6 +292,10 @@ public class 
LoadTreeStatementDataTypeConvertExecutionVisitor
 
   private static boolean handleTSStatus(
       final TSStatus result, final LoadTsFileStatement loadTsFileStatement) {
+    return handleTSStatus(result, (Object) loadTsFileStatement);
+  }
+
+  private static boolean handleTSStatus(final TSStatus result, final Object 
loadTsFileStatement) {
     if (!(result.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
         || result.getCode() == 
TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()
         || result.getCode() == 
TSStatusCode.LOAD_IDEMPOTENT_CONFLICT_EXCEPTION.getStatusCode())) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeTsFileTabletIterator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeTsFileTabletIterator.java
index e65ba74244b..41c88ce4752 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeTsFileTabletIterator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeTsFileTabletIterator.java
@@ -117,7 +117,9 @@ class LoadTreeTsFileTabletIterator
         if (recoverFromIteratorFailure(e)) {
           continue;
         }
-        close();
+        if (!shouldRethrow(e)) {
+          close();
+        }
         throw toRuntimeException(e);
       }
     }
@@ -139,7 +141,9 @@ class LoadTreeTsFileTabletIterator
         if (recoverFromIteratorFailure(e)) {
           continue;
         }
-        close();
+        if (!shouldRethrow(e)) {
+          close();
+        }
         throw toRuntimeException(e);
       }
     }
@@ -167,6 +171,7 @@ class LoadTreeTsFileTabletIterator
         return;
       } catch (final Exception e) {
         if (shouldRethrow(e)) {
+          scanInitialized = false;
           throw toRuntimeException(e);
         }
         if (!switchFromScanToQuery(e)) {
@@ -363,6 +368,8 @@ class LoadTreeTsFileTabletIterator
         return true;
       } catch (final Exception e) {
         if (shouldRethrow(e)) {
+          pendingQueryTasks.addFirst(activeQueryTask);
+          activeQueryTask = null;
           throw toRuntimeException(e);
         }
         LOGGER.warn(
@@ -399,6 +406,9 @@ class LoadTreeTsFileTabletIterator
   }
 
   private boolean shouldRethrow(final Exception e) {
+    if (LoadTsFileDataTypeConverter.isMemoryPressureException(e)) {
+      return true;
+    }
     Throwable current = e;
     while (Objects.nonNull(current)) {
       if (current instanceof InterruptedException
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java
index 6fc1f1bf6d9..d9723a22343 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTsFileDataTypeConverter.java
@@ -58,7 +58,7 @@ public class LoadTsFileDataTypeConverter {
     return TabletConversionSemaphoreHolder.INSTANCE;
   }
 
-  private static int getTabletConversionPermitCount() {
+  static int getTabletConversionPermitCount() {
     final int configuredThreadCount =
         Math.max(
             1,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManager.java
new file mode 100644
index 00000000000..08d7c2223c4
--- /dev/null
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManager.java
@@ -0,0 +1,541 @@
+/*
+ * 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.storageengine.load.converter;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import com.google.common.annotations.VisibleForTesting;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.Iterator;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Supplier;
+
+/**
+ * Deduplicates pipe TsFile conversion tasks. The status table is bounded, and 
only a bounded number
+ * of tasks retain an in-memory parser checkpoint. The active-load directory 
remains the durable
+ * source of each receiver-owned file.
+ */
+public final class PipeTsFileConversionTaskManager {
+
+  private static final Logger LOGGER =
+      LoggerFactory.getLogger(PipeTsFileConversionTaskManager.class);
+
+  public enum State {
+    PENDING,
+    RUNNING,
+    PAUSED,
+    SUCCESS,
+    FAILED
+  }
+
+  public static final class Task {
+    private final String taskId;
+    private final boolean asyncLoadOnTypeMismatch;
+    private volatile State state = State.PENDING;
+    private volatile TSStatus status;
+    private volatile boolean typeMismatchDetected;
+    private volatile boolean receiverOwned;
+    private boolean retrySealAllowed;
+    private Object conversionContext;
+
+    private Task(final String taskId, final boolean asyncLoadOnTypeMismatch) {
+      this.taskId = taskId;
+      this.asyncLoadOnTypeMismatch = asyncLoadOnTypeMismatch;
+    }
+
+    public String getTaskId() {
+      return taskId;
+    }
+
+    public boolean isAsyncLoadOnTypeMismatch() {
+      return asyncLoadOnTypeMismatch;
+    }
+
+    public State getState() {
+      return state;
+    }
+
+    public TSStatus getStatus() {
+      return status;
+    }
+
+    public boolean isTypeMismatchDetected() {
+      return typeMismatchDetected;
+    }
+
+    public boolean isReceiverOwned() {
+      return receiverOwned;
+    }
+  }
+
+  private static final class UnretainedContext {
+    private final String taskId;
+    private final Object context;
+
+    private UnretainedContext(final String taskId, final Object context) {
+      this.taskId = taskId;
+      this.context = context;
+    }
+  }
+
+  private static final int MAX_TASKS = 4096;
+  private static final int MAX_CONTEXTS =
+      LoadTsFileDataTypeConverter.getTabletConversionPermitCount();
+  private static final Map<String, Task> TASKS = new LinkedHashMap<>(128, 
0.75F, true);
+  private static final AtomicInteger RETAINED_CONTEXT_COUNT = new 
AtomicInteger();
+  private static final ThreadLocal<String> CURRENT_TASK_ID = new 
ThreadLocal<>();
+  // Keeps legacy seal requests (which predate conversion task ids) eligible 
for receiver takeover.
+  private static final ThreadLocal<Boolean> CURRENT_TYPE_MISMATCH = new 
ThreadLocal<>();
+  private static final ThreadLocal<UnretainedContext> 
CURRENT_UNRETAINED_CONTEXT =
+      new ThreadLocal<>();
+
+  private PipeTsFileConversionTaskManager() {
+    // utility class
+  }
+
+  public static Task registerIfAbsent(final String taskId) {
+    return registerIfAbsent(taskId, true);
+  }
+
+  public static Task registerIfAbsent(final String taskId, final boolean 
asyncLoadOnTypeMismatch) {
+    if (taskId == null || taskId.isEmpty()) {
+      return null;
+    }
+    synchronized (TASKS) {
+      Task task = TASKS.get(taskId);
+      if (task == null) {
+        if (!hasTaskCapacity()) {
+          return null;
+        }
+        task = new Task(taskId, asyncLoadOnTypeMismatch);
+        TASKS.put(taskId, task);
+      }
+      return task;
+    }
+  }
+
+  public static Task get(final String taskId) {
+    if (taskId == null || taskId.isEmpty()) {
+      return null;
+    }
+    synchronized (TASKS) {
+      return TASKS.get(taskId);
+    }
+  }
+
+  /** Returns a response for a duplicate seal, or {@code null} when no task is 
known. */
+  public static TSStatus getDuplicateStatus(
+      final String taskId, final boolean asyncLoadOnTypeMismatch) {
+    if (taskId == null || taskId.isEmpty()) {
+      return null;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      return task == null ? null : getDuplicateStatus(task, 
asyncLoadOnTypeMismatch);
+    }
+  }
+
+  /** Atomically claims a new/retryable seal or returns the status of its 
existing task. */
+  public static TSStatus registerAndGetDuplicateStatus(
+      final String taskId, final boolean asyncLoadOnTypeMismatch) {
+    if (taskId == null || taskId.isEmpty()) {
+      return null;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task == null) {
+        if (!hasTaskCapacity()) {
+          return createReceiverTemporaryUnavailableStatus(null);
+        }
+        TASKS.put(taskId, new Task(taskId, asyncLoadOnTypeMismatch));
+        return null;
+      }
+      if (task.retrySealAllowed && !task.receiverOwned) {
+        task.retrySealAllowed = false;
+        task.status = null;
+        task.state = State.PENDING;
+        return null;
+      }
+      return getDuplicateStatus(task, asyncLoadOnTypeMismatch);
+    }
+  }
+
+  private static TSStatus getDuplicateStatus(
+      final Task task, final boolean asyncLoadOnTypeMismatch) {
+    if (asyncLoadOnTypeMismatch || task.isAsyncLoadOnTypeMismatch()) {
+      if (task.receiverOwned || task.state == State.SUCCESS) {
+        // Once the receiver owns the file, the sender must not create a 
second conversion task.
+        return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+      }
+      if (task.state == State.FAILED && task.status != null) {
+        return toReceiverStatus(task.status);
+      }
+      if (task.state == State.PAUSED && task.status != null) {
+        return toReceiverStatus(task.status);
+      }
+      return createReceiverTemporaryUnavailableStatus(null);
+    }
+    if (task.state == State.SUCCESS) {
+      return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+    }
+    if ((task.state == State.PAUSED || task.state == State.FAILED) && 
task.status != null) {
+      return toReceiverStatus(task.status);
+    }
+    return createReceiverTemporaryUnavailableStatus(null);
+  }
+
+  private static TSStatus toReceiverStatus(final TSStatus status) {
+    if (status == null
+        || status.getCode() != 
TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()) {
+      return status;
+    }
+    return createReceiverTemporaryUnavailableStatus(status.getMessage());
+  }
+
+  private static TSStatus createReceiverTemporaryUnavailableStatus(final 
String message) {
+    return new 
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+        .setMessage(message);
+  }
+
+  public static void enter(final String taskId) {
+    CURRENT_TYPE_MISMATCH.set(false);
+    if (taskId != null && !taskId.isEmpty()) {
+      final String previousTaskId = CURRENT_TASK_ID.get();
+      if (previousTaskId != null && !previousTaskId.equals(taskId)) {
+        clearCurrentUnretainedContext(previousTaskId);
+      }
+      CURRENT_TASK_ID.set(taskId);
+    } else {
+      clearCurrentUnretainedContext(CURRENT_TASK_ID.get());
+      CURRENT_TASK_ID.remove();
+    }
+  }
+
+  public static String getCurrentTaskId() {
+    return CURRENT_TASK_ID.get();
+  }
+
+  @SuppressWarnings("unchecked")
+  public static <T> T getOrCreateCurrentContext(final Supplier<T> supplier) {
+    final String currentTaskId = CURRENT_TASK_ID.get();
+    if (currentTaskId == null) {
+      return supplier.get();
+    }
+
+    final UnretainedContext currentUnretainedContext = 
CURRENT_UNRETAINED_CONTEXT.get();
+    if (currentUnretainedContext != null) {
+      if (currentTaskId.equals(currentUnretainedContext.taskId)) {
+        return (T) currentUnretainedContext.context;
+      }
+      clearCurrentUnretainedContext(currentUnretainedContext.taskId);
+    }
+
+    Object evictedContext = null;
+    final T context;
+    boolean retained = false;
+    synchronized (TASKS) {
+      final Task task = TASKS.get(currentTaskId);
+      if (task == null) {
+        context = supplier.get();
+      } else if (task.state == State.SUCCESS || task.state == State.FAILED) {
+        // A parser callback that races with terminal completion may finish 
its current call, but
+        // it must not recreate a checkpoint for a task that is already 
terminal.
+        context = supplier.get();
+      } else if (task.conversionContext != null) {
+        context = (T) task.conversionContext;
+        retained = true;
+      } else {
+        context = supplier.get();
+        final ContextReservation reservation = reserveContextSlot(task);
+        evictedContext = reservation.evictedContext;
+        if (reservation.slotAvailable) {
+          retainContext(task, context);
+          retained = true;
+        }
+      }
+    }
+    closeContext(evictedContext);
+    if (!retained) {
+      CURRENT_UNRETAINED_CONTEXT.set(new UnretainedContext(currentTaskId, 
context));
+    }
+    return context;
+  }
+
+  public static void clearCurrentContext() {
+    clearContext(CURRENT_TASK_ID.get());
+  }
+
+  public static void clearContext(final String taskId) {
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    Object context = null;
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null) {
+        context = detachContext(task);
+      }
+    }
+    closeContext(context);
+    clearCurrentUnretainedContext(taskId);
+  }
+
+  private static ContextReservation reserveContextSlot(final Task currentTask) 
{
+    if (RETAINED_CONTEXT_COUNT.get() < MAX_CONTEXTS) {
+      return new ContextReservation(true, null);
+    }
+    for (final Task task : TASKS.values()) {
+      if (task != currentTask && task.conversionContext != null && task.state 
!= State.RUNNING) {
+        return new ContextReservation(true, detachContext(task));
+      }
+    }
+    return new ContextReservation(false, null);
+  }
+
+  private static final class ContextReservation {
+    private final boolean slotAvailable;
+    private final Object evictedContext;
+
+    private ContextReservation(final boolean slotAvailable, final Object 
evictedContext) {
+      this.slotAvailable = slotAvailable;
+      this.evictedContext = evictedContext;
+    }
+  }
+
+  private static void retainContext(final Task task, final Object context) {
+    task.conversionContext = context;
+    RETAINED_CONTEXT_COUNT.incrementAndGet();
+  }
+
+  private static Object detachContext(final Task task) {
+    final Object context = task.conversionContext;
+    if (context != null) {
+      task.conversionContext = null;
+      RETAINED_CONTEXT_COUNT.decrementAndGet();
+    }
+    return context;
+  }
+
+  private static void closeContext(final Object context) {
+    if (!(context instanceof AutoCloseable)) {
+      return;
+    }
+    try {
+      ((AutoCloseable) context).close();
+    } catch (final Exception e) {
+      LOGGER.warn("Failed to close pipe TsFile conversion context.", e);
+    }
+  }
+
+  public static void leave() {
+    clearCurrentUnretainedContext(CURRENT_TASK_ID.get());
+    CURRENT_TASK_ID.remove();
+    CURRENT_TYPE_MISMATCH.remove();
+  }
+
+  private static void clearCurrentUnretainedContext(final String taskId) {
+    if (taskId == null) {
+      return;
+    }
+    final UnretainedContext context = CURRENT_UNRETAINED_CONTEXT.get();
+    if (context != null && taskId.equals(context.taskId)) {
+      CURRENT_UNRETAINED_CONTEXT.remove();
+      closeContext(context.context);
+    }
+  }
+
+  public static void markTypeMismatchDetected() {
+    CURRENT_TYPE_MISMATCH.set(true);
+    final String taskId = CURRENT_TASK_ID.get();
+    if (taskId == null) {
+      return;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null && task.state != State.SUCCESS && task.state != 
State.FAILED) {
+        task.typeMismatchDetected = true;
+        task.status = null;
+        task.retrySealAllowed = false;
+        task.state = State.RUNNING;
+      }
+    }
+  }
+
+  public static boolean isTypeMismatchDetected(final String taskId) {
+    if (taskId == null || taskId.isEmpty()) {
+      return Boolean.TRUE.equals(CURRENT_TYPE_MISMATCH.get());
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      return task != null && task.typeMismatchDetected;
+    }
+  }
+
+  public static void markReceiverOwned(final String taskId) {
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null) {
+        task.receiverOwned = true;
+        task.retrySealAllowed = false;
+      }
+    }
+  }
+
+  public static void markPending(final String taskId) {
+    update(taskId, State.PENDING, null);
+  }
+
+  /**
+   * Moves a locally running task back to pending before its file becomes 
visible to active load.
+   * The caller must invoke this before moving the file so an active-load 
worker cannot be running
+   * the same task concurrently.
+   */
+  public static void prepareForActiveLoad(final String taskId) {
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null && task.state != State.SUCCESS && task.state != 
State.FAILED) {
+        task.status = null;
+        task.retrySealAllowed = false;
+        task.state = State.PENDING;
+      }
+    }
+  }
+
+  public static void markRunning(final String taskId) {
+    update(taskId, State.RUNNING, null);
+  }
+
+  public static void markPaused(final String taskId, final TSStatus status) {
+    update(taskId, State.PAUSED, status);
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null && !task.receiverOwned && task.state == State.PAUSED) {
+        // A retry claims this same task and reuses its retained parser 
checkpoint.
+        task.retrySealAllowed = true;
+      }
+    }
+  }
+
+  public static void markRetryable(final String taskId, final TSStatus status) 
{
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null
+          && !task.receiverOwned
+          && task.state != State.SUCCESS
+          && task.state != State.FAILED) {
+        task.status = status;
+        task.retrySealAllowed = true;
+        task.state = State.PAUSED;
+      }
+    }
+  }
+
+  public static void markSuccess(final String taskId) {
+    complete(taskId, State.SUCCESS, new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+  }
+
+  public static void markFailed(final String taskId, final TSStatus status) {
+    complete(taskId, State.FAILED, status);
+  }
+
+  private static void update(final String taskId, final State state, final 
TSStatus status) {
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task == null) {
+        return;
+      }
+      if (task.state == State.SUCCESS || task.state == State.FAILED) {
+        return;
+      }
+      if (state == State.PENDING && task.state == State.RUNNING) {
+        return;
+      }
+      task.status = status;
+      if (state == State.PENDING || state == State.RUNNING) {
+        task.retrySealAllowed = false;
+      }
+      task.state = state;
+    }
+  }
+
+  private static void complete(final String taskId, final State state, final 
TSStatus status) {
+    if (taskId == null || taskId.isEmpty()) {
+      return;
+    }
+    Object context = null;
+    synchronized (TASKS) {
+      final Task task = TASKS.get(taskId);
+      if (task != null && task.state != State.SUCCESS && task.state != 
State.FAILED) {
+        task.status = status;
+        task.retrySealAllowed = false;
+        context = detachContext(task);
+        task.state = state;
+      }
+    }
+    closeContext(context);
+    clearCurrentUnretainedContext(taskId);
+  }
+
+  @VisibleForTesting
+  static int getMaxRetainedContextCount() {
+    return MAX_CONTEXTS;
+  }
+
+  @VisibleForTesting
+  static int getRetainedContextCount() {
+    return RETAINED_CONTEXT_COUNT.get();
+  }
+
+  private static void evictCompletedTasksIfNecessary() {
+    final Iterator<Map.Entry<String, Task>> iterator = 
TASKS.entrySet().iterator();
+    while (iterator.hasNext()) {
+      final Task task = iterator.next().getValue();
+      if (task.state == State.SUCCESS || task.state == State.FAILED) {
+        iterator.remove();
+        if (TASKS.size() < MAX_TASKS) {
+          return;
+        }
+      }
+    }
+  }
+
+  private static boolean hasTaskCapacity() {
+    if (TASKS.size() >= MAX_TASKS) {
+      evictCompletedTasksIfNecessary();
+    }
+    return TASKS.size() < MAX_TASKS;
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
index 3bdadc04cf4..e4e7fea03e0 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
@@ -19,7 +19,10 @@
 
 package org.apache.iotdb.db.pipe.sink;
 
+import org.apache.iotdb.commons.consensus.index.impl.IoTProgressIndex;
 import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
+import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.request.PipeRequestType;
 import 
org.apache.iotdb.commons.pipe.sink.payload.thrift.response.PipeTransferFilePieceResp;
 import 
org.apache.iotdb.db.pipe.processor.twostage.exchange.payload.CombineRequest;
@@ -54,6 +57,7 @@ import org.apache.tsfile.write.record.Tablet;
 import org.apache.tsfile.write.schema.MeasurementSchema;
 import org.junit.Assert;
 import org.junit.Test;
+import org.mockito.Mockito;
 
 import java.io.DataOutputStream;
 import java.io.IOException;
@@ -522,6 +526,32 @@ public class PipeDataNodeThriftRequestTest {
     Assert.assertTrue(deserializeReq.shouldWaitForSchemaBeforeLoad());
   }
 
+  @Test
+  public void 
testPipeTransferTsFileSealConversionTaskIdDistinguishesProgressIndexes() {
+    final CommitterKey committerKey = new CommitterKey("pipe", 1L, 1, 0);
+    final EnrichedEvent firstEvent = Mockito.mock(EnrichedEvent.class);
+    Mockito.when(firstEvent.getCommitterKey()).thenReturn(committerKey);
+    
Mockito.when(firstEvent.getCommitIds()).thenReturn(Collections.singletonList(1L));
+    Mockito.when(firstEvent.getProgressIndex()).thenReturn(new 
IoTProgressIndex(1, 1L));
+
+    final EnrichedEvent secondEvent = Mockito.mock(EnrichedEvent.class);
+    Mockito.when(secondEvent.getCommitterKey()).thenReturn(committerKey);
+    
Mockito.when(secondEvent.getCommitIds()).thenReturn(Collections.singletonList(1L));
+    Mockito.when(secondEvent.getProgressIndex()).thenReturn(new 
IoTProgressIndex(1, 100L));
+
+    final String firstTaskId =
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(firstEvent), "root.db", 0);
+    Assert.assertEquals(
+        firstTaskId,
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(firstEvent), "root.db", 0));
+    Assert.assertNotEquals(
+        firstTaskId,
+        PipeTransferTsFileSealWithModReq.generateConversionTaskId(
+            "sink-task", Collections.singletonList(secondEvent), "root.db", 
0));
+  }
+
   @Test
   public void testPipeTransferSchemaSnapshotSealReq() throws IOException {
     final String mTreeSnapshotName = "mtree.snapshot";
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSinkTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSinkTest.java
index db6169e9831..fbb985a40b6 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSinkTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/airgap/IoTDBDataRegionAirGapSinkTest.java
@@ -104,7 +104,8 @@ public class IoTDBDataRegionAirGapSinkTest {
       }
 
       
Assert.assertTrue(requestTypes.contains(PipeRequestType.TRANSFER_TS_FILE_PIECE.getType()));
-      
Assert.assertTrue(requestTypes.contains(PipeRequestType.TRANSFER_TS_FILE_SEAL.getType()));
+      Assert.assertTrue(
+          
requestTypes.contains(PipeRequestType.TRANSFER_TS_FILE_SEAL_WITH_MOD.getType()));
       
Assert.assertFalse(requestTypes.contains(PipeRequestType.TRANSFER_TABLET_RAW.getType()));
       
Assert.assertFalse(requestTypes.contains(PipeRequestType.TRANSFER_TABLET_BATCH.getType()));
       Assert.assertEquals(transferredTsFileBytes, sink.rateLimitedBytes.get());
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java
index 9a805971191..f70fbf3610c 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java
@@ -123,6 +123,19 @@ public class ActiveLoadTsFileLoaderTest {
     Assert.assertTrue(pendingQueue.isEmpty());
   }
 
+  @Test
+  public void testPendingEntryRetainsConversionTaskId() {
+    final ActiveLoadPendingQueue pendingQueue = new ActiveLoadPendingQueue();
+    final String tsFilePath = new File(tempDir, 
"task.tsfile").getAbsolutePath();
+    Assert.assertTrue(
+        pendingQueue.enqueue(tsFilePath, tempDir.getAbsolutePath(), true, 
"conversion-task"));
+
+    final ActiveLoadPendingQueue.ActiveLoadEntry entry = 
pendingQueue.dequeueFromPending();
+    Assert.assertNotNull(entry);
+    Assert.assertEquals("conversion-task", entry.getConversionTaskId());
+    pendingQueue.removeFromLoading(tsFilePath);
+  }
+
   private File createTsFileWithCompanionFiles(final String fileName) throws 
Exception {
     final File tsFile = new File(tempDir, fileName);
     Assert.assertTrue(tsFile.createNewFile());
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtilTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtilTest.java
index c0aaa8c9349..0082830d39f 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtilTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadUtilTest.java
@@ -113,8 +113,105 @@ public class ActiveLoadUtilTest {
     }
   }
 
+  @Test
+  public void testDeterministicTransferFailsWhenSourceAndTargetAreMissing() 
throws Exception {
+    final List<File> sourceFiles = createTsFileAndCompanions();
+    for (final File sourceFile : sourceFiles) {
+      Assert.assertTrue(sourceFile.delete());
+    }
+
+    try {
+      ActiveLoadUtil.transferFilesToActiveDir(targetDir, sourceFiles, true, 
"missing-task");
+      Assert.fail("Expected IOException");
+    } catch (final IOException ignored) {
+      // A missing source and missing deterministic target cannot prove a 
durable handoff.
+    }
+  }
+
+  @Test
+  public void testDeterministicTransferIsIdempotentAfterSourcesAreDeleted() 
throws Exception {
+    final List<File> sourceFiles = createTsFileAndCompanions();
+    ActiveLoadUtil.transferFilesToActiveDir(targetDir, sourceFiles, true, 
"retry-task");
+
+    final File transferDir =
+        new File(targetDir, 
ActiveLoadPathHelper.formatPipeTaskTransferDirectoryName("retry-task"));
+    Assert.assertTrue(transferDir.isDirectory());
+    ActiveLoadUtil.transferFilesToActiveDir(targetDir, sourceFiles, true, 
"retry-task");
+
+    final File[] transferDirs = targetDir.listFiles(File::isDirectory);
+    Assert.assertNotNull(transferDirs);
+    Assert.assertEquals(1, transferDirs.length);
+    Assert.assertEquals(transferDir.getAbsolutePath(), 
transferDirs[0].getAbsolutePath());
+    for (final File sourceFile : sourceFiles) {
+      Assert.assertFalse(sourceFile.exists());
+    }
+  }
+
+  @Test
+  public void testDeterministicTransferUsesTaskIdentityWhenFileNameChanges() 
throws Exception {
+    final List<File> firstSourceFiles = 
createTsFileAndCompanions("1-0-0-0.tsfile");
+    ActiveLoadUtil.transferFilesToActiveDir(targetDir, firstSourceFiles, true, 
"stable-task");
+
+    final List<File> retrySourceFiles = 
createTsFileAndCompanions("2-0-0-0.tsfile");
+    ActiveLoadUtil.transferFilesToActiveDir(targetDir, retrySourceFiles, true, 
"stable-task");
+
+    final File transferDir =
+        new File(
+            targetDir, 
ActiveLoadPathHelper.formatPipeTaskTransferDirectoryName("stable-task"));
+    Assert.assertTrue(new File(transferDir, "1-0-0-0.tsfile").exists());
+    Assert.assertFalse(new File(transferDir, "2-0-0-0.tsfile").exists());
+    for (final File sourceFile : retrySourceFiles) {
+      Assert.assertFalse(sourceFile.exists());
+    }
+  }
+
+  @Test
+  public void testIncompleteDeterministicTargetIsNotDeleted() throws Exception 
{
+    final List<File> sourceFiles = createTsFileAndCompanions();
+    final File transferDir =
+        new File(
+            targetDir, 
ActiveLoadPathHelper.formatPipeTaskTransferDirectoryName("partial-task"));
+    Assert.assertTrue(transferDir.mkdirs());
+    Assert.assertTrue(new File(transferDir, 
sourceFiles.get(0).getName()).createNewFile());
+
+    try {
+      ActiveLoadUtil.transferFilesToActiveDir(targetDir, sourceFiles, true, 
"partial-task");
+      Assert.fail("Expected IOException");
+    } catch (final IOException ignored) {
+      // Do not remove a target which may already be visible to active load.
+    }
+    Assert.assertTrue(transferDir.exists());
+    for (final File sourceFile : sourceFiles) {
+      Assert.assertTrue(sourceFile.exists());
+    }
+  }
+
+  @Test
+  public void testDeterministicTargetWithOnlyTsFileIsIncomplete() throws 
Exception {
+    final List<File> sourceFiles = createTsFileAndCompanions();
+    final File transferDir =
+        new File(targetDir, 
ActiveLoadPathHelper.formatPipeTaskTransferDirectoryName("partial-ts"));
+    Assert.assertTrue(transferDir.mkdirs());
+    final File tsFile = sourceFiles.get(sourceFiles.size() - 1);
+    Files.copy(tsFile.toPath(), new File(transferDir, 
tsFile.getName()).toPath());
+
+    try {
+      ActiveLoadUtil.transferFilesToActiveDir(targetDir, sourceFiles, true, 
"partial-ts");
+      Assert.fail("Expected IOException");
+    } catch (final IOException ignored) {
+      // expected
+    }
+    for (final File sourceFile : sourceFiles) {
+      Assert.assertTrue(sourceFile.exists());
+    }
+  }
+
   private List<File> createTsFileAndCompanions() throws Exception {
-    final File tsFile = new File(sourceDir, "1-0-0-0.tsfile");
+    return createTsFileAndCompanions("1-0-0-0.tsfile");
+  }
+
+  private List<File> createTsFileAndCompanions(final String tsFileName) throws 
Exception {
+    final File tsFile = new File(sourceDir, tsFileName);
     final File resourceFile = new File(tsFile.getAbsolutePath() + 
TsFileResource.RESOURCE_SUFFIX);
     final File modsFile = new File(tsFile.getAbsolutePath() + 
ModificationFile.FILE_SUFFIX);
     final List<File> sourceFiles = Arrays.asList(resourceFile, modsFile, 
tsFile);
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManagerTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManagerTest.java
new file mode 100644
index 00000000000..1c0c4b8fdde
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/PipeTsFileConversionTaskManagerTest.java
@@ -0,0 +1,381 @@
+/*
+ * 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.storageengine.load.converter;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+public class PipeTsFileConversionTaskManagerTest {
+
+  @Test
+  public void testDuplicateStatusRespectsTakeoverMode() {
+    final String synchronousTaskId = "sync-" + System.nanoTime();
+    PipeTsFileConversionTaskManager.registerIfAbsent(synchronousTaskId, false);
+    Assert.assertEquals(
+        
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+        PipeTsFileConversionTaskManager.getDuplicateStatus(synchronousTaskId, 
false).getCode());
+    final TSStatus pausedStatus =
+        new 
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+            .setMessage("paused conversion");
+    PipeTsFileConversionTaskManager.markPaused(synchronousTaskId, 
pausedStatus);
+    Assert.assertSame(
+        pausedStatus, 
PipeTsFileConversionTaskManager.getDuplicateStatus(synchronousTaskId, false));
+    PipeTsFileConversionTaskManager.markSuccess(synchronousTaskId);
+    Assert.assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+        PipeTsFileConversionTaskManager.getDuplicateStatus(synchronousTaskId, 
false).getCode());
+
+    final String asynchronousTaskId = "async-" + System.nanoTime();
+    PipeTsFileConversionTaskManager.registerIfAbsent(asynchronousTaskId, true);
+    PipeTsFileConversionTaskManager.markPaused(
+        asynchronousTaskId,
+        new 
TSStatus(TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()));
+    Assert.assertEquals(
+        
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+        PipeTsFileConversionTaskManager.getDuplicateStatus(asynchronousTaskId, 
true).getCode());
+    PipeTsFileConversionTaskManager.markReceiverOwned(asynchronousTaskId);
+    Assert.assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+        PipeTsFileConversionTaskManager.getDuplicateStatus(asynchronousTaskId, 
true).getCode());
+  }
+
+  @Test
+  public void testContextIsRetainedUntilTerminalState() {
+    final String taskId = "context-" + System.nanoTime();
+    final AtomicBoolean closed = new AtomicBoolean(false);
+    PipeTsFileConversionTaskManager.registerIfAbsent(taskId, true);
+    PipeTsFileConversionTaskManager.enter(taskId);
+    try {
+      final AutoCloseable context =
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () -> (AutoCloseable) () -> closed.set(true));
+      PipeTsFileConversionTaskManager.markPaused(
+          taskId, new 
TSStatus(TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()));
+      Assert.assertSame(
+          context,
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () -> (AutoCloseable) () -> closed.set(true)));
+      Assert.assertFalse(closed.get());
+      PipeTsFileConversionTaskManager.markSuccess(taskId);
+      Assert.assertTrue(closed.get());
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
+  }
+
+  @Test
+  public void testPrepareForActiveLoadDoesNotDowngradeTerminalTask() {
+    final String taskId = "handoff-" + System.nanoTime();
+    PipeTsFileConversionTaskManager.registerIfAbsent(taskId, true);
+    PipeTsFileConversionTaskManager.markRunning(taskId);
+    PipeTsFileConversionTaskManager.prepareForActiveLoad(taskId);
+    Assert.assertEquals(
+        PipeTsFileConversionTaskManager.State.PENDING,
+        PipeTsFileConversionTaskManager.get(taskId).getState());
+
+    PipeTsFileConversionTaskManager.markSuccess(taskId);
+    PipeTsFileConversionTaskManager.prepareForActiveLoad(taskId);
+    Assert.assertEquals(
+        PipeTsFileConversionTaskManager.State.SUCCESS,
+        PipeTsFileConversionTaskManager.get(taskId).getState());
+
+    final AtomicBoolean lateContextClosed = new AtomicBoolean(false);
+    PipeTsFileConversionTaskManager.enter(taskId);
+    try {
+      PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+          () -> (AutoCloseable) () -> lateContextClosed.set(true));
+      Assert.assertEquals(0, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
+    Assert.assertTrue(lateContextClosed.get());
+  }
+
+  @Test
+  public void testRegisterAndGetDuplicateStatusClaimsTaskAtomically() throws 
Exception {
+    final String taskId = "atomic-register-" + System.nanoTime();
+    final int concurrency = 16;
+    final ExecutorService executor = Executors.newFixedThreadPool(concurrency);
+    final CountDownLatch ready = new CountDownLatch(concurrency);
+    final CountDownLatch start = new CountDownLatch(1);
+    final List<Future<TSStatus>> results = new ArrayList<>();
+    try {
+      for (int i = 0; i < concurrency; i++) {
+        results.add(
+            executor.submit(
+                () -> {
+                  ready.countDown();
+                  start.await();
+                  return 
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(
+                      taskId, false);
+                }));
+      }
+
+      Assert.assertTrue(ready.await(10, TimeUnit.SECONDS));
+      start.countDown();
+      int claimantCount = 0;
+      for (final Future<TSStatus> result : results) {
+        final TSStatus status = result.get(10, TimeUnit.SECONDS);
+        if (status == null) {
+          claimantCount++;
+        } else {
+          Assert.assertEquals(
+              
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+              status.getCode());
+        }
+      }
+      Assert.assertEquals(1, claimantCount);
+    } finally {
+      start.countDown();
+      executor.shutdownNow();
+      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+      PipeTsFileConversionTaskManager.markSuccess(taskId);
+    }
+  }
+
+  @Test
+  public void testFailedHandoffCanBeReclaimedOnlyOnce() {
+    final String taskId = "retryable-handoff-" + System.nanoTime();
+    final TSStatus pausedStatus =
+        new 
TSStatus(TSStatusCode.LOAD_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+            .setMessage("active-load handoff failed");
+
+    
Assert.assertNull(PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId,
 true));
+    PipeTsFileConversionTaskManager.markRunning(taskId);
+    PipeTsFileConversionTaskManager.markRetryable(taskId, pausedStatus);
+
+    
Assert.assertNull(PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId,
 true));
+    Assert.assertEquals(
+        PipeTsFileConversionTaskManager.State.PENDING,
+        PipeTsFileConversionTaskManager.get(taskId).getState());
+    Assert.assertNull(PipeTsFileConversionTaskManager.get(taskId).getStatus());
+    Assert.assertEquals(
+        
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode(),
+        PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId, 
true).getCode());
+
+    PipeTsFileConversionTaskManager.markRetryable(taskId, pausedStatus);
+    PipeTsFileConversionTaskManager.markReceiverOwned(taskId);
+    Assert.assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+        PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId, 
true).getCode());
+    PipeTsFileConversionTaskManager.markSuccess(taskId);
+  }
+
+  @Test
+  public void testReceiverOwnedStatusRespectsTakeoverMode() {
+    final String asynchronousTaskId = "receiver-owned-async-" + 
System.nanoTime();
+    final TSStatus failedStatus =
+        new TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode())
+            .setMessage("active load failed after handoff");
+    Assert.assertNull(
+        
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(asynchronousTaskId,
 true));
+    PipeTsFileConversionTaskManager.markReceiverOwned(asynchronousTaskId);
+    PipeTsFileConversionTaskManager.markFailed(asynchronousTaskId, 
failedStatus);
+    Assert.assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+        
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(asynchronousTaskId,
 true)
+            .getCode());
+
+    final String synchronousTaskId = "receiver-owned-sync-" + 
System.nanoTime();
+    final TSStatus pausedStatus =
+        new 
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode())
+            .setMessage("conversion paused after handoff");
+    Assert.assertNull(
+        
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(synchronousTaskId,
 false));
+    PipeTsFileConversionTaskManager.markReceiverOwned(synchronousTaskId);
+    PipeTsFileConversionTaskManager.markPaused(synchronousTaskId, 
pausedStatus);
+    Assert.assertSame(
+        pausedStatus,
+        
PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(synchronousTaskId,
 false));
+    PipeTsFileConversionTaskManager.markSuccess(synchronousTaskId);
+  }
+
+  @Test
+  public void testTerminalStateAndStatusTransitionAtomically() throws 
Exception {
+    final int taskCount = 256;
+    final String[] taskIds = new String[taskCount];
+    final TSStatus[] failedStatuses = new TSStatus[taskCount];
+    for (int i = 0; i < taskCount; i++) {
+      taskIds[i] = "terminal-race-" + System.nanoTime() + '-' + i;
+      failedStatuses[i] =
+          new 
TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode()).setMessage(taskIds[i]);
+      PipeTsFileConversionTaskManager.registerIfAbsent(taskIds[i], false);
+    }
+
+    final ExecutorService executor = Executors.newFixedThreadPool(2);
+    final CyclicBarrier barrier = new CyclicBarrier(3);
+    try {
+      final Future<Void> successFuture =
+          executor.submit(
+              () -> {
+                for (final String taskId : taskIds) {
+                  barrier.await();
+                  PipeTsFileConversionTaskManager.markSuccess(taskId);
+                  barrier.await();
+                }
+                return null;
+              });
+      final Future<Void> failureFuture =
+          executor.submit(
+              () -> {
+                for (int i = 0; i < taskCount; i++) {
+                  barrier.await();
+                  PipeTsFileConversionTaskManager.markFailed(taskIds[i], 
failedStatuses[i]);
+                  barrier.await();
+                }
+                return null;
+              });
+
+      for (int i = 0; i < taskCount; i++) {
+        barrier.await(10, TimeUnit.SECONDS);
+        barrier.await(10, TimeUnit.SECONDS);
+        final PipeTsFileConversionTaskManager.Task task =
+            PipeTsFileConversionTaskManager.get(taskIds[i]);
+        if (task.getState() == PipeTsFileConversionTaskManager.State.SUCCESS) {
+          Assert.assertEquals(
+              TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
task.getStatus().getCode());
+        } else {
+          Assert.assertEquals(PipeTsFileConversionTaskManager.State.FAILED, 
task.getState());
+          Assert.assertSame(failedStatuses[i], task.getStatus());
+        }
+      }
+      successFuture.get(10, TimeUnit.SECONDS);
+      failureFuture.get(10, TimeUnit.SECONDS);
+    } finally {
+      executor.shutdownNow();
+      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+    }
+  }
+
+  @Test
+  public void testRetainedContextCountHasHardLimit() {
+    Assert.assertEquals(0, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+    final int maxContextCount = 
PipeTsFileConversionTaskManager.getMaxRetainedContextCount();
+    final List<String> retainedTaskIds = new ArrayList<>();
+    final List<AtomicBoolean> retainedContextClosed = new ArrayList<>();
+    final String overflowTaskId = "context-overflow-" + System.nanoTime();
+
+    try {
+      for (int i = 0; i < maxContextCount; i++) {
+        final String taskId = "context-retained-" + System.nanoTime() + '-' + 
i;
+        final AtomicBoolean closed = new AtomicBoolean(false);
+        retainedTaskIds.add(taskId);
+        retainedContextClosed.add(closed);
+        PipeTsFileConversionTaskManager.registerIfAbsent(taskId, true);
+        PipeTsFileConversionTaskManager.markRunning(taskId);
+        PipeTsFileConversionTaskManager.enter(taskId);
+        PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+            () -> (AutoCloseable) () -> closed.set(true));
+        PipeTsFileConversionTaskManager.leave();
+        Assert.assertFalse(closed.get());
+      }
+      Assert.assertEquals(
+          maxContextCount, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+
+      final AtomicBoolean overflowContextClosed = new AtomicBoolean(false);
+      PipeTsFileConversionTaskManager.registerIfAbsent(overflowTaskId, true);
+      PipeTsFileConversionTaskManager.markRunning(overflowTaskId);
+      PipeTsFileConversionTaskManager.enter(overflowTaskId);
+      final AutoCloseable overflowContext =
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () -> (AutoCloseable) () -> overflowContextClosed.set(true));
+      Assert.assertSame(
+          overflowContext,
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () -> (AutoCloseable) () -> Assert.fail("must reuse the 
unretained context")));
+      Assert.assertEquals(
+          maxContextCount, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+      Assert.assertFalse(overflowContextClosed.get());
+
+      PipeTsFileConversionTaskManager.leave();
+      Assert.assertTrue(overflowContextClosed.get());
+      Assert.assertEquals(
+          maxContextCount, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+      PipeTsFileConversionTaskManager.markSuccess(overflowTaskId);
+      retainedTaskIds.forEach(PipeTsFileConversionTaskManager::markSuccess);
+    }
+
+    for (final AtomicBoolean closed : retainedContextClosed) {
+      Assert.assertTrue(closed.get());
+    }
+    Assert.assertEquals(0, 
PipeTsFileConversionTaskManager.getRetainedContextCount());
+  }
+
+  @Test
+  public void testPausedTaskRetryReusesCheckpointWithoutCreatingTask() {
+    final String taskId = "paused-retry-" + System.nanoTime();
+    final AtomicBoolean closed = new AtomicBoolean(false);
+    
Assert.assertNull(PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId,
 false));
+    PipeTsFileConversionTaskManager.enter(taskId);
+    final Object context;
+    try {
+      context =
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () -> (AutoCloseable) () -> closed.set(true));
+      PipeTsFileConversionTaskManager.markPaused(
+          taskId,
+          new 
TSStatus(TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION.getStatusCode()));
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
+
+    
Assert.assertNull(PipeTsFileConversionTaskManager.registerAndGetDuplicateStatus(taskId,
 false));
+    PipeTsFileConversionTaskManager.enter(taskId);
+    try {
+      Assert.assertSame(
+          context,
+          PipeTsFileConversionTaskManager.getOrCreateCurrentContext(
+              () ->
+                  (AutoCloseable)
+                      () -> Assert.fail("a new conversion context must not be 
created")));
+      PipeTsFileConversionTaskManager.markSuccess(taskId);
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
+    Assert.assertTrue(closed.get());
+  }
+
+  @Test
+  public void testLegacyTaskCanReportTypeMismatchWithoutTaskId() {
+    PipeTsFileConversionTaskManager.enter(null);
+    try {
+      
Assert.assertFalse(PipeTsFileConversionTaskManager.isTypeMismatchDetected(null));
+      PipeTsFileConversionTaskManager.markTypeMismatchDetected();
+      
Assert.assertTrue(PipeTsFileConversionTaskManager.isTypeMismatchDetected(null));
+    } finally {
+      PipeTsFileConversionTaskManager.leave();
+    }
+    
Assert.assertFalse(PipeTsFileConversionTaskManager.isTypeMismatchDetected(null));
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java
index f5727634d19..ac01c68706d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java
@@ -124,6 +124,13 @@ public class PipeSinkConstant {
       "sink.exception.data.convert-on-type-mismatch";
   public static final boolean 
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_DEFAULT_VALUE =
       true;
+  public static final String
+      CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY =
+          
"connector.exception.data.convert-on-type-mismatch.tsfile.async-load";
+  public static final String 
SINK_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY =
+      "sink.exception.data.convert-on-type-mismatch.tsfile.async-load";
+  public static final boolean
+      
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_DEFAULT_VALUE
 = true;
 
   public static final String CONNECTOR_EXCEPTION_CONFLICT_RESOLVE_STRATEGY_KEY 
=
       "connector.exception.conflict.resolve-strategy";
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/plugin/env/PipeTaskSinkRuntimeEnvironment.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/plugin/env/PipeTaskSinkRuntimeEnvironment.java
index b8382891348..26081d9c78a 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/plugin/env/PipeTaskSinkRuntimeEnvironment.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/plugin/env/PipeTaskSinkRuntimeEnvironment.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.commons.pipe.config.plugin.env;
 
 public class PipeTaskSinkRuntimeEnvironment extends PipeTaskRuntimeEnvironment 
{
   private String attributeSortedString;
+  private String sinkTaskId;
 
   public PipeTaskSinkRuntimeEnvironment(
       final String pipeName, final long creationTime, final int regionId) {
@@ -34,4 +35,12 @@ public class PipeTaskSinkRuntimeEnvironment extends 
PipeTaskRuntimeEnvironment {
   public void setAttributeSortedString(String attributeSortedString) {
     this.attributeSortedString = attributeSortedString;
   }
+
+  public String getSinkTaskId() {
+    return sinkTaskId;
+  }
+
+  public void setSinkTaskId(final String sinkTaskId) {
+    this.sinkTaskId = sinkTaskId;
+  }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
index cb293de1444..44f04706eae 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/receiver/IoTDBFileReceiver.java
@@ -780,6 +780,7 @@ public abstract class IoTDBFileReceiver implements 
IoTDBReceiver {
 
   protected final TPipeTransferResp handleTransferFileSealV2(final 
PipeTransferFileSealReqV2 req) {
     final List<String> fileNames = req.getFileNames();
+    TSStatus loadStatus = null;
     try {
       final List<File> files =
           fileNames.stream()
@@ -844,6 +845,7 @@ public abstract class IoTDBFileReceiver implements 
IoTDBReceiver {
           
files.stream().map(File::getAbsolutePath).collect(Collectors.toList());
 
       final TSStatus status = loadFileV2(req, fileAbsolutePaths);
+      loadStatus = status;
       if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
         LOGGER.debug(
             "Receiver id = {}: Seal file {} successfully.", receiverId.get(), 
fileAbsolutePaths);
@@ -878,10 +880,22 @@ public abstract class IoTDBFileReceiver implements 
IoTDBReceiver {
       closeCurrentWritingFileWriter(false);
       // Clear the directory instead of only deleting the referenced files in 
seal request
       // to avoid previously undeleted file being redundant when transferring 
multi files
-      
IoTDBReceiverAgent.cleanPipeReceiverDir(receiverFileDirWithIdSuffix.get());
+      if (shouldDeleteSealedFilesOnFailure(req, loadStatus)) {
+        
IoTDBReceiverAgent.cleanPipeReceiverDir(receiverFileDirWithIdSuffix.get());
+      }
     }
   }
 
+  /**
+   * Decides whether files in the receiver's staging directory should be 
removed after a V2 seal.
+   * The default keeps the historical behavior. A receiver may retain files 
when a conversion task
+   * is retryable so that a sender retry can resume the same task without 
losing its input.
+   */
+  protected boolean shouldDeleteSealedFilesOnFailure(
+      final PipeTransferFileSealReqV2 req, final TSStatus loadStatus) {
+    return true;
+  }
+
   private TPipeTransferResp checkNonFinalFileSeal(
       final File file, final String fileName, final long fileLength) throws 
IOException {
     if (!file.exists()) {
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBSink.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBSink.java
index 1d929411929..7914bb65281 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBSink.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/sink/protocol/IoTDBSink.java
@@ -72,6 +72,8 @@ import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CON
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_CONFLICT_RETRY_MAX_TIME_SECONDS_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_DEFAULT_VALUE;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_KEY;
+import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_DEFAULT_VALUE;
+import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_OTHERS_RECORD_IGNORED_DATA_DEFAULT_VALUE;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_OTHERS_RECORD_IGNORED_DATA_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_EXCEPTION_OTHERS_RETRY_MAX_TIME_SECONDS_DEFAULT_VALUE;
@@ -118,6 +120,7 @@ import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SIN
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_CONFLICT_RESOLVE_STRATEGY_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_CONFLICT_RETRY_MAX_TIME_SECONDS_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_KEY;
+import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_OTHERS_RECORD_IGNORED_DATA_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_EXCEPTION_OTHERS_RETRY_MAX_TIME_SECONDS_KEY;
 import static 
org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_FORMAT_KEY;
@@ -176,9 +179,12 @@ public abstract class IoTDBSink implements PipeConnector, 
PipeConnectorWithEvent
   protected boolean shouldReceiverConvertOnTypeMismatch =
       CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_DEFAULT_VALUE;
 
+  protected boolean shouldAsyncLoadTsFileOnTypeMismatch =
+      
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_DEFAULT_VALUE;
   private final AtomicLong totalUncompressedSize = new AtomicLong(0);
   private final AtomicLong totalCompressedSize = new AtomicLong(0);
   protected String attributeSortedString;
+  protected String sinkTaskId;
   protected Timer compressionTimer;
   protected boolean isRealtimeFirst;
 
@@ -378,6 +384,7 @@ public abstract class IoTDBSink implements PipeConnector, 
PipeConnectorWithEvent
     if (environment instanceof PipeTaskSinkRuntimeEnvironment) {
       attributeSortedString =
           ((PipeTaskSinkRuntimeEnvironment) 
environment).getAttributeSortedString();
+      sinkTaskId = ((PipeTaskSinkRuntimeEnvironment) 
environment).getSinkTaskId();
     }
 
     nodeUrls.clear();
@@ -468,10 +475,20 @@ public abstract class IoTDBSink implements PipeConnector, 
PipeConnectorWithEvent
                 CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_KEY,
                 SINK_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_KEY),
             CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_DEFAULT_VALUE);
+    shouldAsyncLoadTsFileOnTypeMismatch =
+        parameters.getBooleanOrDefault(
+            Arrays.asList(
+                
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY,
+                
SINK_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY),
+            
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_DEFAULT_VALUE);
     LOGGER.info(
         "IoTDBSink {} = {}",
         CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_KEY,
         shouldReceiverConvertOnTypeMismatch);
+    LOGGER.info(
+        "IoTDBSink {} = {}",
+        
CONNECTOR_EXCEPTION_DATA_CONVERT_ON_TYPE_MISMATCH_TSFILE_ASYNC_LOAD_KEY,
+        shouldAsyncLoadTsFileOnTypeMismatch);
     isRealtimeFirst =
         parameters.getBooleanOrDefault(
             Arrays.asList(
@@ -638,6 +655,14 @@ public abstract class IoTDBSink implements PipeConnector, 
PipeConnectorWithEvent
     return shouldWaitForSchemaBeforeLoad;
   }
 
+  public boolean shouldAsyncLoadTsFileOnTypeMismatch() {
+    return shouldAsyncLoadTsFileOnTypeMismatch;
+  }
+
+  public String getSinkTaskId() {
+    return sinkTaskId;
+  }
+
   public void setTabletBatchSizeHistogram(Histogram tabletBatchSizeHistogram) {
     // do nothing by default
   }

Reply via email to