This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 05a608f8ea1 Pipe: Make each connector subtask inject cron event at fix
rate to avoid random time of batch transmission (#11501)
05a608f8ea1 is described below
commit 05a608f8ea1ab641413325561645ea767459ece5
Author: Caideyipi <[email protected]>
AuthorDate: Thu Nov 9 14:24:38 2023 +0800
Pipe: Make each connector subtask inject cron event at fix rate to avoid
random time of batch transmission (#11501)
Now parallel connectors run the same time, thus the heartbeat events are
not sure to trigger the general event transfer function, causing potentially
such as the random delay of the batch transmission. Therefore, here we inject
cron events when no event can be pulled.
---------
Co-authored-by: Steve Yurong Su <[email protected]>
---
.../pipe/agent/runtime/PipeCronEventInjector.java | 4 +-
.../thrift/async/IoTDBThriftAsyncConnector.java | 2 +-
.../subtask/connector/PipeConnectorSubtask.java | 50 ++++++++++++++++------
.../apache/iotdb/commons/conf/CommonConfig.java | 11 +++++
.../iotdb/commons/conf/CommonDescriptor.java | 5 +++
.../iotdb/commons/pipe/config/PipeConfig.java | 7 +++
6 files changed, 64 insertions(+), 15 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeCronEventInjector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeCronEventInjector.java
index 2a9ae717247..2198d7608b4 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeCronEventInjector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeCronEventInjector.java
@@ -22,6 +22,7 @@ package org.apache.iotdb.db.pipe.agent.runtime;
import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.concurrent.ThreadName;
import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
+import org.apache.iotdb.commons.pipe.config.PipeConfig;
import
org.apache.iotdb.db.pipe.extractor.realtime.listener.PipeInsertionDataNodeListener;
import org.slf4j.Logger;
@@ -35,7 +36,8 @@ public class PipeCronEventInjector {
private static final Logger LOGGER =
LoggerFactory.getLogger(PipeCronEventInjector.class);
- private static final int CRON_EVENT_INJECTOR_INTERVAL_SECONDS = 30;
+ private static final long CRON_EVENT_INJECTOR_INTERVAL_SECONDS =
+
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds();
private static final ScheduledExecutorService CRON_EVENT_INJECTOR_EXECUTOR =
IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
index fdfaca268ee..9354e10aecd 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
@@ -511,7 +511,7 @@ public class IoTDBThriftAsyncConnector extends
IoTDBConnector {
// requestCommitId can not be generated by commitIdGenerator because the
commit id must
// be bind to a specific InsertTabletEvent or TsFileInsertionEvent,
otherwise the commit
- // process will be stuck.
+ // process will stuck.
final long requestCommitId = tabletBatchBuilder.getLastCommitId();
final PipeTransferTabletBatchEventHandler
pipeTransferTabletBatchEventHandler =
new PipeTransferTabletBatchEventHandler(tabletBatchBuilder, this);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/subtask/connector/PipeConnectorSubtask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/subtask/connector/PipeConnectorSubtask.java
index e2fafe00daf..c9b8c20a3db 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/subtask/connector/PipeConnectorSubtask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/subtask/connector/PipeConnectorSubtask.java
@@ -67,6 +67,16 @@ public class PipeConnectorSubtask extends PipeSubtask {
private final String attributeSortedString;
private final int connectorIndex;
+ // Now parallel connectors run the same time, thus the heartbeat events are
not sure
+ // to trigger the general event transfer function, causing potentially such
as
+ // the random delay of the batch transmission. Therefore, here we inject
cron events
+ // when no event can be pulled.
+ private static final PipeHeartbeatEvent CRON_HEARTBEAT_EVENT =
+ new PipeHeartbeatEvent("cron", false);
+ private static final long CRON_HEARTBEAT_EVENT_INJECT_INTERVAL_SECONDS =
+
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds();
+ private long lastHeartbeatEventInjectTime = System.currentTimeMillis();
+
public PipeConnectorSubtask(
String taskID,
long creationTime,
@@ -112,11 +122,16 @@ public class PipeConnectorSubtask extends PipeSubtask {
final Event event = lastEvent != null ? lastEvent :
inputPendingQueue.waitedPoll();
// Record this event for retrying on connection failure or other exceptions
setLastEvent(event);
- if (event == null) {
- return false;
- }
try {
+ if (event == null) {
+ if (System.currentTimeMillis() - lastHeartbeatEventInjectTime
+ > CRON_HEARTBEAT_EVENT_INJECT_INTERVAL_SECONDS) {
+ transferHeartbeatEvent(CRON_HEARTBEAT_EVENT);
+ }
+ return false;
+ }
+
if (event instanceof TabletInsertionEvent) {
outputPipeConnector.transfer((TabletInsertionEvent) event);
PipeConnectorMetrics.getInstance().markTabletEvent(taskID);
@@ -124,16 +139,7 @@ public class PipeConnectorSubtask extends PipeSubtask {
outputPipeConnector.transfer((TsFileInsertionEvent) event);
PipeConnectorMetrics.getInstance().markTsFileEvent(taskID);
} else if (event instanceof PipeHeartbeatEvent) {
- try {
- outputPipeConnector.heartbeat();
- outputPipeConnector.transfer(event);
- } catch (Exception e) {
- throw new PipeConnectionException(
- "PipeConnector: " + outputPipeConnector.getClass().getName() + "
heartbeat failed",
- e);
- }
- ((PipeHeartbeatEvent) event).onTransferred();
- PipeConnectorMetrics.getInstance().markPipeHeartbeatEvent(taskID);
+ transferHeartbeatEvent((PipeHeartbeatEvent) event);
} else {
outputPipeConnector.transfer(event);
}
@@ -162,6 +168,24 @@ public class PipeConnectorSubtask extends PipeSubtask {
return true;
}
+ private void transferHeartbeatEvent(PipeHeartbeatEvent event) {
+ try {
+ outputPipeConnector.heartbeat();
+ outputPipeConnector.transfer(event);
+ } catch (Exception e) {
+ throw new PipeConnectionException(
+ "PipeConnector: "
+ + outputPipeConnector.getClass().getName()
+ + " heartbeat failed, or encountered failure when transferring
generic event.",
+ e);
+ }
+
+ lastHeartbeatEventInjectTime = System.currentTimeMillis();
+
+ event.onTransferred();
+ PipeConnectorMetrics.getInstance().markPipeHeartbeatEvent(taskID);
+ }
+
@Override
public synchronized void onSuccess(Boolean hasAtLeastOneEventProcessed) {
isSubmitted = false;
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 48a3a30dda9..3678df5c5eb 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
@@ -156,6 +156,7 @@ public class CommonConfig {
private int pipeSubtaskExecutorBasicCheckPointIntervalByConsumedEventCount =
10_000;
private long pipeSubtaskExecutorBasicCheckPointIntervalByTimeDuration = 10 *
1000L;
private long pipeSubtaskExecutorPendingQueueMaxBlockingTimeMs = 1000;
+ private long pipeSubtaskExecutorCronHeartbeatEventIntervalSeconds = 30;
private int pipeExtractorAssignerDisruptorRingBufferSize = 65536;
private long pipeExtractorAssignerDisruptorRingBufferEntrySizeInBytes = 50;
// 50B
@@ -716,6 +717,16 @@ public class CommonConfig {
pipeSubtaskExecutorPendingQueueMaxBlockingTimeMs;
}
+ public long getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds() {
+ return pipeSubtaskExecutorCronHeartbeatEventIntervalSeconds;
+ }
+
+ public void setPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds(
+ long pipeSubtaskExecutorCronHeartbeatEventIntervalSeconds) {
+ this.pipeSubtaskExecutorCronHeartbeatEventIntervalSeconds =
+ pipeSubtaskExecutorCronHeartbeatEventIntervalSeconds;
+ }
+
public void setPipeAirGapReceiverEnabled(boolean pipeAirGapReceiverEnabled) {
this.pipeAirGapReceiverEnabled = pipeAirGapReceiverEnabled;
}
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 cc794954a78..7bda519fc49 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
@@ -287,6 +287,11 @@ public class CommonDescriptor {
properties.getProperty(
"pipe_subtask_executor_pending_queue_max_blocking_time_ms",
String.valueOf(config.getPipeSubtaskExecutorPendingQueueMaxBlockingTimeMs()))));
+ config.setPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds(
+ Long.parseLong(
+ properties.getProperty(
+ "pipe_subtask_executor_cron_heartbeat_event_interval_seconds",
+
String.valueOf(config.getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds()))));
config.setPipeExtractorAssignerDisruptorRingBufferSize(
Integer.parseInt(
diff --git
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
index 47e3f1d5003..081a58a084f 100644
---
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/PipeConfig.java
@@ -71,6 +71,10 @@ public class PipeConfig {
return COMMON_CONFIG.getPipeSubtaskExecutorPendingQueueMaxBlockingTimeMs();
}
+ public long getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds() {
+ return
COMMON_CONFIG.getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds();
+ }
+
/////////////////////////////// Extractor ///////////////////////////////
public int getPipeExtractorAssignerDisruptorRingBufferSize() {
@@ -205,6 +209,9 @@ public class PipeConfig {
LOGGER.info(
"PipeSubtaskExecutorPendingQueueMaxBlockingTimeMs: {}",
getPipeSubtaskExecutorPendingQueueMaxBlockingTimeMs());
+ LOGGER.info(
+ "PipeSubtaskExecutorCronHeartbeatEventIntervalSeconds: {}",
+ getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds());
LOGGER.info(
"PipeExtractorAssignerDisruptorRingBufferSize: {}",