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: {}",