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

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 5f72265e23a branch-4.1: [fix](job) Avoid adaptive batching for Kafka 
tasks with low lag (#66863)
5f72265e23a is described below

commit 5f72265e23a4aea3896d4879e7498add42236d20
Author: hui lai <[email protected]>
AuthorDate: Thu Aug 20 09:22:28 2026 +0800

    branch-4.1: [fix](job) Avoid adaptive batching for Kafka tasks with low lag 
(#66863)
    
    pick #66099.
---
 .../load/routineload/KafkaRoutineLoadJob.java      |  19 +++
 .../doris/load/routineload/KafkaTaskInfo.java      |  16 +-
 .../load/routineload/KafkaRoutineLoadJobTest.java  | 162 +++++++++++++++++++++
 .../test_routine_load_adaptive_param.groovy        |  16 +-
 4 files changed, 202 insertions(+), 11 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaRoutineLoadJob.java
 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaRoutineLoadJob.java
index 579f860b46e..4a8f4401ae6 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaRoutineLoadJob.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaRoutineLoadJob.java
@@ -913,6 +913,25 @@ public class KafkaRoutineLoadJob extends RoutineLoadJob {
         }
     }
 
+    boolean isTaskLagGreaterThanMaxBatchRows(Map<Integer, Long> 
partitionIdToOffset) {
+        long remainingRows = Math.max(maxBatchRows, 
RoutineLoadJob.DEFAULT_MAX_BATCH_ROWS);
+        for (Map.Entry<Integer, Long> entry : partitionIdToOffset.entrySet()) {
+            Long latestOffset = 
cachedPartitionWithLatestOffsets.get(entry.getKey());
+            if (latestOffset == null) {
+                continue;
+            }
+            long partitionLag = latestOffset - entry.getValue();
+            if (partitionLag <= 0) {
+                continue;
+            }
+            if (partitionLag > remainingRows) {
+                return true;
+            }
+            remainingRows -= partitionLag;
+        }
+        return false;
+    }
+
     // check if given partitions has more data to consume.
     // 'partitionIdToOffset' to the offset to be consumed.
     public boolean hasMoreDataToConsume(UUID taskId, Map<Integer, Long> 
partitionIdToOffset) throws UserException {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaTaskInfo.java 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaTaskInfo.java
index eec36a5af9b..6b8c7636a00 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaTaskInfo.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/KafkaTaskInfo.java
@@ -23,6 +23,7 @@ import org.apache.doris.catalog.OlapTable;
 import org.apache.doris.catalog.Table;
 import org.apache.doris.common.Config;
 import org.apache.doris.common.UserException;
+import org.apache.doris.common.util.DebugPointUtil;
 import org.apache.doris.common.util.DebugUtil;
 import org.apache.doris.nereids.load.NereidsLoadTaskInfo;
 import org.apache.doris.nereids.load.NereidsStreamLoadPlanner;
@@ -56,6 +57,9 @@ public class KafkaTaskInfo extends RoutineLoadTaskInfo {
     // <partitionId, offset to be consumed>
     private Map<Integer, Long> partitionIdToOffset;
 
+    private int adaptiveMinBatchInterval;
+    private boolean isAdaptiveBatch;
+
     public KafkaTaskInfo(UUID id, long jobId,
                          long timeoutMs, Map<Integer, Long> 
partitionIdToOffset, boolean isMultiTable,
                          long lastScheduledTime, boolean isEof) {
@@ -122,9 +126,13 @@ public class KafkaTaskInfo extends RoutineLoadTaskInfo {
 
     @Override
     public void updateAdaptiveTimeout(RoutineLoadJob routineLoadJob) {
-        if (!isEof) {
+        adaptiveMinBatchInterval = 
Config.routine_load_adaptive_min_batch_interval_sec;
+        KafkaRoutineLoadJob kafkaRoutineLoadJob = (KafkaRoutineLoadJob) 
routineLoadJob;
+        isAdaptiveBatch = 
DebugPointUtil.isEnable("KafkaTaskInfo.shouldUseAdaptiveBatch")
+                || 
kafkaRoutineLoadJob.isTaskLagGreaterThanMaxBatchRows(partitionIdToOffset);
+        if (isAdaptiveBatch) {
             long maxBatchIntervalS = 
Math.max(routineLoadJob.getMaxBatchIntervalS(),
-                    Config.routine_load_adaptive_min_batch_interval_sec);
+                    adaptiveMinBatchInterval);
             long timeoutSec = maxBatchIntervalS * 
Config.routine_load_task_timeout_multiplier;
             long realTimeoutSec = Math.max(timeoutSec, 
Config.routine_load_task_min_timeout_sec);
             this.timeoutMs = realTimeoutSec * 1000;
@@ -137,8 +145,8 @@ public class KafkaTaskInfo extends RoutineLoadTaskInfo {
         long maxBatchIntervalS = routineLoadJob.getMaxBatchIntervalS();
         long maxBatchRows = routineLoadJob.getMaxBatchRows();
         long maxBatchSize = routineLoadJob.getMaxBatchSizeBytes();
-        if (!isEof) {
-            maxBatchIntervalS = Math.max(maxBatchIntervalS, 
Config.routine_load_adaptive_min_batch_interval_sec);
+        if (isAdaptiveBatch) {
+            maxBatchIntervalS = Math.max(maxBatchIntervalS, 
adaptiveMinBatchInterval);
             maxBatchRows = Math.max(maxBatchRows, 
RoutineLoadJob.DEFAULT_MAX_BATCH_ROWS);
             maxBatchSize = Math.max(maxBatchSize, 
RoutineLoadJob.DEFAULT_MAX_BATCH_SIZE);
         }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java
index 19ebc50ffbb..3487114ebbf 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java
@@ -48,6 +48,7 @@ import 
org.apache.doris.nereids.trees.plans.commands.load.LoadSeparator;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.system.SystemInfoService;
 import org.apache.doris.thrift.TResourceInfo;
+import org.apache.doris.thrift.TRoutineLoadTask;
 import org.apache.doris.transaction.BeginTransactionException;
 import org.apache.doris.transaction.GlobalTransactionMgr;
 
@@ -66,6 +67,8 @@ import org.apache.logging.log4j.Logger;
 import org.junit.Assert;
 import org.junit.Before;
 import org.junit.Test;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
 
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -396,6 +399,165 @@ public class KafkaRoutineLoadJobTest {
         Assert.assertTrue(newTask.needDedalySchedule());
     }
 
+    @Test
+    public void testAdaptiveBatchUsesTaskLagThreshold() {
+        RoutineLoadManager routineLoadManager = 
Mockito.mock(RoutineLoadManager.class);
+        Env env = Mockito.mock(Env.class);
+        int previousAdaptiveIntervalSec = 
Config.routine_load_adaptive_min_batch_interval_sec;
+
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class)) {
+            Config.routine_load_adaptive_min_batch_interval_sec = 360;
+            envStatic.when(Env::getCurrentEnv).thenReturn(env);
+            
Mockito.when(env.getRoutineLoadManager()).thenReturn(routineLoadManager);
+
+            KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, 
"kafka_routine_load_job", 1L,
+                    1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN);
+            Deencapsulation.setField(routineLoadJob, "maxBatchIntervalS", 20L);
+            Deencapsulation.setField(routineLoadJob, "maxBatchRows", 200000L);
+            Deencapsulation.setField(routineLoadJob, "maxBatchSizeBytes", 100L 
* 1024 * 1024);
+            
Mockito.when(routineLoadManager.getJob(1L)).thenReturn(routineLoadJob);
+
+            Map<Integer, Long> taskProgress = Maps.newHashMap();
+            taskProgress.put(1, 10L);
+            taskProgress.put(2, 20L);
+
+            KafkaTaskInfo taskWithUnknownLag = new KafkaTaskInfo(new UUID(1, 
1), 1L, 20000,
+                    taskProgress, false, 1000, false);
+            taskWithUnknownLag.updateAdaptiveTimeout(routineLoadJob);
+            TRoutineLoadTask unknownLagThriftTask = new TRoutineLoadTask();
+            Deencapsulation.invoke(
+                    taskWithUnknownLag, "adaptiveBatchParam", 
unknownLagThriftTask, routineLoadJob);
+            Assert.assertEquals(20L, unknownLagThriftTask.getMaxIntervalS());
+
+            Map<Integer, Long> latestOffsets = Maps.newHashMap();
+            latestOffsets.put(1, 10_000_010L);
+            latestOffsets.put(2, 10_000_020L);
+            latestOffsets.put(3, 100_000_000L);
+            Deencapsulation.setField(routineLoadJob, 
"cachedPartitionWithLatestOffsets", latestOffsets);
+
+            KafkaTaskInfo taskAtLagThreshold = new KafkaTaskInfo(new UUID(1, 
2), 1L, 20000,
+                    taskProgress, false, 1000, false);
+            taskAtLagThreshold.updateAdaptiveTimeout(routineLoadJob);
+            latestOffsets.put(2, 10_000_021L);
+            TRoutineLoadTask thresholdThriftTask = new TRoutineLoadTask();
+            Deencapsulation.invoke(
+                    taskAtLagThreshold, "adaptiveBatchParam", 
thresholdThriftTask, routineLoadJob);
+            Assert.assertEquals(20L, thresholdThriftTask.getMaxIntervalS());
+            Assert.assertEquals(200000L, 
thresholdThriftTask.getMaxBatchRows());
+            Assert.assertEquals(100L * 1024 * 1024, 
thresholdThriftTask.getMaxBatchSize());
+            Assert.assertEquals(routineLoadJob.getTimeout() * 1000L, 
taskAtLagThreshold.getTimeoutMs());
+
+            KafkaTaskInfo taskAboveLagThreshold = new KafkaTaskInfo(new 
UUID(1, 3), 1L, 20000,
+                    taskProgress, false, 1000, true);
+            taskAboveLagThreshold.updateAdaptiveTimeout(routineLoadJob);
+            TRoutineLoadTask adaptiveThriftTask = new TRoutineLoadTask();
+            Deencapsulation.invoke(
+                    taskAboveLagThreshold, "adaptiveBatchParam", 
adaptiveThriftTask, routineLoadJob);
+            Assert.assertEquals(360L, adaptiveThriftTask.getMaxIntervalS());
+            Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_ROWS,
+                    adaptiveThriftTask.getMaxBatchRows());
+            Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_SIZE,
+                    adaptiveThriftTask.getMaxBatchSize());
+            Assert.assertEquals(360L * 
Config.routine_load_task_timeout_multiplier * 1000,
+                    taskAboveLagThreshold.getTimeoutMs());
+
+            Deencapsulation.setField(routineLoadJob, "maxBatchRows", 
50_000_000L);
+            latestOffsets.put(1, 25_000_010L);
+            latestOffsets.put(2, 25_000_020L);
+            KafkaTaskInfo taskAtConfiguredLagThreshold = new KafkaTaskInfo(new 
UUID(1, 4), 1L, 20000,
+                    taskProgress, false, 1000, false);
+            taskAtConfiguredLagThreshold.updateAdaptiveTimeout(routineLoadJob);
+            latestOffsets.put(2, 25_000_021L);
+            TRoutineLoadTask configuredThresholdThriftTask = new 
TRoutineLoadTask();
+            Deencapsulation.invoke(taskAtConfiguredLagThreshold, 
"adaptiveBatchParam",
+                    configuredThresholdThriftTask, routineLoadJob);
+            Assert.assertEquals(20L, 
configuredThresholdThriftTask.getMaxIntervalS());
+            Assert.assertEquals(50_000_000L, 
configuredThresholdThriftTask.getMaxBatchRows());
+            Assert.assertEquals(routineLoadJob.getTimeout() * 1000L,
+                    taskAtConfiguredLagThreshold.getTimeoutMs());
+
+            KafkaTaskInfo taskAboveConfiguredLagThreshold = new 
KafkaTaskInfo(new UUID(1, 5), 1L, 20000,
+                    taskProgress, false, 1000, false);
+            
taskAboveConfiguredLagThreshold.updateAdaptiveTimeout(routineLoadJob);
+            TRoutineLoadTask configuredAdaptiveThriftTask = new 
TRoutineLoadTask();
+            Deencapsulation.invoke(taskAboveConfiguredLagThreshold, 
"adaptiveBatchParam",
+                    configuredAdaptiveThriftTask, routineLoadJob);
+            Assert.assertEquals(360L, 
configuredAdaptiveThriftTask.getMaxIntervalS());
+            Assert.assertEquals(50_000_000L, 
configuredAdaptiveThriftTask.getMaxBatchRows());
+        } finally {
+            Config.routine_load_adaptive_min_batch_interval_sec = 
previousAdaptiveIntervalSec;
+        }
+    }
+
+    @Test
+    public void testAdaptiveBatchUsesCapturedIntervalAcrossConfigChange() {
+        RoutineLoadManager routineLoadManager = 
Mockito.mock(RoutineLoadManager.class);
+        Env env = Mockito.mock(Env.class);
+        int previousAdaptiveIntervalSec = 
Config.routine_load_adaptive_min_batch_interval_sec;
+
+        try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class)) {
+            Config.routine_load_adaptive_min_batch_interval_sec = 360;
+            envStatic.when(Env::getCurrentEnv).thenReturn(env);
+            
Mockito.when(env.getRoutineLoadManager()).thenReturn(routineLoadManager);
+
+            KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, 
"kafka_routine_load_job", 1L,
+                    1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN);
+            Deencapsulation.setField(routineLoadJob, "maxBatchIntervalS", 30L);
+            Deencapsulation.setField(routineLoadJob, "maxBatchRows", 200000L);
+            Deencapsulation.setField(routineLoadJob, "maxBatchSizeBytes", 100L 
* 1024 * 1024);
+            
Mockito.when(routineLoadManager.getJob(1L)).thenReturn(routineLoadJob);
+
+            Map<Integer, Long> taskProgress = Maps.newHashMap();
+            taskProgress.put(1, 10L);
+            Map<Integer, Long> latestOffsets = Maps.newHashMap();
+            latestOffsets.put(1, 20_000_011L);
+            Deencapsulation.setField(routineLoadJob, 
"cachedPartitionWithLatestOffsets", latestOffsets);
+
+            KafkaTaskInfo scheduledTask = new KafkaTaskInfo(new UUID(1, 6), 
1L, 20000,
+                    taskProgress, false, 1000, false);
+            scheduledTask.updateAdaptiveTimeout(routineLoadJob);
+            long adaptiveTimeoutMs = 360L * 
Config.routine_load_task_timeout_multiplier * 1000L;
+            Assert.assertEquals(adaptiveTimeoutMs, 
scheduledTask.getTimeoutMs());
+
+            Config.routine_load_adaptive_min_batch_interval_sec = 720;
+            TRoutineLoadTask scheduledThriftTask = new TRoutineLoadTask();
+            Deencapsulation.invoke(scheduledTask, "adaptiveBatchParam", 
scheduledThriftTask, routineLoadJob);
+            Assert.assertEquals(360L, scheduledThriftTask.getMaxIntervalS());
+            Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_ROWS, 
scheduledThriftTask.getMaxBatchRows());
+            Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_SIZE, 
scheduledThriftTask.getMaxBatchSize());
+            Assert.assertEquals(adaptiveTimeoutMs, 
scheduledTask.getTimeoutMs());
+
+            KafkaTaskInfo nextSchedulingAttempt = new KafkaTaskInfo(new 
UUID(1, 7), 1L, 20000,
+                    taskProgress, false, 1000, false);
+            nextSchedulingAttempt.updateAdaptiveTimeout(routineLoadJob);
+            TRoutineLoadTask nextThriftTask = new TRoutineLoadTask();
+            Deencapsulation.invoke(nextSchedulingAttempt, 
"adaptiveBatchParam", nextThriftTask, routineLoadJob);
+            Assert.assertEquals(720L, nextThriftTask.getMaxIntervalS());
+            Assert.assertEquals(720L * 
Config.routine_load_task_timeout_multiplier * 1000L,
+                    nextSchedulingAttempt.getTimeoutMs());
+
+            for (int nonPositiveInterval : new int[] {0, -1}) {
+                Config.routine_load_adaptive_min_batch_interval_sec = 
nonPositiveInterval;
+                KafkaTaskInfo nonPositiveConfigTask = new KafkaTaskInfo(new 
UUID(1, 8), 1L, 20000,
+                        taskProgress, false, 1000, false);
+                nonPositiveConfigTask.updateAdaptiveTimeout(routineLoadJob);
+                TRoutineLoadTask nonPositiveConfigThriftTask = new 
TRoutineLoadTask();
+                Deencapsulation.invoke(nonPositiveConfigTask, 
"adaptiveBatchParam",
+                        nonPositiveConfigThriftTask, routineLoadJob);
+                Assert.assertEquals(30L, 
nonPositiveConfigThriftTask.getMaxIntervalS());
+                Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_ROWS,
+                        nonPositiveConfigThriftTask.getMaxBatchRows());
+                Assert.assertEquals(RoutineLoadJob.DEFAULT_MAX_BATCH_SIZE,
+                        nonPositiveConfigThriftTask.getMaxBatchSize());
+                long normalTimeoutMs = Math.max(30L * 
Config.routine_load_task_timeout_multiplier,
+                        Config.routine_load_task_min_timeout_sec) * 1000L;
+                Assert.assertEquals(normalTimeoutMs, 
nonPositiveConfigTask.getTimeoutMs());
+            }
+        } finally {
+            Config.routine_load_adaptive_min_batch_interval_sec = 
previousAdaptiveIntervalSec;
+        }
+    }
+
     @Test
     public void testUpdateLagRebuildsConvertedPropertiesAfterReplay(@Mocked 
Env env) throws UserException {
         KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, 
"kafka_routine_load_job", 1L,
diff --git 
a/regression-test/suites/load_p0/routine_load/test_routine_load_adaptive_param.groovy
 
b/regression-test/suites/load_p0/routine_load/test_routine_load_adaptive_param.groovy
index 49d901c31c9..8652d018780 100644
--- 
a/regression-test/suites/load_p0/routine_load/test_routine_load_adaptive_param.groovy
+++ 
b/regression-test/suites/load_p0/routine_load/test_routine_load_adaptive_param.groovy
@@ -65,31 +65,33 @@ suite("test_routine_load_adaptive_param","nonConcurrent") {
                 );
             """
 
-            def injection = "RoutineLoadTaskInfo.judgeEof"
+            def eofInjection = "RoutineLoadTaskInfo.judgeEof"
+            def adaptiveBatchInjection = "KafkaTaskInfo.shouldUseAdaptiveBatch"
             try {
-                GetDebugPoint().enableDebugPointForAllFEs(injection)
+                GetDebugPoint().enableDebugPointForAllFEs(eofInjection)
+                
GetDebugPoint().enableDebugPointForAllFEs(adaptiveBatchInjection)
 
                 RoutineLoadTestUtils.sendTestDataToKafka(producer, 
kafkaCsvTpoics)
                 RoutineLoadTestUtils.waitForTaskFinish(runSql, job, tableName, 
0)
                 
                 logger.info("---test adaptively increase---")
                 RoutineLoadTestUtils.sendTestDataToKafka(producer, 
kafkaCsvTpoics)
-                // Drive data each round so an isEof=false task keeps being 
scheduled. The converged
-                // adaptive timeout (3600) lives on the renewed idle task 
(txnId == -1), so both checks
+                // Drive data each round so a task keeps being scheduled. The 
converged adaptive timeout
+                // (3600) lives on the renewed idle task (txnId == -1), so 
both checks
                 // poll by value (task timeout col, and the committed txn's 
persisted timeout looked up
                 // by task-UUID label) instead of racing a sub-second running 
task.
                 RoutineLoadTestUtils.checkTaskTimeoutWithData(runSql, 
producer, kafkaCsvTpoics, job, "3600")
                 RoutineLoadTestUtils.checkTxnTimeoutMatchesTaskTimeout(runSql, 
producer, kafkaCsvTpoics, job, "3600000")
                 RoutineLoadTestUtils.waitForTaskFinish(runSql, job, tableName, 
2)
             } finally {
-                GetDebugPoint().disableDebugPointForAllFEs(injection)
+                
GetDebugPoint().disableDebugPointForAllFEs(adaptiveBatchInjection)
+                GetDebugPoint().disableDebugPointForAllFEs(eofInjection)
             }
 
             logger.info("---test restore adaptively---")
             RoutineLoadTestUtils.sendTestDataToKafka(producer, kafkaCsvTpoics)
             RoutineLoadTestUtils.waitForTaskFinish(runSql, job, tableName, 4)
-            // After EOF the adaptive timeout only converges when an isEof 
task is scheduled with
-            // data, so keep feeding small batches until the task timeout 
restores to the job timeout.
+            // Keep feeding small batches until the low task lag restores the 
timeout to the job timeout.
             RoutineLoadTestUtils.checkTaskTimeoutWithData(runSql, producer, 
kafkaCsvTpoics, job, "100")
         } finally {
             sql "stop routine load for ${job}"


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to