This is an automated email from the ASF dual-hosted git repository.
morrySnow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 190eb6aa992 [improvement](statistics) Share partitionUpdateRows map
between analyze job and its tasks (#66488)
190eb6aa992 is described below
commit 190eb6aa992969fd3ee58e0252321068c9e8ee4b
Author: yujun <[email protected]>
AuthorDate: Wed Aug 19 16:43:09 2026 +0800
[improvement](statistics) Share partitionUpdateRows map between analyze job
and its tasks (#66488)
### What problem does this PR solve?
Problem Summary:
Every column task deep-copies the whole `partitionUpdateRows` map from
the analyze job through the `AnalysisInfo` copy constructor. With a
large number of partitions this yields an O(Task x Partition) memory
amplification: the copied maps are retained by all history records in
`analysisJobInfoMap` / `analysisTaskInfoMap` (up to
`analyze_record_limit` each), so a 4200-partition table analyzed with
many columns keeps tens of millions of `ConcurrentHashMap` nodes alive
in the FE heap.
This change makes tasks share the job's map reference instead of
deep-copying it, and clears the shared map once at the job terminal
state, so the memory held by every job/task record is released at once.
---
.../org/apache/doris/statistics/AnalysisInfo.java | 18 +-
.../org/apache/doris/statistics/AnalysisJob.java | 8 +
.../apache/doris/statistics/AnalysisManager.java | 71 +++++--
.../apache/doris/statistics/BaseAnalysisTask.java | 12 +-
.../apache/doris/statistics/OlapAnalysisTask.java | 8 +-
.../doris/statistics/StatisticsAutoCollector.java | 5 +-
.../apache/doris/statistics/AnalysisJobTest.java | 22 +++
.../doris/statistics/AnalysisManagerTest.java | 218 +++++++++++++++++++++
.../statistics/StatisticsAutoCollectorTest.java | 31 +++
9 files changed, 366 insertions(+), 27 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisInfo.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisInfo.java
index d2c293b33b1..a77c93eced6 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisInfo.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisInfo.java
@@ -189,7 +189,15 @@ public class AnalysisInfo implements Writable {
@SerializedName("tv")
public final long tableVersion;
- public final Map<Long, Long> partitionUpdateRows = new
ConcurrentHashMap<>();
+ /**
+ * Number of updated rows per partition. This map is shared between a job
and all
+ * of its tasks: each task's AnalysisInfo holds the same reference as the
job's,
+ * instead of a deep copy, so that a single clear at the job terminal
state (see
+ * AnalysisManager#updateTaskStatus) releases the memory retained by every
task
+ * record at once. It is null when the job is recovered from the image
because the
+ * field is not persisted.
+ */
+ public final Map<Long, Long> partitionUpdateRows;
@SerializedName("tblUpdateTime")
public final long tblUpdateTime;
@@ -254,9 +262,11 @@ public class AnalysisInfo implements Writable {
this.updateRows = updateRows;
this.tableVersion = tableVersion;
this.priority = priority;
- if (partitionUpdateRows != null) {
- this.partitionUpdateRows.putAll(partitionUpdateRows);
- }
+ // Share the same partitionUpdateRows map reference between the job
and all its
+ // tasks instead of deep-copying it per task, so that clearing it once
at the
+ // job terminal state releases the memory retained by every task
record at once.
+ this.partitionUpdateRows = partitionUpdateRows == null
+ ? new ConcurrentHashMap<>() : partitionUpdateRows;
this.enablePartition = enablePartition;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisJob.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisJob.java
index 8d0f3b892b9..379b0a1f9a0 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisJob.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisJob.java
@@ -182,6 +182,14 @@ public class AnalysisJob {
updateTaskState(AnalysisState.FAILED, reason);
cancel();
} finally {
+ // All task threads have been marked killed by cancel(), so the
guarded
+ // partitionUpdateRows writers have stopped. Clear the shared map
here, after
+ // the cancels (the terminal-state clear inside updateTaskStatus
runs while
+ // sibling threads may still be executing and can miss in-flight
writes), so
+ // the job and all its task records release the map at once.
+ if (jobInfo.partitionUpdateRows != null) {
+ jobInfo.partitionUpdateRows.clear();
+ }
deregisterJob();
}
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisManager.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisManager.java
index faac932743f..fbc1c33881d 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisManager.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/AnalysisManager.java
@@ -429,7 +429,9 @@ public class AnalysisManager implements Writable {
}
infoBuilder.setTableVersion(version);
infoBuilder.setPriority(JobPriority.MANUAL);
- infoBuilder.setPartitionUpdateRows(tableStatsStatus == null ? null :
tableStatsStatus.partitionUpdateRows);
+ // Must not alias TableStatsMeta.partitionUpdateRows: it is cleared at
the job terminal state.
+ infoBuilder.setPartitionUpdateRows(tableStatsStatus == null ? null
+ : new
ConcurrentHashMap<>(tableStatsStatus.partitionUpdateRows));
infoBuilder.setEnablePartition(StatisticsUtil.enablePartitionAnalyze());
return infoBuilder.build();
}
@@ -548,6 +550,13 @@ public class AnalysisManager implements Writable {
if (MetricRepo.isInit) {
MetricRepo.COUNTER_STATISTICS_FAILED_ANALYZE_JOB.increase(1L);
}
+ // The job reached a terminal state: all tasks share the
job's
+ // partitionUpdateRows map, so clearing it here releases
the memory
+ // retained by every task record in the history at once.
The success
+ // path clears it inside updateTableStats.
+ if (job.partitionUpdateRows != null) {
+ job.partitionUpdateRows.clear();
+ }
} else {
job.markFinished();
if (MetricRepo.isInit) {
@@ -567,22 +576,27 @@ public class AnalysisManager implements Writable {
@VisibleForTesting
public void updateTableStats(AnalysisInfo jobInfo) {
- TableIf tbl = StatisticsUtil.findTable(jobInfo.catalogId,
jobInfo.dbId, jobInfo.tblId);
- TableStatsMeta tableStats = findTableStatsStatus(tbl.getId());
- if (tableStats == null) {
- updateTableStatsStatus(new TableStatsMeta(jobInfo.rowCount,
jobInfo, tbl));
- } else {
- tableStats.update(jobInfo, tbl);
- logCreateTableStats(tableStats);
- }
- if (jobInfo.jobColumns != null) {
- jobInfo.jobColumns.clear();
- }
- if (jobInfo.partitionNames != null) {
- jobInfo.partitionNames.clear();
- }
- if (jobInfo.partitionUpdateRows != null) {
- jobInfo.partitionUpdateRows.clear();
+ // Clear the shared maps even when the stats update fails, so the job
and all its
+ // task records release their memory once the job reaches a terminal
state.
+ try {
+ TableIf tbl = StatisticsUtil.findTable(jobInfo.catalogId,
jobInfo.dbId, jobInfo.tblId);
+ TableStatsMeta tableStats = findTableStatsStatus(tbl.getId());
+ if (tableStats == null) {
+ updateTableStatsStatus(new TableStatsMeta(jobInfo.rowCount,
jobInfo, tbl));
+ } else {
+ tableStats.update(jobInfo, tbl);
+ logCreateTableStats(tableStats);
+ }
+ } finally {
+ if (jobInfo.jobColumns != null) {
+ jobInfo.jobColumns.clear();
+ }
+ if (jobInfo.partitionNames != null) {
+ jobInfo.partitionNames.clear();
+ }
+ if (jobInfo.partitionUpdateRows != null) {
+ jobInfo.partitionUpdateRows.clear();
+ }
}
}
@@ -992,11 +1006,22 @@ public class AnalysisManager implements Writable {
return;
}
checkPriv(anyTask);
- logKilled(analysisJobInfoMap.get(anyTask.getJobId()));
+ AnalysisInfo job = analysisJobInfoMap.get(anyTask.getJobId());
+ // The job record may have been evicted from analysisJobInfoMap while
the job was
+ // still running; its shared partitionUpdateRows map was already
cleared at
+ // eviction time (see replayCreateAnalysisJob), so skip the job-level
update.
+ if (job != null) {
+ logKilled(job);
+ }
for (BaseAnalysisTask taskInfo : analysisTaskMap.values()) {
taskInfo.cancel();
logKilled(taskInfo.info);
}
+ // The job reached a terminal state: all tasks share the job's
partitionUpdateRows
+ // map, so clearing it here releases the memory retained by every task
record.
+ if (job != null && job.partitionUpdateRows != null) {
+ job.partitionUpdateRows.clear();
+ }
}
private void logKilled(AnalysisInfo info) {
@@ -1035,7 +1060,15 @@ public class AnalysisManager implements Writable {
public void replayCreateAnalysisJob(AnalysisInfo jobInfo) {
synchronized (analysisJobInfoMap) {
while (analysisJobInfoMap.size() >= Config.analyze_record_limit) {
-
analysisJobInfoMap.remove(analysisJobInfoMap.pollFirstEntry().getKey());
+ // pollFirstEntry removes the oldest entry from the map.
+ AnalysisInfo evicted =
analysisJobInfoMap.pollFirstEntry().getValue();
+ // The evicted job may still be running: its updateTaskStatus
will return
+ // early on the "job == null" guard and never reach the
terminal-state
+ // clear, so clear the shared partitionUpdateRows map here to
release the
+ // memory held by the evicted job and all its task records at
once.
+ if (evicted.partitionUpdateRows != null) {
+ evicted.partitionUpdateRows.clear();
+ }
}
if (jobInfo.message != null && jobInfo.message.length() >=
StatisticConstants.MSG_LEN_UPPER_BOUND) {
jobInfo.message = jobInfo.message.substring(0,
StatisticConstants.MSG_LEN_UPPER_BOUND);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java
index 6ad286562ce..60e244d240a 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/BaseAnalysisTask.java
@@ -528,14 +528,22 @@ public abstract class BaseAnalysisTask {
if (partitionRowCount > hugePartitionThreshold &&
AnalysisInfo.JobType.SYSTEM.equals(info.jobType)) {
hasHughPartition = true;
// -1 means it's skipped because this partition is too
large.
- jobInfo.partitionUpdateRows.putIfAbsent(partition.getId(),
-1L);
+ // Skip the write when the task is cancelled: the job may
already be in a
+ // terminal state and have cleared the shared map, so a
late write would
+ // re-populate it and keep the memory alive.
+ if (!killed) {
+
jobInfo.partitionUpdateRows.putIfAbsent(partition.getId(), -1L);
+ }
LOG.info("Partition {} in table {} is too large, skip
it.", part, tbl.getName());
continue;
}
batchRowCount += partitionRowCount;
// For cluster upgrade compatible (older version metadata
doesn't have partition update rows map)
// and insert before first analyze, set partition update rows
to 0.
- jobInfo.partitionUpdateRows.putIfAbsent(partition.getId(), 0L);
+ // Skip the write when the task is cancelled: see the comment
above.
+ if (!killed) {
+ jobInfo.partitionUpdateRows.putIfAbsent(partition.getId(),
0L);
+ }
}
params.put("partId", partition == null ? "-1" :
String.valueOf(partition.getId()));
// Skip partitions that not changed after last analyze.
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/OlapAnalysisTask.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/OlapAnalysisTask.java
index a7ac364c984..5d49b271851 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/statistics/OlapAnalysisTask.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/OlapAnalysisTask.java
@@ -416,7 +416,13 @@ public class OlapAnalysisTask extends BaseAnalysisTask {
if (partition == null) {
columnStats.partitionUpdateRows.remove(partId);
tableStats.partitionUpdateRows.remove(partId);
- jobInfo.partitionUpdateRows.remove(partId);
+ // Skip the write when the task is cancelled: the job may
already be in a
+ // terminal state and have cleared the shared map, so no write
should touch
+ // it afterwards (a remove is a no-op on the cleared map, but
keep the
+ // invariant that cancelled tasks never write to the shared
map).
+ if (!killed) {
+ jobInfo.partitionUpdateRows.remove(partId);
+ }
expiredPartition.add(partId);
if (expiredPartition.size() ==
Config.max_allowed_in_element_num_of_delete) {
String partitionCondition = " AND part_id in (" +
Joiner.on(", ").join(expiredPartition) + ")";
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsAutoCollector.java
b/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsAutoCollector.java
index 793b43a8c93..f66dec86dc7 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsAutoCollector.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/statistics/StatisticsAutoCollector.java
@@ -45,6 +45,7 @@ import java.util.Map.Entry;
import java.util.Optional;
import java.util.Set;
import java.util.StringJoiner;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
@@ -275,7 +276,9 @@ public class StatisticsAutoCollector extends MasterDaemon {
.setUpdateRows(tableStatsStatus == null ? 0 :
tableStatsStatus.updatedRows.get())
.setTableVersion(version)
.setPriority(priority)
- .setPartitionUpdateRows(tableStatsStatus == null ? null :
tableStatsStatus.partitionUpdateRows)
+ // Must not alias TableStatsMeta.partitionUpdateRows: cleared
at the job terminal state.
+ .setPartitionUpdateRows(tableStatsStatus == null ? null
+ : new
ConcurrentHashMap<>(tableStatsStatus.partitionUpdateRows))
.setEnablePartition(StatisticsUtil.enablePartitionAnalyze())
.build();
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisJobTest.java
b/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisJobTest.java
index 95b6e378ce1..9230d0db2b3 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisJobTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisJobTest.java
@@ -17,6 +17,7 @@
package org.apache.doris.statistics;
+import org.apache.doris.statistics.AnalysisInfo.AnalysisType;
import org.apache.doris.statistics.util.StatisticsUtil;
import org.junit.jupiter.api.Assertions;
@@ -28,6 +29,8 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;
public class AnalysisJobTest {
@@ -124,6 +127,25 @@ public class AnalysisJobTest {
Assertions.assertEquals(1, writeBufInvokeTimes.get());
}
+ @Test
+ public void testTaskFailedClearsSharedPartitionUpdateRows() {
+ ConcurrentMap<Long, Long> partitionRows = new ConcurrentHashMap<>();
+ partitionRows.put(1L, 10L);
+ partitionRows.put(2L, 20L);
+ AnalysisInfo jobInfo = new
AnalysisInfoBuilder().setJobId(1).setPartitionUpdateRows(partitionRows)
+
.setState(AnalysisState.RUNNING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(AnalysisInfo.JobType.MANUAL).build();
+ BaseAnalysisTask task = Mockito.mock(BaseAnalysisTask.class);
+
+ AnalysisJob job = Mockito.spy(new AnalysisJob(jobInfo,
Collections.singletonList(task)));
+ job.analysisManager = Mockito.mock(AnalysisManager.class);
+ Mockito.doNothing().when(job).deregisterJob();
+
+ job.taskFailed(task, "boom");
+ Assertions.assertTrue(partitionRows.isEmpty(),
+ "taskFailed must clear the shared partitionUpdateRows map
after cancelling tasks");
+ }
+
@Test
public void testUpdateTaskState() {
AnalysisInfo info = Mockito.mock(AnalysisInfo.class);
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisManagerTest.java
b/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisManagerTest.java
index 30b7bd55ed5..59ebbc4f522 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisManagerTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/statistics/AnalysisManagerTest.java
@@ -20,6 +20,7 @@ package org.apache.doris.statistics;
import org.apache.doris.analysis.AnalyzeProperties;
import org.apache.doris.catalog.Column;
import org.apache.doris.catalog.Database;
+import org.apache.doris.catalog.DatabaseIf;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.OlapTable;
import org.apache.doris.catalog.PrimitiveType;
@@ -28,16 +29,24 @@ import org.apache.doris.catalog.info.TableNameInfo;
import org.apache.doris.common.Config;
import org.apache.doris.common.DdlException;
import org.apache.doris.common.Pair;
+import org.apache.doris.datasource.CatalogIf;
import org.apache.doris.datasource.InternalCatalog;
+import org.apache.doris.mysql.privilege.AccessControllerManager;
import org.apache.doris.nereids.trees.expressions.ExprId;
import org.apache.doris.nereids.trees.expressions.Slot;
import org.apache.doris.nereids.trees.expressions.SlotReference;
import org.apache.doris.nereids.trees.plans.commands.AnalyzeTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.KillAnalyzeJobCommand;
import org.apache.doris.nereids.types.IntegerType;
+import org.apache.doris.nereids.util.MoreFieldsThread;
+import org.apache.doris.persist.EditLog;
+import org.apache.doris.qe.ConnectContext;
import org.apache.doris.statistics.AnalysisInfo.AnalysisMethod;
import org.apache.doris.statistics.AnalysisInfo.AnalysisType;
import org.apache.doris.statistics.AnalysisInfo.JobType;
import org.apache.doris.statistics.AnalysisInfo.ScheduleType;
+import org.apache.doris.statistics.util.DBObjects;
+import org.apache.doris.statistics.util.StatisticsUtil;
import org.apache.doris.thrift.TQueryColumn;
import com.google.common.collect.ImmutableList;
@@ -53,6 +62,8 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;
// CHECKSTYLE OFF
@@ -254,6 +265,213 @@ public class AnalysisManagerTest {
Assertions.assertEquals(skip2, ti2.message);
}
+ @Test
+ public void testUpdateTaskStatusClearsSharedPartitionUpdateRowsOnFailure()
{
+ BaseAnalysisTask task1 = Mockito.mock(BaseAnalysisTask.class);
+ BaseAnalysisTask task2 = Mockito.mock(BaseAnalysisTask.class);
+
+ AnalysisManager manager = Mockito.spy(new AnalysisManager());
+ Mockito.doNothing().when(manager).logCreateAnalysisTask(Mockito.any());
+ Mockito.doNothing().when(manager).logCreateAnalysisJob(Mockito.any());
+ Mockito.doNothing().when(manager).updateTableStats(Mockito.any());
+
+ ConcurrentMap<Long, Long> partitionRows = new ConcurrentHashMap<>();
+ partitionRows.put(1L, 10L);
+ partitionRows.put(2L, 20L);
+ AnalysisInfo job = new
AnalysisInfoBuilder().setJobId(1).setPartitionUpdateRows(partitionRows)
+
.setState(AnalysisState.PENDING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(AnalysisInfo.JobType.MANUAL).build();
+ AnalysisInfo taskInfo1 = new
AnalysisInfoBuilder(job).setJobId(1).setTaskId(2)
+
.setJobType(JobType.MANUAL).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setState(AnalysisState.PENDING).build();
+ AnalysisInfo taskInfo2 = new
AnalysisInfoBuilder(job).setJobId(1).setTaskId(3)
+
.setJobType(JobType.MANUAL).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setState(AnalysisState.PENDING).build();
+ // Tasks share the job's partitionUpdateRows map instead of
deep-copying it.
+ Assertions.assertSame(job.partitionUpdateRows,
taskInfo1.partitionUpdateRows);
+ Assertions.assertSame(job.partitionUpdateRows,
taskInfo2.partitionUpdateRows);
+
+ manager.replayCreateAnalysisJob(job);
+ manager.replayCreateAnalysisTask(taskInfo1);
+ manager.replayCreateAnalysisTask(taskInfo2);
+ task1.info = taskInfo1;
+ task2.info = taskInfo2;
+ Map<Long, BaseAnalysisTask> tasks = new HashMap<>();
+ tasks.put(2L, task1);
+ tasks.put(3L, task2);
+ manager.addToJobIdTasksMap(1, tasks);
+
+ // Job still running (task2 pending): the shared map must survive.
+ manager.updateTaskStatus(taskInfo1, AnalysisState.FAILED, "err1", 0);
+ Assertions.assertFalse(job.partitionUpdateRows.isEmpty());
+ // Job reached terminal (FAILED) state: the shared map is cleared
exactly once,
+ // releasing the memory held by the job and all its task records at
once.
+ manager.updateTaskStatus(taskInfo2, AnalysisState.FAILED, "err2", 0);
+ Assertions.assertEquals(AnalysisState.FAILED, job.state);
+ Assertions.assertTrue(job.partitionUpdateRows.isEmpty());
+ }
+
+ @Test
+ public void testPartitionUpdateRowsJobSnapshotNotAliased() {
+ // The job must own an independent snapshot of the table stats'
partition update
+ // rows: tasks share it and it is cleared at the job terminal state,
so it must
+ // not alias the TableStatsMeta map or that clear would wipe it.
+ AnalysisManager manager = Mockito.spy(new AnalysisManager());
+ ConcurrentMap<Long, Long> tableStatsRows = new ConcurrentHashMap<>();
+ tableStatsRows.put(1L, 10L);
+ tableStatsRows.put(2L, 20L);
+ TableStatsMeta tableStatsStatus = new TableStatsMeta();
+ tableStatsStatus.partitionUpdateRows = tableStatsRows;
+
Mockito.doReturn(tableStatsStatus).when(manager).findTableStatsStatus(30001L);
+
+ Env env = Mockito.mock(Env.class);
+ try (MockedStatic<Env> envMockedStatic =
Mockito.mockStatic(Env.class)) {
+ envMockedStatic.when(Env::getCurrentEnv).thenReturn(env);
+ Mockito.when(env.getNextId()).thenReturn(1L);
+
+ AnalysisInfo job = manager.buildAnalysisJobInfo(
+ mockAnalyzeCommand(AnalysisMethod.FULL, ScheduleType.ONCE,
false, false));
+ Assertions.assertNotSame(tableStatsRows, job.partitionUpdateRows);
+ Assertions.assertEquals(tableStatsRows, job.partitionUpdateRows);
+ }
+ }
+
+ @Test
+ public void testUpdateTableStatsClearsSharedMapOnThrow() {
+ // updateTableStats must clear the shared map even when the stats
update fails
+ // (e.g. the table is dropped), so the terminal job does not leak
memory.
+ AnalysisManager manager = Mockito.spy(new AnalysisManager());
+ ConcurrentMap<Long, Long> partitionRows = new ConcurrentHashMap<>();
+ partitionRows.put(1L, 10L);
+ AnalysisInfo jobInfo = new
AnalysisInfoBuilder().setJobId(1).setPartitionUpdateRows(partitionRows)
+
.setState(AnalysisState.FINISHED).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(AnalysisInfo.JobType.MANUAL).build();
+ try (MockedStatic<StatisticsUtil> suStatic =
Mockito.mockStatic(StatisticsUtil.class)) {
+ suStatic.when(() -> StatisticsUtil.findTable(Mockito.anyLong(),
Mockito.anyLong(), Mockito.anyLong()))
+ .thenThrow(new RuntimeException("table gone"));
+ Assertions.assertThrows(RuntimeException.class, () ->
manager.updateTableStats(jobInfo));
+ }
+ Assertions.assertTrue(partitionRows.isEmpty(),
+ "shared map must be cleared even when updateTableStats
throws");
+ }
+
+ @Test
+ public void testHandleKillAnalyzeJobClearsSharedPartitionUpdateRows()
throws DdlException {
+ AnalysisManager manager = Mockito.spy(new AnalysisManager());
+ ConcurrentMap<Long, Long> partitionRows = new ConcurrentHashMap<>();
+ partitionRows.put(1L, 10L);
+ partitionRows.put(2L, 20L);
+ AnalysisInfo job = new
AnalysisInfoBuilder().setJobId(1).setPartitionUpdateRows(partitionRows)
+ .setCatalogId(10001L).setDBId(20001L).setTblId(30001L)
+
.setState(AnalysisState.RUNNING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(AnalysisInfo.JobType.MANUAL).build();
+ AnalysisInfo taskInfo = new
AnalysisInfoBuilder(job).setJobId(1).setTaskId(2)
+ .setState(AnalysisState.RUNNING).setJobType(JobType.MANUAL)
+ .setAnalysisType(AnalysisType.FUNDAMENTALS).build();
+ // Tasks share the job's partitionUpdateRows map.
+ Assertions.assertSame(job.partitionUpdateRows,
taskInfo.partitionUpdateRows);
+
+ BaseAnalysisTask task = Mockito.mock(BaseAnalysisTask.class);
+ task.info = taskInfo;
+ Mockito.when(task.getJobId()).thenReturn(1L);
+ Map<Long, BaseAnalysisTask> tasks = new HashMap<>();
+ tasks.put(2L, task);
+ manager.addToJobIdTasksMap(1, tasks);
+ manager.replayCreateAnalysisJob(job);
+ manager.replayCreateAnalysisTask(taskInfo);
+
+ ConnectContext ctx = new ConnectContext();
+ MoreFieldsThread.setConnectContext(ctx);
+ Env env = Mockito.mock(Env.class);
+ AccessControllerManager accessManager =
Mockito.mock(AccessControllerManager.class);
+ EditLog editLog = Mockito.mock(EditLog.class);
+
Mockito.when(accessManager.checkTblPriv(Mockito.any(ConnectContext.class),
Mockito.any(), Mockito.any(),
+ Mockito.any(), Mockito.any())).thenReturn(true);
+ Mockito.when(env.getAccessManager()).thenReturn(accessManager);
+ Mockito.when(env.getEditLog()).thenReturn(editLog);
+ DBObjects dbObjects = new DBObjects(Mockito.mock(CatalogIf.class),
Mockito.mock(DatabaseIf.class),
+ Mockito.mock(TableIf.class));
+ try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+ MockedStatic<StatisticsUtil> suStatic =
Mockito.mockStatic(StatisticsUtil.class)) {
+ envStatic.when(Env::getCurrentEnv).thenReturn(env);
+ suStatic.when(() -> StatisticsUtil.convertIdToObjects(10001L,
20001L, 30001L))
+ .thenReturn(dbObjects);
+ manager.handleKillAnalyzeJob(new KillAnalyzeJobCommand(1));
+ }
+ MoreFieldsThread.removeConnectContext();
+ Assertions.assertEquals(AnalysisState.FAILED, job.state);
+ Assertions.assertTrue(partitionRows.isEmpty(),
+ "shared map must be cleared when a running job is killed");
+ }
+
+ @Test
+ public void testHandleKillAnalyzeJobWithEvictedJobRecord() throws
DdlException {
+ // The job record may have been evicted from analysisJobInfoMap while
the job was
+ // still running: killing it must not NPE and must still cancel its
running tasks.
+ AnalysisManager manager = Mockito.spy(new AnalysisManager());
+ AnalysisInfo job = new
AnalysisInfoBuilder().setJobId(1).setPartitionUpdateRows(new
ConcurrentHashMap<>())
+ .setCatalogId(10001L).setDBId(20001L).setTblId(30001L)
+
.setState(AnalysisState.RUNNING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(AnalysisInfo.JobType.MANUAL).build();
+ AnalysisInfo taskInfo = new
AnalysisInfoBuilder(job).setJobId(1).setTaskId(2)
+ .setState(AnalysisState.RUNNING).setJobType(JobType.MANUAL)
+ .setAnalysisType(AnalysisType.FUNDAMENTALS).build();
+
+ BaseAnalysisTask task = Mockito.mock(BaseAnalysisTask.class);
+ task.info = taskInfo;
+ Mockito.when(task.getJobId()).thenReturn(1L);
+ Map<Long, BaseAnalysisTask> tasks = new HashMap<>();
+ tasks.put(2L, task);
+ manager.addToJobIdTasksMap(1, tasks);
+
+ ConnectContext ctx = new ConnectContext();
+ MoreFieldsThread.setConnectContext(ctx);
+ Env env = Mockito.mock(Env.class);
+ AccessControllerManager accessManager =
Mockito.mock(AccessControllerManager.class);
+ EditLog editLog = Mockito.mock(EditLog.class);
+
Mockito.when(accessManager.checkTblPriv(Mockito.any(ConnectContext.class),
Mockito.any(), Mockito.any(),
+ Mockito.any(), Mockito.any())).thenReturn(true);
+ Mockito.when(env.getAccessManager()).thenReturn(accessManager);
+ Mockito.when(env.getEditLog()).thenReturn(editLog);
+ DBObjects dbObjects = new DBObjects(Mockito.mock(CatalogIf.class),
Mockito.mock(DatabaseIf.class),
+ Mockito.mock(TableIf.class));
+ try (MockedStatic<Env> envStatic = Mockito.mockStatic(Env.class);
+ MockedStatic<StatisticsUtil> suStatic =
Mockito.mockStatic(StatisticsUtil.class)) {
+ envStatic.when(Env::getCurrentEnv).thenReturn(env);
+ suStatic.when(() -> StatisticsUtil.convertIdToObjects(10001L,
20001L, 30001L))
+ .thenReturn(dbObjects);
+ manager.handleKillAnalyzeJob(new KillAnalyzeJobCommand(1));
+ }
+ MoreFieldsThread.removeConnectContext();
+ Assertions.assertEquals(AnalysisState.FAILED, taskInfo.state,
+ "tasks of an evicted job must still be marked as killed");
+ }
+
+ @Test
+ public void testReplayCreateAnalysisJobEvictionClearsSharedMap() {
+ Config.analyze_record_limit = 2;
+ AnalysisManager manager = new AnalysisManager();
+ // The evicted job may still be running: its updateTaskStatus would
return early
+ // on the "job == null" guard and never reach the terminal-state
clear, so the
+ // shared map must be cleared at eviction time.
+ ConcurrentMap<Long, Long> evictedRows = new ConcurrentHashMap<>();
+ evictedRows.put(1L, 10L);
+ manager.replayCreateAnalysisJob(new AnalysisInfoBuilder().setJobId(1)
+
.setPartitionUpdateRows(evictedRows).setState(AnalysisState.RUNNING)
+
.setAnalysisType(AnalysisType.FUNDAMENTALS).setJobType(JobType.MANUAL).build());
+ manager.replayCreateAnalysisJob(new AnalysisInfoBuilder().setJobId(2)
+
.setState(AnalysisState.RUNNING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(JobType.MANUAL).build());
+ manager.replayCreateAnalysisJob(new AnalysisInfoBuilder().setJobId(3)
+
.setState(AnalysisState.RUNNING).setAnalysisType(AnalysisType.FUNDAMENTALS)
+ .setJobType(JobType.MANUAL).build());
+ Assertions.assertEquals(2, manager.analysisJobInfoMap.size());
+ Assertions.assertTrue(manager.analysisJobInfoMap.containsKey(2L));
+ Assertions.assertTrue(manager.analysisJobInfoMap.containsKey(3L));
+ Assertions.assertTrue(evictedRows.isEmpty(),
+ "the evicted job's shared map must be cleared at eviction");
+ }
+
@Test
public void testRecordLimit1() {
Config.analyze_record_limit = 2;
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/statistics/StatisticsAutoCollectorTest.java
b/fe/fe-core/src/test/java/org/apache/doris/statistics/StatisticsAutoCollectorTest.java
index 3312159573a..1857b536426 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/statistics/StatisticsAutoCollectorTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/statistics/StatisticsAutoCollectorTest.java
@@ -42,6 +42,8 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map.Entry;
import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
public class StatisticsAutoCollectorTest {
@@ -194,6 +196,35 @@ public class StatisticsAutoCollectorTest {
Assertions.assertEquals(10, analyzeJobForTbl.tableVersion);
}
+ @Test
+ public void testCreateAnalyzeJobForTblSnapshotPartitionUpdateRows() {
+ StatisticsAutoCollector collector = new StatisticsAutoCollector();
+ OlapTable table = Mockito.mock(OlapTable.class);
+ Database db = Mockito.mock(Database.class);
+ InternalCatalog catalog = Mockito.mock(InternalCatalog.class);
+ Mockito.when(table.getDatabase()).thenReturn(db);
+ Mockito.when(db.getCatalog()).thenReturn(catalog);
+ Mockito.when(db.getId()).thenReturn(100L);
+ Mockito.when(catalog.getId()).thenReturn(10L);
+
+ Set<Pair<String, String>> jobColumns = Sets.newHashSet();
+ jobColumns.add(Pair.of("a", "b"));
+
+ // The job must own an independent snapshot of the table stats'
partition update
+ // rows: it is cleared at the job terminal state, so it must not alias
the
+ // TableStatsMeta map or that clear would wipe it.
+ ConcurrentMap<Long, Long> tableStatsRows = new ConcurrentHashMap<>();
+ tableStatsRows.put(1L, 10L);
+ tableStatsRows.put(2L, 20L);
+ TableStatsMeta tableStatsStatus = new TableStatsMeta();
+ tableStatsStatus.partitionUpdateRows = tableStatsRows;
+
+ AnalysisInfo job = collector.createAnalyzeJobForTbl(table, jobColumns,
JobPriority.HIGH,
+ AnalysisMethod.SAMPLE, 100, tableStatsStatus, 10);
+ Assertions.assertNotSame(tableStatsRows, job.partitionUpdateRows);
+ Assertions.assertEquals(tableStatsRows, job.partitionUpdateRows);
+ }
+
@Test
public void testReadyToSample() {
StatisticsAutoCollector collector = new StatisticsAutoCollector();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]