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

Reply via email to