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

rong 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 4d290b54479 Subscription: avoid the file reading when the memory is 
not enough for tsfile slicing & implement all-out effort global timeout control 
& support client resume from breakpoint for tsfile consumption (#13992) (#14072)
4d290b54479 is described below

commit 4d290b54479a4164f0ad6925a93553e4b12a4011
Author: V_Galaxy <[email protected]>
AuthorDate: Wed Nov 13 16:03:37 2024 +0800

    Subscription: avoid the file reading when the memory is not enough for 
tsfile slicing & implement all-out effort global timeout control & support 
client resume from breakpoint for tsfile consumption (#13992) (#14072)
---
 .../SubscriptionPipeTimeoutException.java          |   2 +-
 ....java => SubscriptionPollTimeoutException.java} |  12 +-
 ...tion.java => SubscriptionTimeoutException.java} |  14 +-
 .../payload/poll/SubscriptionPollRequest.java      |   2 +-
 .../consumer/SubscriptionConsumer.java             | 173 ++++++++-----
 .../SubscriptionExecutorServiceManager.java        |  11 +-
 .../consumer/SubscriptionProvider.java             |  13 +-
 .../subtask/processor/PipeProcessorSubtask.java    |   4 +-
 .../evolvable/batch/PipeTabletEventBatch.java      |   5 +-
 .../common/tsfile/PipeTsFileInsertionEvent.java    |  55 +++--
 .../db/pipe/resource/memory/PipeMemoryManager.java |   6 +
 .../agent/SubscriptionReceiverAgent.java           |  13 +
 .../broker/SubscriptionPrefetchingQueue.java       |  22 +-
 .../broker/SubscriptionPrefetchingTabletQueue.java |  22 +-
 .../broker/SubscriptionPrefetchingTsFileQueue.java |  35 +--
 .../db/subscription/event/SubscriptionEvent.java   |   6 +-
 .../event/batch/SubscriptionPipeEventBatch.java    |   6 +-
 .../batch/SubscriptionPipeTabletEventBatch.java    |   6 +-
 .../batch/SubscriptionPipeTsFileEventBatch.java    |   8 +-
 .../SubscriptionEventExtendableResponse.java       |  10 -
 .../event/response/SubscriptionEventResponse.java  |   4 +-
 .../response/SubscriptionEventSingleResponse.java  |   2 +-
 .../response/SubscriptionEventTabletResponse.java  |  10 +
 .../response/SubscriptionEventTsFileResponse.java  |  73 +++++-
 .../receiver/SubscriptionReceiver.java             |   2 +
 .../receiver/SubscriptionReceiverV1.java           | 271 +++++++++++----------
 .../execution/SubscriptionSubtaskExecutor.java     |  50 ++++
 ...utor.java => SubscriptionSubtaskScheduler.java} |  22 +-
 .../task/subtask/SubscriptionConnectorSubtask.java |  44 +++-
 .../subtask/SubscriptionReceiverSubtask.java}      |  15 +-
 .../apache/iotdb/commons/conf/CommonConfig.java    |  22 +-
 .../iotdb/commons/conf/CommonDescriptor.java       |  11 +-
 .../agent/task/execution/PipeSubtaskExecutor.java  |  32 ++-
 .../task/subtask/PipeAbstractConnectorSubtask.java |  12 +-
 .../pipe/agent/task/subtask/PipeSubtask.java       |   4 +-
 .../subscription/config/SubscriptionConfig.java    |  14 +-
 36 files changed, 661 insertions(+), 352 deletions(-)

diff --git 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
index f9ecc0f1b8b..26231e48396 100644
--- 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
+++ 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.rpc.subscription.exception;
 
 import java.util.Objects;
 
-public class SubscriptionPipeTimeoutException extends 
SubscriptionRuntimeNonCriticalException {
+public class SubscriptionPipeTimeoutException extends 
SubscriptionTimeoutException {
 
   public SubscriptionPipeTimeoutException(final String message) {
     super(message);
diff --git 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPollTimeoutException.java
similarity index 74%
copy from 
iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
copy to 
iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPollTimeoutException.java
index f9ecc0f1b8b..67c460ecfb3 100644
--- 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
+++ 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPollTimeoutException.java
@@ -21,21 +21,21 @@ package org.apache.iotdb.rpc.subscription.exception;
 
 import java.util.Objects;
 
-public class SubscriptionPipeTimeoutException extends 
SubscriptionRuntimeNonCriticalException {
+public class SubscriptionPollTimeoutException extends 
SubscriptionTimeoutException {
 
-  public SubscriptionPipeTimeoutException(final String message) {
+  public SubscriptionPollTimeoutException(final String message) {
     super(message);
   }
 
-  public SubscriptionPipeTimeoutException(final String message, final 
Throwable cause) {
+  public SubscriptionPollTimeoutException(final String message, final 
Throwable cause) {
     super(message, cause);
   }
 
   @Override
   public boolean equals(final Object obj) {
-    return obj instanceof SubscriptionPipeTimeoutException
-        && Objects.equals(getMessage(), ((SubscriptionPipeTimeoutException) 
obj).getMessage())
-        && Objects.equals(getTimeStamp(), ((SubscriptionPipeTimeoutException) 
obj).getTimeStamp());
+    return obj instanceof SubscriptionPollTimeoutException
+        && Objects.equals(getMessage(), ((SubscriptionPollTimeoutException) 
obj).getMessage())
+        && Objects.equals(getTimeStamp(), ((SubscriptionPollTimeoutException) 
obj).getTimeStamp());
   }
 
   @Override
diff --git 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionTimeoutException.java
similarity index 66%
copy from 
iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
copy to 
iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionTimeoutException.java
index f9ecc0f1b8b..f6471a74aba 100644
--- 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionPipeTimeoutException.java
+++ 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/exception/SubscriptionTimeoutException.java
@@ -21,21 +21,23 @@ package org.apache.iotdb.rpc.subscription.exception;
 
 import java.util.Objects;
 
-public class SubscriptionPipeTimeoutException extends 
SubscriptionRuntimeNonCriticalException {
+public abstract class SubscriptionTimeoutException extends 
SubscriptionRuntimeNonCriticalException {
 
-  public SubscriptionPipeTimeoutException(final String message) {
+  public static String KEYWORD = "TimeoutException";
+
+  public SubscriptionTimeoutException(final String message) {
     super(message);
   }
 
-  public SubscriptionPipeTimeoutException(final String message, final 
Throwable cause) {
+  public SubscriptionTimeoutException(final String message, final Throwable 
cause) {
     super(message, cause);
   }
 
   @Override
   public boolean equals(final Object obj) {
-    return obj instanceof SubscriptionPipeTimeoutException
-        && Objects.equals(getMessage(), ((SubscriptionPipeTimeoutException) 
obj).getMessage())
-        && Objects.equals(getTimeStamp(), ((SubscriptionPipeTimeoutException) 
obj).getTimeStamp());
+    return obj instanceof SubscriptionTimeoutException
+        && Objects.equals(getMessage(), ((SubscriptionTimeoutException) 
obj).getMessage())
+        && Objects.equals(getTimeStamp(), ((SubscriptionTimeoutException) 
obj).getTimeStamp());
   }
 
   @Override
diff --git 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/payload/poll/SubscriptionPollRequest.java
 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/payload/poll/SubscriptionPollRequest.java
index b9337e5779f..3337887b185 100644
--- 
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/payload/poll/SubscriptionPollRequest.java
+++ 
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/subscription/payload/poll/SubscriptionPollRequest.java
@@ -36,7 +36,7 @@ public class SubscriptionPollRequest {
 
   private final transient SubscriptionPollPayload payload;
 
-  private final transient long timeoutMs; // unused now
+  private final transient long timeoutMs;
 
   /** The maximum size, in bytes, for the response payload. */
   private final transient long maxBytes;
diff --git 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionConsumer.java
 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionConsumer.java
index 1bf0aea9c1e..70738a5aece 100644
--- 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionConsumer.java
+++ 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionConsumer.java
@@ -26,8 +26,10 @@ import org.apache.iotdb.rpc.subscription.config.TopicConfig;
 import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionConnectionException;
 import org.apache.iotdb.rpc.subscription.exception.SubscriptionException;
 import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionPipeTimeoutException;
+import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionPollTimeoutException;
 import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionRuntimeCriticalException;
 import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionRuntimeNonCriticalException;
+import 
org.apache.iotdb.rpc.subscription.exception.SubscriptionTimeoutException;
 import org.apache.iotdb.rpc.subscription.payload.poll.ErrorPayload;
 import org.apache.iotdb.rpc.subscription.payload.poll.FileInitPayload;
 import org.apache.iotdb.rpc.subscription.payload.poll.FilePiecePayload;
@@ -111,9 +113,9 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
 
   private final String fileSaveDir;
   private final boolean fileSaveFsync;
+  private final Set<SubscriptionCommitContext> inFlightFilesCommitContextSet = 
new HashSet<>();
 
   private final int thriftMaxFrameSize;
-
   private final int maxPollParallelism;
 
   @SuppressWarnings("java:S3077")
@@ -171,7 +173,6 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     this.fileSaveFsync = builder.fileSaveFsync;
 
     this.thriftMaxFrameSize = builder.thriftMaxFrameSize;
-
     this.maxPollParallelism = builder.maxPollParallelism;
   }
 
@@ -399,33 +400,47 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
   }
 
   private Path getFilePath(
+      final SubscriptionCommitContext commitContext,
       final String topicName,
       final String fileName,
       final boolean allowFileAlreadyExistsException,
       final boolean allowInvalidPathException)
       throws SubscriptionException {
     try {
-      final Path filePath = getFileDir(topicName).resolve(fileName);
-      Files.createFile(filePath);
-      return filePath;
-    } catch (final FileAlreadyExistsException fileAlreadyExistsException) {
-      if (allowFileAlreadyExistsException) {
-        final String suffix = RandomStringGenerator.generate(16);
-        LOGGER.warn(
-            "Detect already existed file {} when polling topic {}, add random 
suffix {} to filename",
-            fileName,
-            topicName,
-            suffix);
-        return getFilePath(topicName, fileName + "." + suffix, false, true);
+      final Path filePath;
+      try {
+        filePath = getFileDir(topicName).resolve(fileName);
+      } catch (final InvalidPathException invalidPathException) {
+        if (allowInvalidPathException) {
+          return getFilePath(commitContext, URLEncoder.encode(topicName), 
fileName, true, false);
+        }
+        throw new SubscriptionRuntimeNonCriticalException(
+            invalidPathException.getMessage(), invalidPathException);
       }
-      throw new SubscriptionRuntimeNonCriticalException(
-          fileAlreadyExistsException.getMessage(), fileAlreadyExistsException);
-    } catch (final InvalidPathException invalidPathException) {
-      if (allowInvalidPathException) {
-        return getFilePath(URLEncoder.encode(topicName), fileName, true, 
false);
+
+      try {
+        Files.createFile(filePath);
+        return filePath;
+      } catch (final FileAlreadyExistsException fileAlreadyExistsException) {
+        if (allowFileAlreadyExistsException) {
+          if (inFlightFilesCommitContextSet.contains(commitContext)) {
+            LOGGER.info(
+                "Detect already existed file {} when polling topic {}, resume 
consumption",
+                fileName,
+                topicName);
+            return filePath;
+          }
+          final String suffix = RandomStringGenerator.generate(16);
+          LOGGER.warn(
+              "Detect already existed file {} when polling topic {}, add 
random suffix {} to filename",
+              fileName,
+              topicName,
+              suffix);
+          return getFilePath(commitContext, topicName, fileName + "." + 
suffix, false, true);
+        }
+        throw new SubscriptionRuntimeNonCriticalException(
+            fileAlreadyExistsException.getMessage(), 
fileAlreadyExistsException);
       }
-      throw new SubscriptionRuntimeNonCriticalException(
-          invalidPathException.getMessage(), invalidPathException);
     } catch (final IOException e) {
       throw new SubscriptionRuntimeNonCriticalException(e.getMessage(), e);
     }
@@ -579,7 +594,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
         final List<SubscriptionMessage> currentMessages = new ArrayList<>();
         try {
           currentResponses.clear();
-          currentResponses = pollInternal(topicNames);
+          currentResponses = pollInternal(topicNames, timer.remainingMs());
           for (final SubscriptionPollResponse response : currentResponses) {
             final short responseType = response.getResponseType();
             if 
(!SubscriptionPollResponseType.isValidatedResponseType(responseType)) {
@@ -676,11 +691,14 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     final SubscriptionCommitContext commitContext = 
response.getCommitContext();
     final String fileName = ((FileInitPayload) 
response.getPayload()).getFileName();
     final String topicName = commitContext.getTopicName();
-    final Path filePath = getFilePath(topicName, fileName, true, true);
+    final Path filePath = getFilePath(commitContext, topicName, fileName, 
true, true);
     final File file = filePath.toFile();
     try (final RandomAccessFile fileWriter = new RandomAccessFile(file, "rw")) 
{
       return Optional.of(pollFileInternal(commitContext, fileName, file, 
fileWriter, timer));
     } catch (final Exception e) {
+      if (!(e instanceof SubscriptionPollTimeoutException)) {
+        inFlightFilesCommitContextSet.remove(commitContext);
+      }
       // construct temporary message to nack
       nack(
           Collections.singletonList(
@@ -696,26 +714,31 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
       final RandomAccessFile fileWriter,
       final PollTimer timer)
       throws IOException, SubscriptionException {
+    long writingOffset = fileWriter.length();
+
     LOGGER.info(
-        "{} start to poll file {} with commit context {}",
+        "{} start to poll file {} with commit context {} at offset {}",
         this,
         file.getAbsolutePath(),
-        commitContext);
+        commitContext,
+        writingOffset);
 
-    long writingOffset = fileWriter.length();
+    fileWriter.seek(writingOffset);
     while (true) {
       timer.update();
-      if (timer.isExpired()) {
-        final String errorMessage =
+      if (timer.isExpired(TIMER_DELTA_MS)) {
+        // resume from breakpoint if timeout happened when polling files
+        inFlightFilesCommitContextSet.add(commitContext);
+        final String message =
             String.format(
-                "timeout while poll file %s with commit context: %s, consumer: 
%s",
-                file.getAbsolutePath(), commitContext, this);
-        LOGGER.warn(errorMessage);
-        throw new SubscriptionRuntimeNonCriticalException(errorMessage);
+                "Timeout occurred when SubscriptionConsumer %s polling file %s 
with commit context %s, record writing offset %s for subsequent poll",
+                this, file.getAbsolutePath(), commitContext, writingOffset);
+        LOGGER.info(message);
+        throw new SubscriptionRuntimeNonCriticalException(message);
       }
 
       final List<SubscriptionPollResponse> responses =
-          pollFileInternal(commitContext, writingOffset);
+          pollFileInternal(commitContext, writingOffset, timer.remainingMs());
 
       // It's agreed that the server will always return at least one response, 
even in case of
       // failure.
@@ -831,6 +854,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
                 commitContext);
 
             // generate subscription message
+            inFlightFilesCommitContextSet.remove(commitContext);
             return new SubscriptionMessage(commitContext, 
file.getAbsolutePath());
           }
         case ERROR:
@@ -839,17 +863,30 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
 
             final String errorMessage = ((ErrorPayload) 
payload).getErrorMessage();
             final boolean critical = ((ErrorPayload) payload).isCritical();
-            LOGGER.warn(
-                "Error occurred when SubscriptionConsumer {} polling file {} 
with commit context {}: {}, critical: {}",
-                this,
-                file.getAbsolutePath(),
-                commitContext,
-                errorMessage,
-                critical);
-            if (critical) {
-              throw new SubscriptionRuntimeCriticalException(errorMessage);
+            if (!critical
+                && Objects.nonNull(errorMessage)
+                && 
errorMessage.contains(SubscriptionTimeoutException.KEYWORD)) {
+              // resume from breakpoint if timeout happened when polling files
+              inFlightFilesCommitContextSet.add(commitContext);
+              final String message =
+                  String.format(
+                      "Timeout occurred when SubscriptionConsumer %s polling 
file %s with commit context %s, record writing offset %s for subsequent poll",
+                      this, file.getAbsolutePath(), commitContext, 
writingOffset);
+              LOGGER.info(message);
+              throw new SubscriptionPollTimeoutException(message);
             } else {
-              throw new SubscriptionRuntimeNonCriticalException(errorMessage);
+              LOGGER.warn(
+                  "Error occurred when SubscriptionConsumer {} polling file {} 
with commit context {}: {}, critical: {}",
+                  this,
+                  file.getAbsolutePath(),
+                  commitContext,
+                  errorMessage,
+                  critical);
+              if (critical) {
+                throw new SubscriptionRuntimeCriticalException(errorMessage);
+              } else {
+                throw new 
SubscriptionRuntimeNonCriticalException(errorMessage);
+              }
             }
           }
         default:
@@ -880,16 +917,6 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
 
     int nextOffset = ((TabletsPayload) 
initialResponse.getPayload()).getNextOffset();
     while (true) {
-      timer.update();
-      if (timer.isExpired()) {
-        final String errorMessage =
-            String.format(
-                "timeout while poll tablets with commit context: %s, consumer: 
%s",
-                commitContext, this);
-        LOGGER.warn(errorMessage);
-        throw new SubscriptionRuntimeNonCriticalException(errorMessage);
-      }
-
       if (nextOffset < 0) {
         if (!Objects.equals(tablets.size(), -nextOffset)) {
           final String errorMessage =
@@ -902,8 +929,18 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
         return Optional.of(new SubscriptionMessage(commitContext, tablets));
       }
 
+      timer.update();
+      if (timer.isExpired(TIMER_DELTA_MS)) {
+        final String errorMessage =
+            String.format(
+                "timeout while poll tablets with commit context: %s, consumer: 
%s",
+                commitContext, this);
+        LOGGER.warn(errorMessage);
+        throw new SubscriptionRuntimeNonCriticalException(errorMessage);
+      }
+
       final List<SubscriptionPollResponse> responses =
-          pollTabletsInternal(commitContext, nextOffset);
+          pollTabletsInternal(commitContext, nextOffset, timer.remainingMs());
 
       // It's agreed that the server will always return at least one response, 
even in case of
       // failure.
@@ -972,8 +1009,8 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     }
   }
 
-  private List<SubscriptionPollResponse> pollInternal(final Set<String> 
topicNames)
-      throws SubscriptionException {
+  private List<SubscriptionPollResponse> pollInternal(
+      final Set<String> topicNames, final long timeoutMs) throws 
SubscriptionException {
     providers.acquireReadLock();
     try {
       final SubscriptionProvider provider = 
providers.getNextAvailableProvider();
@@ -988,7 +1025,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
       }
       // ignore SubscriptionConnectionException to improve poll auto retry
       try {
-        return provider.poll(topicNames);
+        return provider.poll(topicNames, timeoutMs);
       } catch (final SubscriptionConnectionException ignored) {
         return Collections.emptyList();
       }
@@ -998,7 +1035,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
   }
 
   private List<SubscriptionPollResponse> pollFileInternal(
-      final SubscriptionCommitContext commitContext, final long writingOffset)
+      final SubscriptionCommitContext commitContext, final long writingOffset, 
final long timeoutMs)
       throws SubscriptionException {
     final int dataNodeId = commitContext.getDataNodeId();
     providers.acquireReadLock();
@@ -1015,7 +1052,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
       }
       // ignore SubscriptionConnectionException to improve poll auto retry
       try {
-        return provider.pollFile(commitContext, writingOffset);
+        return provider.pollFile(commitContext, writingOffset, timeoutMs);
       } catch (final SubscriptionConnectionException ignored) {
         return Collections.emptyList();
       }
@@ -1025,7 +1062,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
   }
 
   private List<SubscriptionPollResponse> pollTabletsInternal(
-      final SubscriptionCommitContext commitContext, final int offset)
+      final SubscriptionCommitContext commitContext, final int offset, final 
long timeoutMs)
       throws SubscriptionException {
     final int dataNodeId = commitContext.getDataNodeId();
     providers.acquireReadLock();
@@ -1042,7 +1079,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
       }
       // ignore SubscriptionConnectionException to improve poll auto retry
       try {
-        return provider.pollTablets(commitContext, offset);
+        return provider.pollTablets(commitContext, offset, timeoutMs);
       } catch (final SubscriptionConnectionException ignored) {
         return Collections.emptyList();
       }
@@ -1061,7 +1098,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
           .computeIfAbsent(message.getCommitContext().getDataNodeId(), (id) -> 
new ArrayList<>())
           .add(message.getCommitContext());
     }
-    for (final Map.Entry<Integer, List<SubscriptionCommitContext>> entry :
+    for (final Entry<Integer, List<SubscriptionCommitContext>> entry :
         dataNodeIdToSubscriptionCommitContexts.entrySet()) {
       commitInternal(entry.getKey(), entry.getValue(), false);
     }
@@ -1073,7 +1110,10 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     for (final SubscriptionMessage message : messages) {
       // make every effort to delete stale intermediate file
       if (Objects.equals(
-          SubscriptionMessageType.TS_FILE_HANDLER.getType(), 
message.getMessageType())) {
+              SubscriptionMessageType.TS_FILE_HANDLER.getType(), 
message.getMessageType())
+          &&
+          // do not delete file that can resume from breakpoint
+          !inFlightFilesCommitContextSet.contains(message.getCommitContext())) 
{
         try {
           message.getTsFileHandler().deleteFile();
         } catch (final Exception ignored) {
@@ -1083,7 +1123,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
           .computeIfAbsent(message.getCommitContext().getDataNodeId(), (id) -> 
new ArrayList<>())
           .add(message.getCommitContext());
     }
-    for (final Map.Entry<Integer, List<SubscriptionCommitContext>> entry :
+    for (final Entry<Integer, List<SubscriptionCommitContext>> entry :
         dataNodeIdToSubscriptionCommitContexts.entrySet()) {
       commitInternal(entry.getKey(), entry.getValue(), true);
     }
@@ -1093,7 +1133,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     final Map<Integer, List<SubscriptionCommitContext>> 
dataNodeIdToSubscriptionCommitContexts =
         new HashMap<>();
     for (final SubscriptionPollResponse response : responses) {
-      // the actual file name cannot be obtained here through the response...
+      // there is no stale intermediate file here
       dataNodeIdToSubscriptionCommitContexts
           .computeIfAbsent(response.getCommitContext().getDataNodeId(), (id) 
-> new ArrayList<>())
           .add(response.getCommitContext());
@@ -1338,7 +1378,6 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     protected boolean fileSaveFsync = 
ConsumerConstant.FILE_SAVE_FSYNC_DEFAULT_VALUE;
 
     protected int thriftMaxFrameSize = SessionConfig.DEFAULT_MAX_FRAME_SIZE;
-
     protected int maxPollParallelism = 
ConsumerConstant.MAX_POLL_PARALLELISM_DEFAULT_VALUE;
 
     public Builder host(final String host) {
@@ -1424,6 +1463,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     result.put("consumerGroupId", consumerGroupId);
     result.put("isClosed", isClosed.toString());
     result.put("fileSaveDir", fileSaveDir);
+    result.put("inFlightFilesCommitContextSet", 
inFlightFilesCommitContextSet.toString());
     result.put("subscribedTopicNames", subscribedTopics.keySet().toString());
     return result;
   }
@@ -1439,6 +1479,7 @@ abstract class SubscriptionConsumer implements 
AutoCloseable {
     result.put("isReleased", isReleased.toString());
     result.put("fileSaveDir", fileSaveDir);
     result.put("fileSaveFsync", String.valueOf(fileSaveFsync));
+    result.put("inFlightFilesCommitContextSet", 
inFlightFilesCommitContextSet.toString());
     result.put("thriftMaxFrameSize", String.valueOf(thriftMaxFrameSize));
     result.put("maxPollParallelism", String.valueOf(maxPollParallelism));
     result.put("subscribedTopics", subscribedTopics.toString());
diff --git 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionExecutorServiceManager.java
 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionExecutorServiceManager.java
index 34588654242..b8f35b392b0 100644
--- 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionExecutorServiceManager.java
+++ 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionExecutorServiceManager.java
@@ -31,7 +31,6 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.ScheduledFuture;
-import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
 
 public final class SubscriptionExecutorServiceManager {
@@ -285,10 +284,12 @@ public final class SubscriptionExecutorServiceManager {
       if (!isShutdown()) {
         synchronized (this) {
           if (!isShutdown()) {
-            return Math.max(
-                ((ThreadPoolExecutor) this.executor).getPoolSize()
-                    - ((ThreadPoolExecutor) this.executor).getActiveCount(),
-                0);
+            // TODO: temporarily disable multiple poll
+            return 0;
+            // return Math.max(
+            //    ((ThreadPoolExecutor) this.executor).getCorePoolSize()
+            //        - ((ThreadPoolExecutor) this.executor).getActiveCount(),
+            //    0);
           }
         }
       }
diff --git 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionProvider.java
 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionProvider.java
index ef03b6f20cb..dc717e6a4c6 100644
--- 
a/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionProvider.java
+++ 
b/iotdb-client/session/src/main/java/org/apache/iotdb/session/subscription/consumer/SubscriptionProvider.java
@@ -299,34 +299,35 @@ final class SubscriptionProvider extends 
SubscriptionSession {
     return unsubscribeResp.getTopics();
   }
 
-  List<SubscriptionPollResponse> poll(final Set<String> topicNames) throws 
SubscriptionException {
+  List<SubscriptionPollResponse> poll(final Set<String> topicNames, final long 
timeoutMs)
+      throws SubscriptionException {
     return poll(
         new SubscriptionPollRequest(
             SubscriptionPollRequestType.POLL.getType(),
             new PollPayload(topicNames),
-            0L,
+            timeoutMs,
             thriftMaxFrameSize));
   }
 
   List<SubscriptionPollResponse> pollFile(
-      final SubscriptionCommitContext commitContext, final long writingOffset)
+      final SubscriptionCommitContext commitContext, final long writingOffset, 
final long timeoutMs)
       throws SubscriptionException {
     return poll(
         new SubscriptionPollRequest(
             SubscriptionPollRequestType.POLL_FILE.getType(),
             new PollFilePayload(commitContext, writingOffset),
-            0L,
+            timeoutMs,
             thriftMaxFrameSize));
   }
 
   List<SubscriptionPollResponse> pollTablets(
-      final SubscriptionCommitContext commitContext, final int offset)
+      final SubscriptionCommitContext commitContext, final int offset, final 
long timeoutMs)
       throws SubscriptionException {
     return poll(
         new SubscriptionPollRequest(
             SubscriptionPollRequestType.POLL_TABLETS.getType(),
             new PollTabletsPayload(commitContext, offset),
-            0L,
+            timeoutMs,
             thriftMaxFrameSize));
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
index 7f29b3c55f9..e5383bfaae4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
@@ -47,7 +47,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.util.Objects;
-import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.atomic.AtomicReference;
 
 public class PipeProcessorSubtask extends PipeReportableSubtask {
@@ -96,7 +96,7 @@ public class PipeProcessorSubtask extends 
PipeReportableSubtask {
   @Override
   public void bindExecutors(
       final ListeningExecutorService subtaskWorkerThreadPoolExecutor,
-      final ExecutorService ignored,
+      final ScheduledExecutorService ignored,
       final PipeSubtaskScheduler subtaskScheduler) {
     this.subtaskWorkerThreadPoolExecutor = subtaskWorkerThreadPoolExecutor;
     this.subtaskScheduler = subtaskScheduler;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/batch/PipeTabletEventBatch.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/batch/PipeTabletEventBatch.java
index c7b8e1f0615..44374adb22a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/batch/PipeTabletEventBatch.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/payload/evolvable/batch/PipeTabletEventBatch.java
@@ -25,7 +25,6 @@ import 
org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALPipeExcepti
 import org.apache.iotdb.pipe.api.event.Event;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
 
-import org.apache.tsfile.exception.write.WriteProcessException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -58,7 +57,7 @@ public abstract class PipeTabletEventBatch implements 
AutoCloseable {
    * @return {@code true} if the batch can be transferred
    */
   public synchronized boolean onEvent(final TabletInsertionEvent event)
-      throws WALPipeException, IOException, WriteProcessException {
+      throws WALPipeException, IOException {
     if (isClosed || !(event instanceof EnrichedEvent)) {
       return false;
     }
@@ -94,7 +93,7 @@ public abstract class PipeTabletEventBatch implements 
AutoCloseable {
    *     exceptions and do not return {@code false} here.
    */
   protected abstract boolean constructBatch(final TabletInsertionEvent event)
-      throws WALPipeException, IOException, WriteProcessException;
+      throws WALPipeException, IOException;
 
   public boolean shouldEmit() {
     return totalBufferSize >= getMaxBatchSizeInBytes()
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 0a19cf2855c..ca8b3178560 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -396,14 +396,19 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent implements TsFileIns
   /////////////////////////// TsFileInsertionEvent ///////////////////////////
 
   @Override
-  public Iterable<TabletInsertionEvent> toTabletInsertionEvents() {
+  public Iterable<TabletInsertionEvent> toTabletInsertionEvents() throws 
PipeException {
+    return toTabletInsertionEvents(Long.MAX_VALUE);
+  }
+
+  public Iterable<TabletInsertionEvent> toTabletInsertionEvents(final long 
timeoutMs)
+      throws PipeException {
     try {
       if (!waitForTsFileClose()) {
         LOGGER.warn(
             "Pipe skipping temporary TsFile's parsing which shouldn't be 
transferred: {}", tsFile);
         return Collections.emptyList();
       }
-      waitForResourceEnough4Parsing();
+      waitForResourceEnough4Parsing(timeoutMs);
       return initDataContainer().toTabletInsertionEvents();
     } catch (final InterruptedException e) {
       Thread.currentThread().interrupt();
@@ -417,31 +422,49 @@ public class PipeTsFileInsertionEvent extends 
EnrichedEvent implements TsFileIns
     }
   }
 
-  private void waitForResourceEnough4Parsing() throws InterruptedException {
+  private void waitForResourceEnough4Parsing(final long timeoutMs) throws 
InterruptedException {
     final PipeMemoryManager memoryManager = 
PipeDataNodeResourceManager.memory();
     if (memoryManager.isEnough4TabletParsing()) {
       return;
     }
 
+    final long startTime = System.currentTimeMillis();
+    long lastRecordTime = startTime;
+
     final long memoryCheckIntervalMs =
         
PipeConfig.getInstance().getPipeTsFileParserCheckMemoryEnoughIntervalMs();
-    final long startTime = System.currentTimeMillis();
     while (!memoryManager.isEnough4TabletParsing()) {
       Thread.sleep(memoryCheckIntervalMs);
-    }
 
-    final double waitTimeSeconds = (System.currentTimeMillis() - startTime) / 
1000.0;
-    if (waitTimeSeconds > 1.0) {
-      LOGGER.info(
-          "Wait for resource enough for parsing {} for {} seconds.",
-          resource != null ? resource.getTsFilePath() : "tsfile",
-          waitTimeSeconds);
-    } else if (LOGGER.isDebugEnabled()) {
-      LOGGER.debug(
-          "Wait for resource enough for parsing {} for {} seconds.",
-          resource != null ? resource.getTsFilePath() : "tsfile",
-          waitTimeSeconds);
+      final long currentTime = System.currentTimeMillis();
+      final double elapsedRecordTimeSeconds = (currentTime - lastRecordTime) / 
1000.0;
+      final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
+      if (elapsedRecordTimeSeconds > 10.0) {
+        LOGGER.info(
+            "Wait for resource enough for parsing {} for {} seconds.",
+            resource != null ? resource.getTsFilePath() : "tsfile",
+            waitTimeSeconds);
+        lastRecordTime = currentTime;
+      } else if (LOGGER.isDebugEnabled()) {
+        LOGGER.debug(
+            "Wait for resource enough for parsing {} for {} seconds.",
+            resource != null ? resource.getTsFilePath() : "tsfile",
+            waitTimeSeconds);
+      }
+
+      if (waitTimeSeconds * 1000 > timeoutMs) {
+        // should contain 'TimeoutException' in exception message
+        throw new InterruptedException(
+            String.format("TimeoutException: Waited %s seconds", 
waitTimeSeconds));
+      }
     }
+
+    final long currentTime = System.currentTimeMillis();
+    final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
+    LOGGER.info(
+        "Wait for resource enough for parsing {} for {} seconds.",
+        resource != null ? resource.getTsFilePath() : "tsfile",
+        waitTimeSeconds);
   }
 
   /** The method is used to prevent circular replication in PipeConsensus */
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java
index f0caf0a631b..8ae6235099c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java
@@ -59,6 +59,7 @@ public class PipeMemoryManager {
       
PipeConfig.getInstance().getPipeDataStructureTabletMemoryBlockAllocationRejectThreshold();
   private volatile long usedMemorySizeInBytesOfTablets;
 
+  // Used to control the memory allocated for managing slice tsfile in 
subscription module.
   private static final double TS_FILE_MEMORY_REJECT_THRESHOLD =
       
PipeConfig.getInstance().getPipeDataStructureTsFileMemoryBlockAllocationRejectThreshold();
   private volatile long usedMemorySizeInBytesOfTsFiles;
@@ -78,6 +79,11 @@ public class PipeMemoryManager {
         < 0.95 * TABLET_MEMORY_REJECT_THRESHOLD * TOTAL_MEMORY_SIZE_IN_BYTES;
   }
 
+  public boolean isEnough4TsFileSlicing() {
+    return (double) usedMemorySizeInBytesOfTsFiles
+        < 0.95 * TS_FILE_MEMORY_REJECT_THRESHOLD * TOTAL_MEMORY_SIZE_IN_BYTES;
+  }
+
   public synchronized PipeMemoryBlock forceAllocate(long sizeInBytes)
       throws PipeRuntimeOutOfMemoryCriticalException {
     return forceAllocate(sizeInBytes, PipeMemoryBlockType.NORMAL);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
index 71f319374cd..0eb32e21f5d 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionReceiverAgent.java
@@ -20,6 +20,7 @@
 package org.apache.iotdb.db.subscription.agent;
 
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.db.subscription.receiver.SubscriptionReceiver;
 import org.apache.iotdb.db.subscription.receiver.SubscriptionReceiverV1;
 import org.apache.iotdb.rpc.RpcUtils;
@@ -69,6 +70,18 @@ public class SubscriptionReceiverAgent {
     }
   }
 
+  public long remainingMs() {
+    return remainingMs(PipeSubscribeRequestVersion.VERSION_1.getVersion()); // 
default to VERSION_1
+  }
+
+  public long remainingMs(final byte reqVersion) {
+    if (RECEIVER_CONSTRUCTORS.containsKey(reqVersion)) {
+      return getReceiver(reqVersion).remainingMs();
+    } else {
+      return 
SubscriptionConfig.getInstance().getSubscriptionDefaultTimeoutInMs();
+    }
+  }
+
   private SubscriptionReceiver getReceiver(final byte reqVersion) {
     if (receiverThreadLocal.get() == null) {
       return setAndGetReceiver(reqVersion);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
index 105ed26082e..433052409d9 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingQueue.java
@@ -23,11 +23,14 @@ import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
+import 
org.apache.iotdb.db.pipe.agent.task.execution.PipeSubtaskExecutorManager;
 import org.apache.iotdb.db.pipe.event.UserDefinedEnrichedEvent;
 import org.apache.iotdb.db.pipe.event.common.terminate.PipeTerminateEvent;
 import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import 
org.apache.iotdb.db.subscription.event.batch.SubscriptionPipeEventBatches;
+import 
org.apache.iotdb.db.subscription.task.subtask.SubscriptionReceiverSubtask;
 import org.apache.iotdb.pipe.api.event.Event;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
@@ -160,6 +163,15 @@ public abstract class SubscriptionPrefetchingQueue {
     lock.writeLock().unlock();
   }
 
+  /////////////////////////////// subtask ///////////////////////////////
+
+  protected void executeReceiverSubtask(
+      final SubscriptionReceiverSubtask subtask, final long timeoutMs) throws 
Exception {
+    PipeSubtaskExecutorManager.getInstance()
+        .getSubscriptionExecutor()
+        .executeReceiverSubtask(subtask, timeoutMs);
+  }
+
   /////////////////////////////// poll ///////////////////////////////
 
   public SubscriptionEvent poll(final String consumerId) {
@@ -176,7 +188,15 @@ public abstract class SubscriptionPrefetchingQueue {
 
     if (prefetchingQueue.isEmpty()) {
       states.markMissingPrefetch();
-      tryPrefetch();
+      try {
+        executeReceiverSubtask(
+            () -> {
+              tryPrefetch();
+              return null;
+            },
+            SubscriptionAgent.receiver().remainingMs());
+      } catch (final Exception ignored) {
+      }
     }
 
     if (prefetchingQueue.isEmpty()) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTabletQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTabletQueue.java
index 09f41cf23cb..e729fc74095 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTabletQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTabletQueue.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.subscription.broker;
 
 import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
 import 
org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionCommitContext;
@@ -145,14 +146,23 @@ public class SubscriptionPrefetchingTabletQueue extends 
SubscriptionPrefetchingQ
 
           // 3. Poll next tablets
           try {
-            ev.fetchNextResponse();
-          } catch (final Exception ignored) {
-            // no exceptions will be thrown
+            executeReceiverSubtask(
+                () -> {
+                  ev.fetchNextResponse(offset);
+                  return null;
+                },
+                SubscriptionAgent.receiver().remainingMs());
+            ev.recordLastPolledTimestamp();
+            eventRef.set(ev);
+          } catch (final Exception e) {
+            final String errorMessage =
+                String.format(
+                    "exception occurred when fetching next response: %s, 
consumer id: %s, commit context: %s, offset: %s, prefetching queue: %s",
+                    e, consumerId, commitContext, offset, this);
+            LOGGER.warn(errorMessage);
+            eventRef.set(generateSubscriptionPollErrorResponse(errorMessage));
           }
 
-          ev.recordLastPolledTimestamp();
-          eventRef.set(ev);
-
           return ev;
         });
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTsFileQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTsFileQueue.java
index 5d12eb490cc..5a7b2f7f31c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTsFileQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/SubscriptionPrefetchingTsFileQueue.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.db.subscription.broker;
 
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import 
org.apache.iotdb.db.subscription.event.pipe.SubscriptionPipeTsFilePlainEvent;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
@@ -143,16 +144,7 @@ public class SubscriptionPrefetchingTsFileQueue extends 
SubscriptionPrefetchingQ
                 
eventRef.set(generateSubscriptionPollErrorResponse(errorMessage));
                 return ev;
               }
-              // check offset
-              if (writingOffset != 0) {
-                final String errorMessage =
-                    String.format(
-                        "inconsistent offset, current: %s, incoming: %s, 
consumer: %s, file name: %s, prefetching queue: %s",
-                        0, writingOffset, consumerId, fileName, this);
-                LOGGER.warn(errorMessage);
-                
eventRef.set(generateSubscriptionPollErrorResponse(errorMessage));
-                return ev;
-              }
+              // no need to check offset for resume from breakpoint
               break;
             case FILE_PIECE:
               // check file name
@@ -206,26 +198,23 @@ public class SubscriptionPrefetchingTsFileQueue extends 
SubscriptionPrefetchingQ
 
           // 3. Poll tsfile piece or tsfile seal
           try {
-            ev.fetchNextResponse();
+            executeReceiverSubtask(
+                () -> {
+                  ev.fetchNextResponse(writingOffset);
+                  return null;
+                },
+                SubscriptionAgent.receiver().remainingMs());
+            ev.recordLastPolledTimestamp();
+            eventRef.set(ev);
           } catch (final Exception e) {
-            LOGGER.warn(
-                "Exception occurred when SubscriptionPrefetchingTsFileQueue {} 
transferring file (with event {}) to consumer {}",
-                this,
-                ev,
-                consumerId,
-                e);
             final String errorMessage =
                 String.format(
-                    "Exception occurred when 
SubscriptionPrefetchingTsFileQueue %s transferring file (with event %s) to 
consumer %s: %s",
-                    this, ev, consumerId, e);
+                    "exception occurred when fetching next response: %s, 
consumer id: %s, commit context: %s, writing offset: %s, prefetching queue: %s",
+                    e, consumerId, commitContext, writingOffset, this);
             LOGGER.warn(errorMessage);
             eventRef.set(generateSubscriptionPollErrorResponse(errorMessage));
-            return ev;
           }
 
-          ev.recordLastPolledTimestamp();
-          eventRef.set(ev);
-
           return ev;
         });
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/SubscriptionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/SubscriptionEvent.java
index bd71db5991c..ec0ed1fc079 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/SubscriptionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/SubscriptionEvent.java
@@ -220,12 +220,12 @@ public class SubscriptionEvent {
 
   //////////////////////////// prefetch & fetch ////////////////////////////
 
-  public void prefetchRemainingResponses() throws IOException {
+  public void prefetchRemainingResponses() throws Exception {
     response.prefetchRemainingResponses();
   }
 
-  public void fetchNextResponse() throws IOException {
-    response.fetchNextResponse();
+  public void fetchNextResponse(final long offset) throws Exception {
+    response.fetchNextResponse(offset);
   }
 
   //////////////////////////// byte buffer ////////////////////////////
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeEventBatch.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeEventBatch.java
index a396ff8ee2a..dbc06881cae 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeEventBatch.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeEventBatch.java
@@ -94,7 +94,7 @@ public abstract class SubscriptionPipeEventBatch {
       final @NonNull EnrichedEvent event, final Consumer<SubscriptionEvent> 
consumer)
       throws Exception {
     if (event instanceof TabletInsertionEvent) {
-      onTabletInsertionEvent((TabletInsertionEvent) event); // no exceptions 
will be thrown
+      onTabletInsertionEvent((TabletInsertionEvent) event);
       enrichedEvents.add(event);
     } else if (event instanceof TsFileInsertionEvent) {
       onTsFileInsertionEvent((TsFileInsertionEvent) event);
@@ -108,9 +108,9 @@ public abstract class SubscriptionPipeEventBatch {
 
   /////////////////////////////// utility ///////////////////////////////
 
-  protected abstract void onTabletInsertionEvent(final TabletInsertionEvent 
event) throws Exception;
+  protected abstract void onTabletInsertionEvent(final TabletInsertionEvent 
event);
 
-  protected abstract void onTsFileInsertionEvent(final TsFileInsertionEvent 
event) throws Exception;
+  protected abstract void onTsFileInsertionEvent(final TsFileInsertionEvent 
event);
 
   protected abstract boolean shouldEmit();
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTabletEventBatch.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTabletEventBatch.java
index b54607c868b..2100d3839de 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTabletEventBatch.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTabletEventBatch.java
@@ -22,7 +22,9 @@ package org.apache.iotdb.db.subscription.event.batch;
 import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
 import 
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
 import 
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
 import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import 
org.apache.iotdb.db.subscription.broker.SubscriptionPrefetchingTabletQueue;
 import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
 import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
@@ -112,7 +114,9 @@ public class SubscriptionPipeTabletEventBatch extends 
SubscriptionPipeEventBatch
 
   @Override
   protected void onTsFileInsertionEvent(final TsFileInsertionEvent event) {
-    for (final TabletInsertionEvent tabletInsertionEvent : 
event.toTabletInsertionEvents()) {
+    for (final TabletInsertionEvent tabletInsertionEvent :
+        ((PipeTsFileInsertionEvent) event)
+            
.toTabletInsertionEvents(SubscriptionAgent.receiver().remainingMs())) {
       onTabletInsertionEvent(tabletInsertionEvent);
     }
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTsFileEventBatch.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTsFileEventBatch.java
index dc96fc476da..514795db23a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTsFileEventBatch.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/batch/SubscriptionPipeTsFileEventBatch.java
@@ -68,8 +68,12 @@ public class SubscriptionPipeTsFileEventBatch extends 
SubscriptionPipeEventBatch
   /////////////////////////////// utility ///////////////////////////////
 
   @Override
-  protected void onTabletInsertionEvent(final TabletInsertionEvent event) 
throws Exception {
-    batch.onEvent(event); // no exceptions will be thrown
+  protected void onTabletInsertionEvent(final TabletInsertionEvent event) {
+    try {
+      batch.onEvent(event);
+    } catch (final Exception ignored) {
+      // no exceptions will be thrown
+    }
     ((EnrichedEvent) event)
         .decreaseReferenceCount(
             SubscriptionPipeTsFileEventBatch.class.getName(),
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventExtendableResponse.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventExtendableResponse.java
index 07dfaabdf4f..dbbe9d898c8 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventExtendableResponse.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventExtendableResponse.java
@@ -57,16 +57,6 @@ public abstract class SubscriptionEventExtendableResponse
     return peekFirst();
   }
 
-  @Override
-  public void fetchNextResponse() throws IOException {
-    prefetchRemainingResponses();
-    if (Objects.isNull(poll())) {
-      LOGGER.warn(
-          "SubscriptionEventExtendableResponse {} is empty when fetching next 
response (broken invariant)",
-          this);
-    }
-  }
-
   @Override
   public void trySerializeCurrentResponse() {
     
SubscriptionPollResponseCache.getInstance().trySerialize(getCurrentResponse());
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventResponse.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventResponse.java
index 211f0905afa..3ded870feed 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventResponse.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventResponse.java
@@ -28,9 +28,9 @@ public interface SubscriptionEventResponse<E> {
 
   E getCurrentResponse();
 
-  void prefetchRemainingResponses() throws IOException;
+  void prefetchRemainingResponses() throws Exception;
 
-  void fetchNextResponse() throws IOException;
+  void fetchNextResponse(final long offset) throws Exception;
 
   /////////////////////////////// byte buffer ///////////////////////////////
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventSingleResponse.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventSingleResponse.java
index dbc48ebda00..7c72e65b6d6 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventSingleResponse.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventSingleResponse.java
@@ -65,7 +65,7 @@ public class SubscriptionEventSingleResponse
   }
 
   @Override
-  public void fetchNextResponse() {
+  public void fetchNextResponse(final long offset) {
     // do nothing
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTabletResponse.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTabletResponse.java
index 5ed3f7fd56d..e2590f5ff32 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTabletResponse.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTabletResponse.java
@@ -77,6 +77,16 @@ public class SubscriptionEventTabletResponse extends 
SubscriptionEventExtendable
     offer(generateNextTabletResponse());
   }
 
+  @Override
+  public void fetchNextResponse(final long offset /* unused */) {
+    prefetchRemainingResponses();
+    if (Objects.isNull(poll())) {
+      LOGGER.warn(
+          "SubscriptionEventTabletResponse {} is empty when fetching next 
response (broken invariant)",
+          this);
+    }
+  }
+
   @Override
   public synchronized void nack() {
     if (nextOffset.get() == 1) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTsFileResponse.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTsFileResponse.java
index 63a4c2d1f93..6c102c21912 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTsFileResponse.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/event/response/SubscriptionEventTsFileResponse.java
@@ -19,9 +19,12 @@
 
 package org.apache.iotdb.db.subscription.event.response;
 
+import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
+import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager;
 import org.apache.iotdb.db.pipe.resource.memory.PipeTsFileMemoryBlock;
+import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import 
org.apache.iotdb.db.subscription.event.cache.CachedSubscriptionPollResponse;
 import org.apache.iotdb.rpc.subscription.exception.SubscriptionException;
 import org.apache.iotdb.rpc.subscription.payload.poll.FileInitPayload;
@@ -67,12 +70,18 @@ public class SubscriptionEventTsFileResponse extends 
SubscriptionEventExtendable
   }
 
   @Override
-  public void prefetchRemainingResponses() throws IOException {
-    if (hasNoMore) {
-      return;
-    }
+  public void prefetchRemainingResponses() {
+    // do nothing
+  }
 
-    generateNextTsFileResponse().ifPresent(super::offer);
+  @Override
+  public void fetchNextResponse(final long offset) throws Exception {
+    generateNextTsFileResponse(offset).ifPresent(super::offer);
+    if (Objects.isNull(poll())) {
+      LOGGER.warn(
+          "SubscriptionEventTsFileResponse {} is empty when fetching next 
response (broken invariant)",
+          this);
+    }
   }
 
   @Override
@@ -103,8 +112,13 @@ public class SubscriptionEventTsFileResponse extends 
SubscriptionEventExtendable
             commitContext));
   }
 
+  private synchronized Optional<CachedSubscriptionPollResponse> 
generateNextTsFileResponse(
+      final long offset) throws SubscriptionException, IOException, 
InterruptedException {
+    return Optional.of(generateResponseWithPieceOrSealPayload(offset));
+  }
+
   private synchronized Optional<CachedSubscriptionPollResponse> 
generateNextTsFileResponse()
-      throws IOException {
+      throws SubscriptionException, IOException, InterruptedException {
     final SubscriptionPollResponse previousResponse = peekLast();
     if (Objects.isNull(previousResponse)) {
       LOGGER.warn(
@@ -137,7 +151,7 @@ public class SubscriptionEventTsFileResponse extends 
SubscriptionEventExtendable
   }
 
   private @NonNull CachedSubscriptionPollResponse 
generateResponseWithPieceOrSealPayload(
-      final long writingOffset) throws IOException {
+      final long writingOffset) throws SubscriptionException, IOException, 
InterruptedException {
     final long tsFileLength = tsFile.length();
     if (writingOffset >= tsFileLength) {
       // generate subscription poll response with seal payload
@@ -159,6 +173,7 @@ public class SubscriptionEventTsFileResponse extends 
SubscriptionEventExtendable
       bufferSize = readFileBufferSize;
     }
 
+    waitForResourceEnough4Slicing(SubscriptionAgent.receiver().remainingMs());
     try (final RandomAccessFile reader = new RandomAccessFile(tsFile, "r")) {
       reader.seek(writingOffset);
 
@@ -185,4 +200,48 @@ public class SubscriptionEventTsFileResponse extends 
SubscriptionEventExtendable
       return response;
     }
   }
+
+  private void waitForResourceEnough4Slicing(final long timeoutMs) throws 
InterruptedException {
+    final PipeMemoryManager memoryManager = 
PipeDataNodeResourceManager.memory();
+    if (memoryManager.isEnough4TsFileSlicing()) {
+      return;
+    }
+
+    final long startTime = System.currentTimeMillis();
+    long lastRecordTime = startTime;
+
+    final long memoryCheckIntervalMs =
+        
PipeConfig.getInstance().getPipeTsFileParserCheckMemoryEnoughIntervalMs();
+    while (!memoryManager.isEnough4TsFileSlicing()) {
+      Thread.sleep(memoryCheckIntervalMs);
+
+      final long currentTime = System.currentTimeMillis();
+      final double elapsedRecordTimeSeconds = (currentTime - lastRecordTime) / 
1000.0;
+      final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
+      if (elapsedRecordTimeSeconds > 10.0) {
+        LOGGER.info(
+            "Wait for resource enough for slicing tsfile {} for {} seconds.",
+            tsFile,
+            waitTimeSeconds);
+        lastRecordTime = currentTime;
+      } else if (LOGGER.isDebugEnabled()) {
+        LOGGER.debug(
+            "Wait for resource enough for slicing tsfile {} for {} seconds.",
+            tsFile,
+            waitTimeSeconds);
+      }
+
+      if (waitTimeSeconds * 1000 > timeoutMs) {
+        // should contain 'TimeoutException' in exception message
+        // see 
org.apache.iotdb.rpc.subscription.exception.SubscriptionTimeoutException.KEYWORD
+        throw new InterruptedException(
+            String.format("TimeoutException: Waited %s seconds", 
waitTimeSeconds));
+      }
+    }
+
+    final long currentTime = System.currentTimeMillis();
+    final double waitTimeSeconds = (currentTime - startTime) / 1000.0;
+    LOGGER.info(
+        "Wait for resource enough for slicing tsfile {} for {} seconds.", 
tsFile, waitTimeSeconds);
+  }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
index 5d2d4b050b1..c636e1b2dc6 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
@@ -30,4 +30,6 @@ public interface SubscriptionReceiver {
   PipeSubscribeRequestVersion getVersion();
 
   void handleExit();
+
+  long remainingMs();
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
index 0a03f28bf7d..73ee10c2854 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiverV1.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.client.IClientManager;
 import org.apache.iotdb.commons.client.exception.ClientManagerException;
 import org.apache.iotdb.commons.consensus.ConfigRegionId;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import org.apache.iotdb.confignode.rpc.thrift.TCloseConsumerReq;
 import org.apache.iotdb.confignode.rpc.thrift.TCreateConsumerReq;
 import org.apache.iotdb.confignode.rpc.thrift.TSubscribeReq;
@@ -68,6 +69,7 @@ import 
org.apache.iotdb.rpc.subscription.payload.response.PipeSubscribeSubscribe
 import 
org.apache.iotdb.rpc.subscription.payload.response.PipeSubscribeUnsubscribeResp;
 import org.apache.iotdb.service.rpc.thrift.TPipeSubscribeReq;
 import org.apache.iotdb.service.rpc.thrift.TPipeSubscribeResp;
+import org.apache.iotdb.session.subscription.util.PollTimer;
 
 import org.apache.thrift.TException;
 import org.slf4j.Logger;
@@ -101,23 +103,7 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
           PipeSubscribeResponseType.ACK.getType());
 
   private final ThreadLocal<ConsumerConfig> consumerConfigThreadLocal = new 
ThreadLocal<>();
-
-  @Override
-  public PipeSubscribeRequestVersion getVersion() {
-    return PipeSubscribeRequestVersion.VERSION_1;
-  }
-
-  @Override
-  public void handleExit() {
-    final ConsumerConfig consumerConfig = consumerConfigThreadLocal.get();
-    if (Objects.nonNull(consumerConfig)) {
-      LOGGER.info(
-          "Subscription: remove consumer config {} when handling exit",
-          consumerConfigThreadLocal.get());
-      // closeConsumer(consumerConfig);
-      consumerConfigThreadLocal.remove();
-    }
-  }
+  private final ThreadLocal<PollTimer> pollTimerThreadLocal = new 
ThreadLocal<>();
 
   @Override
   public final TPipeSubscribeResp handle(final TPipeSubscribeReq req) {
@@ -155,6 +141,32 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
         PipeSubscribeResponseType.ACK.getType());
   }
 
+  @Override
+  public PipeSubscribeRequestVersion getVersion() {
+    return PipeSubscribeRequestVersion.VERSION_1;
+  }
+
+  @Override
+  public void handleExit() {
+    final ConsumerConfig consumerConfig = consumerConfigThreadLocal.get();
+    if (Objects.nonNull(consumerConfig)) {
+      LOGGER.info(
+          "Subscription: remove consumer config {} when handling exit",
+          consumerConfigThreadLocal.get());
+      // closeConsumer(consumerConfig);
+      consumerConfigThreadLocal.remove();
+    }
+  }
+
+  @Override
+  public long remainingMs() {
+    final PollTimer pollTimer = pollTimerThreadLocal.get();
+    if (Objects.isNull(pollTimer)) {
+      return 
SubscriptionConfig.getInstance().getSubscriptionDefaultTimeoutInMs();
+    }
+    return pollTimer.remainingMs();
+  }
+
   private TPipeSubscribeResp handlePipeSubscribeHandshake(final 
PipeSubscribeHandshakeReq req) {
     try {
       return handlePipeSubscribeHandshakeInternal(req);
@@ -342,6 +354,24 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
   }
 
   private TPipeSubscribeResp handlePipeSubscribePoll(final 
PipeSubscribePollReq req) {
+    try {
+      return handlePipeSubscribePollInternal(req);
+    } catch (final Exception e) {
+      LOGGER.warn("Exception occurred when polling with request {}", req, e);
+      final String exceptionMessage =
+          String.format(
+              "Subscription: something unexpected happened when polling with 
request %s: %s",
+              req, e);
+      return PipeSubscribePollResp.toTPipeSubscribeResp(
+          RpcUtils.getStatus(TSStatusCode.SUBSCRIPTION_POLL_ERROR, 
exceptionMessage),
+          Collections.emptyList());
+    } finally {
+      pollTimerThreadLocal.remove();
+    }
+  }
+
+  private TPipeSubscribeResp handlePipeSubscribePollInternal(final 
PipeSubscribePollReq req)
+      throws SubscriptionException {
     final ConsumerConfig consumerConfig = consumerConfigThreadLocal.get();
     if (Objects.isNull(consumerConfig)) {
       LOGGER.warn(
@@ -351,120 +381,113 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
 
     final List<SubscriptionEvent> events;
     final SubscriptionPollRequest request = req.getRequest();
+
+    pollTimerThreadLocal.set(new PollTimer(System.currentTimeMillis(), 
request.getTimeoutMs()));
+
     final long maxBytes = (long) (request.getMaxBytes() * 
POLL_PAYLOAD_SIZE_EXCEED_THRESHOLD);
-    try {
-      final short requestType = request.getRequestType();
-      if (SubscriptionPollRequestType.isValidatedRequestType(requestType)) {
-        switch (SubscriptionPollRequestType.valueOf(requestType)) {
-          case POLL:
-            events =
-                handlePipeSubscribePollInternal(
-                    consumerConfig, (PollPayload) request.getPayload(), 
maxBytes);
-            break;
-          case POLL_FILE:
-            events =
-                handlePipeSubscribePollTsFileInternal(
-                    consumerConfig, (PollFilePayload) request.getPayload());
-            break;
-          case POLL_TABLETS:
-            events =
-                handlePipeSubscribePollTabletsInternal(
-                    consumerConfig, (PollTabletsPayload) request.getPayload());
-            break;
-          default:
-            events = null;
-            break;
-        }
-      } else {
-        events = null;
-      }
-      if (Objects.isNull(events)) {
-        throw new SubscriptionException(String.format("unexpected request 
type: %s", requestType));
+    final short requestType = request.getRequestType();
+    if (SubscriptionPollRequestType.isValidatedRequestType(requestType)) {
+      switch (SubscriptionPollRequestType.valueOf(requestType)) {
+        case POLL:
+          events =
+              handlePipeSubscribePollRequest(
+                  consumerConfig, (PollPayload) request.getPayload(), 
maxBytes);
+          break;
+        case POLL_FILE:
+          events =
+              handlePipeSubscribePollTsFileRequest(
+                  consumerConfig, (PollFilePayload) request.getPayload());
+          break;
+        case POLL_TABLETS:
+          events =
+              handlePipeSubscribePollTabletsRequest(
+                  consumerConfig, (PollTabletsPayload) request.getPayload());
+          break;
+        default:
+          events = null;
+          break;
       }
+    } else {
+      events = null;
+    }
 
-      // generate response
-      final AtomicLong totalSize = new AtomicLong();
-      return PipeSubscribePollResp.toTPipeSubscribeResp(
-          RpcUtils.SUCCESS_STATUS,
-          events.stream()
-              .map(
-                  (event) -> {
-                    final SubscriptionCommitContext commitContext = 
event.getCommitContext();
-                    final SubscriptionPollResponse response = 
event.getCurrentResponse();
-                    if (Objects.isNull(response)) {
-                      LOGGER.warn(
-                          "Subscription: consumer {} poll null response for 
event {} with request: {}",
-                          consumerConfig,
-                          event,
-                          req.getRequest());
-                      // nack
-                      SubscriptionAgent.broker()
-                          .commit(consumerConfig, 
Collections.singletonList(commitContext), true);
-                      return null;
-                    }
+    if (Objects.isNull(events)) {
+      throw new SubscriptionException(String.format("unexpected request type: 
%s", requestType));
+    }
 
-                    try {
-                      final ByteBuffer byteBuffer = 
event.getCurrentResponseByteBuffer();
-
-                      // payload size control
-                      final long size = event.getCurrentResponseSize();
-                      if (totalSize.get() + size > maxBytes) {
-                        throw new SubscriptionPayloadExceedException(
-                            String.format(
-                                "payload size %s byte(s) will exceed the 
threshold %s byte(s)",
-                                totalSize.get() + size, maxBytes));
-                      }
-                      totalSize.getAndAdd(size);
-
-                      SubscriptionPrefetchingQueueMetrics.getInstance()
-                          .mark(
-                              
SubscriptionPrefetchingQueue.generatePrefetchingQueueId(
-                                  commitContext.getConsumerGroupId(), 
commitContext.getTopicName()),
-                              size);
-                      event.invalidateCurrentResponseByteBuffer();
-                      LOGGER.info(
-                          "Subscription: consumer {} poll {} successfully with 
request: {}",
+    // generate response
+    final AtomicLong totalSize = new AtomicLong();
+    return PipeSubscribePollResp.toTPipeSubscribeResp(
+        RpcUtils.SUCCESS_STATUS,
+        events.stream()
+            .map(
+                (event) -> {
+                  final SubscriptionCommitContext commitContext = 
event.getCommitContext();
+                  final SubscriptionPollResponse response = 
event.getCurrentResponse();
+                  if (Objects.isNull(response)) {
+                    LOGGER.warn(
+                        "Subscription: consumer {} poll null response for 
event {} with request: {}",
+                        consumerConfig,
+                        event,
+                        req.getRequest());
+                    // nack
+                    SubscriptionAgent.broker()
+                        .commit(consumerConfig, 
Collections.singletonList(commitContext), true);
+                    return null;
+                  }
+
+                  try {
+                    final ByteBuffer byteBuffer = 
event.getCurrentResponseByteBuffer();
+
+                    // payload size control
+                    final long size = event.getCurrentResponseSize();
+                    if (totalSize.get() + size > maxBytes) {
+                      throw new SubscriptionPayloadExceedException(
+                          String.format(
+                              "payload size %s byte(s) will exceed the 
threshold %s byte(s)",
+                              totalSize.get() + size, maxBytes));
+                    }
+                    totalSize.getAndAdd(size);
+
+                    SubscriptionPrefetchingQueueMetrics.getInstance()
+                        .mark(
+                            
SubscriptionPrefetchingQueue.generatePrefetchingQueueId(
+                                commitContext.getConsumerGroupId(), 
commitContext.getTopicName()),
+                            size);
+                    event.invalidateCurrentResponseByteBuffer();
+                    LOGGER.info(
+                        "Subscription: consumer {} poll {} successfully with 
request: {}",
+                        consumerConfig,
+                        response,
+                        req.getRequest());
+                    return byteBuffer;
+                  } catch (final Exception e) {
+                    if (e instanceof SubscriptionPayloadExceedException) {
+                      LOGGER.error(
+                          "Subscription: consumer {} poll excessive payload {} 
with request: {}, something unexpected happened with parameter configuration or 
payload control...",
+                          consumerConfig,
+                          response,
+                          req.getRequest(),
+                          e);
+                    } else {
+                      LOGGER.warn(
+                          "Subscription: consumer {} poll {} failed with 
request: {}",
                           consumerConfig,
                           response,
-                          req.getRequest());
-                      return byteBuffer;
-                    } catch (final Exception e) {
-                      if (e instanceof SubscriptionPayloadExceedException) {
-                        LOGGER.error(
-                            "Subscription: consumer {} poll excessive payload 
{} with request: {}, something unexpected happened with parameter configuration 
or payload control...",
-                            consumerConfig,
-                            response,
-                            req.getRequest(),
-                            e);
-                      } else {
-                        LOGGER.warn(
-                            "Subscription: consumer {} poll {} failed with 
request: {}",
-                            consumerConfig,
-                            response,
-                            req.getRequest(),
-                            e);
-                      }
-                      // nack
-                      SubscriptionAgent.broker()
-                          .commit(consumerConfig, 
Collections.singletonList(commitContext), true);
-                      return null;
+                          req.getRequest(),
+                          e);
                     }
-                  })
-              .filter(Objects::nonNull)
-              .collect(Collectors.toList()));
-    } catch (final Exception e) {
-      LOGGER.warn("Exception occurred when polling with request {}", req, e);
-      final String exceptionMessage =
-          String.format(
-              "Subscription: something unexpected happened when polling with 
request %s: %s",
-              req, e);
-      return PipeSubscribePollResp.toTPipeSubscribeResp(
-          RpcUtils.getStatus(TSStatusCode.SUBSCRIPTION_POLL_ERROR, 
exceptionMessage),
-          Collections.emptyList());
-    }
+                    // nack
+                    SubscriptionAgent.broker()
+                        .commit(consumerConfig, 
Collections.singletonList(commitContext), true);
+                    return null;
+                  }
+                })
+            .filter(Objects::nonNull)
+            .collect(Collectors.toList()));
   }
 
-  private List<SubscriptionEvent> handlePipeSubscribePollInternal(
+  private List<SubscriptionEvent> handlePipeSubscribePollRequest(
       final ConsumerConfig consumerConfig, final PollPayload messagePayload, 
final long maxBytes) {
     final Set<String> subscribedTopicNames =
         SubscriptionAgent.consumer()
@@ -480,14 +503,14 @@ public class SubscriptionReceiverV1 implements 
SubscriptionReceiver {
     return SubscriptionAgent.broker().poll(consumerConfig, topicNames, 
maxBytes);
   }
 
-  private List<SubscriptionEvent> handlePipeSubscribePollTsFileInternal(
+  private List<SubscriptionEvent> handlePipeSubscribePollTsFileRequest(
       final ConsumerConfig consumerConfig, final PollFilePayload 
messagePayload) {
     return SubscriptionAgent.broker()
         .pollTsFile(
             consumerConfig, messagePayload.getCommitContext(), 
messagePayload.getWritingOffset());
   }
 
-  private List<SubscriptionEvent> handlePipeSubscribePollTabletsInternal(
+  private List<SubscriptionEvent> handlePipeSubscribePollTabletsRequest(
       final ConsumerConfig consumerConfig, final PollTabletsPayload 
messagePayload) {
     return SubscriptionAgent.broker()
         .pollTablets(consumerConfig, messagePayload.getCommitContext(), 
messagePayload.getOffset());
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
index 76b0a7db2f6..f473a56a447 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
@@ -20,14 +20,64 @@
 package org.apache.iotdb.db.subscription.task.execution;
 
 import org.apache.iotdb.commons.concurrent.ThreadName;
+import org.apache.iotdb.commons.pipe.agent.task.execution.PipeSubtaskExecutor;
+import org.apache.iotdb.commons.pipe.agent.task.execution.PipeSubtaskScheduler;
 import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import 
org.apache.iotdb.db.pipe.agent.task.execution.PipeConnectorSubtaskExecutor;
+import 
org.apache.iotdb.db.subscription.task.subtask.SubscriptionReceiverSubtask;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicLong;
 
 public class SubscriptionSubtaskExecutor extends PipeConnectorSubtaskExecutor {
 
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(SubscriptionSubtaskExecutor.class);
+
+  private final AtomicLong submittedReceiverSubtasks = new AtomicLong(0);
+
   public SubscriptionSubtaskExecutor() {
     super(
         
SubscriptionConfig.getInstance().getSubscriptionSubtaskExecutorMaxThreadNum(),
         ThreadName.SUBSCRIPTION_EXECUTOR_POOL);
   }
+
+  @Override
+  protected PipeSubtaskScheduler schedulerSupplier(final PipeSubtaskExecutor 
executor) {
+    return new SubscriptionSubtaskScheduler((SubscriptionSubtaskExecutor) 
executor);
+  }
+
+  public void executeReceiverSubtask(
+      final SubscriptionReceiverSubtask subtask, final long timeoutMs) throws 
Exception {
+    if (!super.hasAvailableThread()) {
+      subtask.call(); // non-strict timeout
+      return;
+    }
+
+    submittedReceiverSubtasks.incrementAndGet();
+    try {
+      final Future<Void> future = 
subtaskWorkerThreadPoolExecutor.submit(subtask);
+      try {
+        future.get(timeoutMs, TimeUnit.MILLISECONDS); // strict timeout
+      } catch (final InterruptedException e) {
+        Thread.currentThread().interrupt(); // restore interrupted state
+        future.cancel(true);
+        throw e;
+      } catch (final ExecutionException | TimeoutException e) {
+        future.cancel(true);
+        throw e;
+      }
+    } finally {
+      submittedReceiverSubtasks.decrementAndGet();
+    }
+  }
+
+  public boolean hasSubmittedReceiverSubtasks() {
+    return submittedReceiverSubtasks.get() > 0;
+  }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskScheduler.java
similarity index 61%
copy from 
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
copy to 
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskScheduler.java
index 76b0a7db2f6..44072d1d823 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskExecutor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/execution/SubscriptionSubtaskScheduler.java
@@ -19,15 +19,21 @@
 
 package org.apache.iotdb.db.subscription.task.execution;
 
-import org.apache.iotdb.commons.concurrent.ThreadName;
-import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
-import 
org.apache.iotdb.db.pipe.agent.task.execution.PipeConnectorSubtaskExecutor;
+import org.apache.iotdb.commons.pipe.agent.task.execution.PipeSubtaskScheduler;
 
-public class SubscriptionSubtaskExecutor extends PipeConnectorSubtaskExecutor {
+public class SubscriptionSubtaskScheduler extends PipeSubtaskScheduler {
 
-  public SubscriptionSubtaskExecutor() {
-    super(
-        
SubscriptionConfig.getInstance().getSubscriptionSubtaskExecutorMaxThreadNum(),
-        ThreadName.SUBSCRIPTION_EXECUTOR_POOL);
+  private final SubscriptionSubtaskExecutor executor;
+
+  public SubscriptionSubtaskScheduler(final SubscriptionSubtaskExecutor 
executor) {
+    super(executor);
+
+    this.executor = executor;
+  }
+
+  @Override
+  public boolean schedule() {
+    // prioritize executing SubscriptionReceiverSubtask over 
SubscriptionConnectorSubtask
+    return !executor.hasSubmittedReceiverSubtasks() && super.schedule();
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionConnectorSubtask.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionConnectorSubtask.java
index 1b4c0425159..c5b841a8c66 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionConnectorSubtask.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionConnectorSubtask.java
@@ -20,14 +20,19 @@
 package org.apache.iotdb.db.subscription.task.subtask;
 
 import 
org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
+import org.apache.iotdb.commons.subscription.config.SubscriptionConfig;
 import 
org.apache.iotdb.db.pipe.agent.task.subtask.connector.PipeConnectorSubtask;
 import org.apache.iotdb.db.subscription.agent.SubscriptionAgent;
 import org.apache.iotdb.pipe.api.PipeConnector;
 import org.apache.iotdb.pipe.api.event.Event;
 
+import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.time.Duration;
+
 public class SubscriptionConnectorSubtask extends PipeConnectorSubtask {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(SubscriptionConnectorSubtask.class);
@@ -55,15 +60,6 @@ public class SubscriptionConnectorSubtask extends 
PipeConnectorSubtask {
     this.consumerGroupId = consumerGroupId;
   }
 
-  @Override
-  protected boolean executeOnce() {
-    if (isClosed.get()) {
-      return false;
-    }
-
-    return SubscriptionAgent.broker().executePrefetch(consumerGroupId, 
topicName);
-  }
-
   public String getTopicName() {
     return topicName;
   }
@@ -76,6 +72,36 @@ public class SubscriptionConnectorSubtask extends 
PipeConnectorSubtask {
     return inputPendingQueue;
   }
 
+  //////////////////////////// execution & callback 
////////////////////////////
+
+  @Override
+  protected void registerCallbackHookAfterSubmit(final 
ListenableFuture<Boolean> future) {
+    final ListenableFuture<Boolean> nextFuture =
+        Futures.withTimeout(
+            future,
+            Duration.ofSeconds(
+                
SubscriptionConfig.getInstance().getSubscriptionDefaultTimeoutInMs()),
+            subtaskCallbackListeningExecutor);
+    Futures.addCallback(nextFuture, this, subtaskCallbackListeningExecutor);
+  }
+
+  @Override
+  public synchronized void onFailure(final Throwable throwable) {
+    isSubmitted = false;
+
+    // just resubmit
+    submitSelf();
+  }
+
+  @Override
+  protected boolean executeOnce() {
+    if (isClosed.get()) {
+      return false;
+    }
+
+    return SubscriptionAgent.broker().executePrefetch(consumerGroupId, 
topicName);
+  }
+
   //////////////////////////// APIs provided for metric framework 
////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionReceiverSubtask.java
similarity index 65%
copy from 
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
copy to 
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionReceiverSubtask.java
index 5d2d4b050b1..a070e6dbe6b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/receiver/SubscriptionReceiver.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionReceiverSubtask.java
@@ -17,17 +17,8 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.subscription.receiver;
+package org.apache.iotdb.db.subscription.task.subtask;
 
-import 
org.apache.iotdb.rpc.subscription.payload.request.PipeSubscribeRequestVersion;
-import org.apache.iotdb.service.rpc.thrift.TPipeSubscribeReq;
-import org.apache.iotdb.service.rpc.thrift.TPipeSubscribeResp;
+import java.util.concurrent.Callable;
 
-public interface SubscriptionReceiver {
-
-  TPipeSubscribeResp handle(TPipeSubscribeReq req);
-
-  PipeSubscribeRequestVersion getVersion();
-
-  void handleExit();
-}
+public interface SubscriptionReceiverSubtask extends Callable<Void> {}
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
index c8b8129546f..cf9fd879f2b 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
@@ -285,12 +285,13 @@ public class CommonConfig {
   private int subscriptionPrefetchTsFileBatchMaxDelayInMs = 5000; // 5s
   private long subscriptionPrefetchTsFileBatchMaxSizeInBytes = 80 * MB;
   private int subscriptionPollMaxBlockingTimeMs = 500;
-  private int subscriptionSerializeMaxBlockingTimeMs = 100;
+  private int subscriptionDefaultTimeoutInMs = 10_000; // 10s
   private long subscriptionLaunchRetryIntervalMs = 1000;
   private int subscriptionRecycleUncommittedEventIntervalMs = 600000; // 600s
   private long subscriptionReadFileBufferSize = 8 * MB;
   private long subscriptionReadTabletBufferSize = 8 * MB;
   private long subscriptionTsFileDeduplicationWindowSeconds = 120; // 120s
+  private volatile long subscriptionTsFileSlicerCheckMemoryEnoughIntervalMs = 
10L;
 
   private long subscriptionMetaSyncerInitialSyncDelayMinutes = 3;
   private long subscriptionMetaSyncerSyncIntervalMinutes = 3;
@@ -1283,13 +1284,12 @@ public class CommonConfig {
     this.subscriptionPollMaxBlockingTimeMs = subscriptionPollMaxBlockingTimeMs;
   }
 
-  public int getSubscriptionSerializeMaxBlockingTimeMs() {
-    return subscriptionSerializeMaxBlockingTimeMs;
+  public int getSubscriptionDefaultTimeoutInMs() {
+    return subscriptionDefaultTimeoutInMs;
   }
 
-  public void setSubscriptionSerializeMaxBlockingTimeMs(
-      int subscriptionSerializeMaxBlockingTimeMs) {
-    this.subscriptionSerializeMaxBlockingTimeMs = 
subscriptionSerializeMaxBlockingTimeMs;
+  public void setSubscriptionDefaultTimeoutInMs(final int 
subscriptionDefaultTimeoutInMs) {
+    this.subscriptionDefaultTimeoutInMs = subscriptionDefaultTimeoutInMs;
   }
 
   public long getSubscriptionLaunchRetryIntervalMs() {
@@ -1336,6 +1336,16 @@ public class CommonConfig {
         subscriptionTsFileDeduplicationWindowSeconds;
   }
 
+  public long getSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs() {
+    return subscriptionTsFileSlicerCheckMemoryEnoughIntervalMs;
+  }
+
+  public void setSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs(
+      long subscriptionTsFileSlicerCheckMemoryEnoughIntervalMs) {
+    this.subscriptionTsFileSlicerCheckMemoryEnoughIntervalMs =
+        subscriptionTsFileSlicerCheckMemoryEnoughIntervalMs;
+  }
+
   public long getSubscriptionMetaSyncerInitialSyncDelayMinutes() {
     return subscriptionMetaSyncerInitialSyncDelayMinutes;
   }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
index 4d2799d178b..dff60d3e100 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonDescriptor.java
@@ -662,11 +662,11 @@ public class CommonDescriptor {
             properties.getProperty(
                 "subscription_poll_max_blocking_time_ms",
                 
String.valueOf(config.getSubscriptionPollMaxBlockingTimeMs()))));
-    config.setSubscriptionSerializeMaxBlockingTimeMs(
+    config.setSubscriptionDefaultTimeoutInMs(
         Integer.parseInt(
             properties.getProperty(
-                "subscription_serialize_max_blocking_time_ms",
-                
String.valueOf(config.getSubscriptionSerializeMaxBlockingTimeMs()))));
+                "subscription_default_timeout_in_ms",
+                String.valueOf(config.getSubscriptionDefaultTimeoutInMs()))));
     config.setSubscriptionLaunchRetryIntervalMs(
         Long.parseLong(
             properties.getProperty(
@@ -692,6 +692,11 @@ public class CommonDescriptor {
             properties.getProperty(
                 "subscription_ts_file_deduplication_window_seconds",
                 
String.valueOf(config.getSubscriptionTsFileDeduplicationWindowSeconds()))));
+    config.setSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs(
+        Long.parseLong(
+            properties.getProperty(
+                "subscription_ts_file_slicer_check_memory_enough_interval_ms",
+                
String.valueOf(config.getSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs()))));
 
     config.setSubscriptionMetaSyncerInitialSyncDelayMinutes(
         Long.parseLong(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java
index 6130292cdf7..c4c96dad3e7 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/execution/PipeSubtaskExecutor.java
@@ -32,15 +32,17 @@ import org.slf4j.LoggerFactory;
 
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 
 public abstract class PipeSubtaskExecutor {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(PipeSubtaskExecutor.class);
 
-  private static final ExecutorService subtaskCallbackListeningExecutor =
-      IoTDBThreadPoolFactory.newSingleThreadExecutor(
+  private static final ScheduledExecutorService 
subtaskCallbackListeningExecutor =
+      IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
           ThreadName.PIPE_SUBTASK_CALLBACK_EXECUTOR_POOL.getName());
+
+  protected final WrappedThreadPoolExecutor underlyingThreadPool;
   protected final ListeningExecutorService subtaskWorkerThreadPoolExecutor;
 
   private final Map<String, PipeSubtask> registeredIdSubtaskMapper;
@@ -50,13 +52,13 @@ public abstract class PipeSubtaskExecutor {
 
   protected PipeSubtaskExecutor(
       final int corePoolSize, final ThreadName threadName, final boolean 
disableLogInThreadPool) {
-    final WrappedThreadPoolExecutor executor =
+    underlyingThreadPool =
         (WrappedThreadPoolExecutor)
             IoTDBThreadPoolFactory.newFixedThreadPool(corePoolSize, 
threadName.getName());
     if (disableLogInThreadPool) {
-      executor.disableErrorLog();
+      underlyingThreadPool.disableErrorLog();
     }
-    subtaskWorkerThreadPoolExecutor = 
MoreExecutors.listeningDecorator(executor);
+    subtaskWorkerThreadPoolExecutor = 
MoreExecutors.listeningDecorator(underlyingThreadPool);
 
     registeredIdSubtaskMapper = new ConcurrentHashMap<>();
 
@@ -74,9 +76,11 @@ public abstract class PipeSubtaskExecutor {
 
     registeredIdSubtaskMapper.put(subtask.getTaskID(), subtask);
     subtask.bindExecutors(
-        subtaskWorkerThreadPoolExecutor,
-        subtaskCallbackListeningExecutor,
-        new PipeSubtaskScheduler(this));
+        subtaskWorkerThreadPoolExecutor, subtaskCallbackListeningExecutor, 
schedulerSupplier(this));
+  }
+
+  protected PipeSubtaskScheduler schedulerSupplier(final PipeSubtaskExecutor 
executor) {
+    return new PipeSubtaskScheduler(executor);
   }
 
   public final synchronized void start(final String subTaskID) {
@@ -160,4 +164,14 @@ public abstract class PipeSubtaskExecutor {
   public final int getRunningSubtaskNumber() {
     return runningSubtaskNumber;
   }
+
+  protected final boolean hasAvailableThread() {
+    // TODO: temporarily disable async receiver subtask execution
+    return false;
+    // return getAvailableThreadCount() > 0;
+  }
+
+  private int getAvailableThreadCount() {
+    return underlyingThreadPool.getCorePoolSize() - 
underlyingThreadPool.getActiveCount();
+  }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractConnectorSubtask.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractConnectorSubtask.java
index f228690a182..58cc142713d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractConnectorSubtask.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeAbstractConnectorSubtask.java
@@ -33,7 +33,7 @@ import 
com.google.common.util.concurrent.ListeningExecutorService;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 
 public abstract class PipeAbstractConnectorSubtask extends 
PipeReportableSubtask {
 
@@ -43,7 +43,7 @@ public abstract class PipeAbstractConnectorSubtask extends 
PipeReportableSubtask
   protected PipeConnector outputPipeConnector;
 
   // For thread pool to execute callbacks
-  protected ExecutorService subtaskCallbackListeningExecutor;
+  protected ScheduledExecutorService subtaskCallbackListeningExecutor;
 
   // For controlling subtask submitting, making sure that
   // a subtask is submitted to only one thread at a time
@@ -62,7 +62,7 @@ public abstract class PipeAbstractConnectorSubtask extends 
PipeReportableSubtask
   @Override
   public void bindExecutors(
       final ListeningExecutorService subtaskWorkerThreadPoolExecutor,
-      final ExecutorService subtaskCallbackListeningExecutor,
+      final ScheduledExecutorService subtaskCallbackListeningExecutor,
       final PipeSubtaskScheduler subtaskScheduler) {
     this.subtaskWorkerThreadPoolExecutor = subtaskWorkerThreadPoolExecutor;
     this.subtaskCallbackListeningExecutor = subtaskCallbackListeningExecutor;
@@ -217,10 +217,14 @@ public abstract class PipeAbstractConnectorSubtask 
extends PipeReportableSubtask
     }
 
     final ListenableFuture<Boolean> nextFuture = 
subtaskWorkerThreadPoolExecutor.submit(this);
-    Futures.addCallback(nextFuture, this, subtaskCallbackListeningExecutor);
+    registerCallbackHookAfterSubmit(nextFuture);
     isSubmitted = true;
   }
 
+  protected void registerCallbackHookAfterSubmit(final 
ListenableFuture<Boolean> future) {
+    Futures.addCallback(future, this, subtaskCallbackListeningExecutor);
+  }
+
   protected synchronized void setLastExceptionEvent(final Event event) {
     lastExceptionEvent = event;
   }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeSubtask.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeSubtask.java
index 2da797c2b3b..1169711b46d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeSubtask.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/subtask/PipeSubtask.java
@@ -29,7 +29,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.util.concurrent.Callable;
-import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 
@@ -65,7 +65,7 @@ public abstract class PipeSubtask
 
   public abstract void bindExecutors(
       ListeningExecutorService subtaskWorkerThreadPoolExecutor,
-      ExecutorService subtaskCallbackListeningExecutor,
+      ScheduledExecutorService subtaskCallbackListeningExecutor,
       PipeSubtaskScheduler subtaskScheduler);
 
   @Override
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
index 6a930e29ff6..fb1f37eae8d 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/subscription/config/SubscriptionConfig.java
@@ -57,8 +57,8 @@ public class SubscriptionConfig {
     return COMMON_CONFIG.getSubscriptionPollMaxBlockingTimeMs();
   }
 
-  public int getSubscriptionSerializeMaxBlockingTimeMs() {
-    return COMMON_CONFIG.getSubscriptionSerializeMaxBlockingTimeMs();
+  public int getSubscriptionDefaultTimeoutInMs() {
+    return COMMON_CONFIG.getSubscriptionDefaultTimeoutInMs();
   }
 
   public long getSubscriptionLaunchRetryIntervalMs() {
@@ -89,6 +89,10 @@ public class SubscriptionConfig {
     return COMMON_CONFIG.getSubscriptionMetaSyncerSyncIntervalMinutes();
   }
 
+  public long getSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs() {
+    return 
COMMON_CONFIG.getSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs();
+  }
+
   /////////////////////////////// Utils ///////////////////////////////
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(SubscriptionConfig.class);
@@ -113,8 +117,7 @@ public class SubscriptionConfig {
         "SubscriptionPrefetchTsFileBatchMaxSizeInBytes: {}",
         getSubscriptionPrefetchTsFileBatchMaxSizeInBytes());
     LOGGER.info("SubscriptionPollMaxBlockingTimeMs: {}", 
getSubscriptionPollMaxBlockingTimeMs());
-    LOGGER.info(
-        "SubscriptionSerializeMaxBlockingTimeMs: {}", 
getSubscriptionSerializeMaxBlockingTimeMs());
+    LOGGER.info("SubscriptionDefaultTimeoutInMs: {}", 
getSubscriptionDefaultTimeoutInMs());
     LOGGER.info("SubscriptionLaunchRetryIntervalMs: {}", 
getSubscriptionLaunchRetryIntervalMs());
     LOGGER.info(
         "SubscriptionRecycleUncommittedEventIntervalMs: {}",
@@ -124,6 +127,9 @@ public class SubscriptionConfig {
     LOGGER.info(
         "SubscriptionTsFileDeduplicationWindowSeconds: {}",
         getSubscriptionTsFileDeduplicationWindowSeconds());
+    LOGGER.info(
+        "SubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs: {}",
+        getSubscriptionTsFileSlicerCheckMemoryEnoughIntervalMs());
 
     LOGGER.info(
         "SubscriptionMetaSyncerInitialSyncDelayMinutes: {}",

Reply via email to