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]