This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch branch-1.2-lts
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-1.2-lts by this push:
new ab2d94fd463 [opt](fe) Optimize calculate load job num metric in FE
(#39268)
ab2d94fd463 is described below
commit ab2d94fd4638e46fe2f8ccb79de1824f37f607b5
Author: htyoung <[email protected]>
AuthorDate: Tue Aug 13 15:43:11 2024 +0800
[opt](fe) Optimize calculate load job num metric in FE (#39268)
cherry pick from https://github.com/apache/doris/pull/31952
https://github.com/apache/doris/pull/34020
Co-authored-by: Lei Zhang <[email protected]>
---
.../org/apache/doris/load/loadv2/LoadManager.java | 61 ++++++++++++++--------
.../java/org/apache/doris/metric/MetricRepo.java | 23 ++++++--
2 files changed, 58 insertions(+), 26 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java
b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java
index cb825ac96e2..e23ee7ecb00 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/load/loadv2/LoadManager.java
@@ -31,6 +31,7 @@ import org.apache.doris.common.DataQualityException;
import org.apache.doris.common.DdlException;
import org.apache.doris.common.LabelAlreadyUsedException;
import org.apache.doris.common.MetaNotFoundException;
+import org.apache.doris.common.Pair;
import org.apache.doris.common.PatternMatcher;
import org.apache.doris.common.PatternMatcherWrapper;
import org.apache.doris.common.UserException;
@@ -51,7 +52,8 @@ import com.google.common.base.Strings;
import com.google.common.collect.Lists;
import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
-import org.apache.commons.lang.StringUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.commons.lang3.time.StopWatch;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@@ -69,6 +71,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantReadWriteLock;
+import java.util.function.Predicate;
import java.util.stream.Collectors;
/**
@@ -324,6 +327,11 @@ public class LoadManager implements Writable {
job.unprotectReadEndOperation(operation);
LOG.info(new LogBuilder(LogKey.LOAD_JOB,
operation.getId()).add("operation", operation)
.add("msg", "replay end load job").build());
+
+ // When idToLoadJob size increase 10000 roughly, we run
removeOldLoadJob to reduce mem used
+ if ((!idToLoadJob.isEmpty()) && (idToLoadJob.size() % 10000 == 0)) {
+ removeOldLoadJob();
+ }
}
/**
@@ -362,14 +370,10 @@ public class LoadManager implements Writable {
/**
* Get load job num, used by metric.
**/
- public long getLoadJobNum(JobState jobState, EtlJobType jobType) {
- readLock();
- try {
- return idToLoadJob.values().stream().filter(j -> j.getState() ==
jobState && j.getJobType() == jobType)
- .count();
- } finally {
- readUnlock();
- }
+ public Map<Pair<EtlJobType, JobState>, Long> getLoadJobNum() {
+ return idToLoadJob.values().stream().collect(Collectors.groupingBy(
+ loadJob -> Pair.of(loadJob.getJobType(), loadJob.getState()),
+ Collectors.counting()));
}
/**
@@ -377,30 +381,43 @@ public class LoadManager implements Writable {
**/
public void removeOldLoadJob() {
long currentTimeMs = System.currentTimeMillis();
+ removeLoadJobIf(job -> job.isExpired(currentTimeMs));
+ }
+ private void jobRemovedTrigger(LoadJob job) {
+ Map<String, List<LoadJob>> map =
dbIdToLabelToLoadJobs.get(job.getDbId());
+ List<LoadJob> list = map.get(job.getLabel());
+ list.remove(job);
+ if (job instanceof SparkLoadJob) {
+ ((SparkLoadJob) job).clearSparkLauncherLog();
+ }
+ if (list.isEmpty()) {
+ map.remove(job.getLabel());
+ }
+ if (map.isEmpty()) {
+ dbIdToLabelToLoadJobs.remove(job.getDbId());
+ }
+ }
+
+ private void removeLoadJobIf(Predicate<LoadJob> pred) {
+ long removeJobNum = 0;
+ StopWatch stopWatch = StopWatch.createStarted();
writeLock();
try {
Iterator<Map.Entry<Long, LoadJob>> iter =
idToLoadJob.entrySet().iterator();
while (iter.hasNext()) {
LoadJob job = iter.next().getValue();
- if (job.isExpired(currentTimeMs)) {
+ if (pred.test(job)) {
iter.remove();
- Map<String, List<LoadJob>> map =
dbIdToLabelToLoadJobs.get(job.getDbId());
- List<LoadJob> list = map.get(job.getLabel());
- list.remove(job);
- if (job instanceof SparkLoadJob) {
- ((SparkLoadJob) job).clearSparkLauncherLog();
- }
- if (list.isEmpty()) {
- map.remove(job.getLabel());
- }
- if (map.isEmpty()) {
- dbIdToLabelToLoadJobs.remove(job.getDbId());
- }
+ jobRemovedTrigger(job);
+ removeJobNum++;
}
}
} finally {
writeUnlock();
+ stopWatch.stop();
+ LOG.info("end to removeOldLoadJob, removeJobNum:{} cost:{} ms",
+ removeJobNum, stopWatch.getTime());
}
}
diff --git a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
index 99b36a69990..cf12a2b57e3 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/metric/MetricRepo.java
@@ -23,6 +23,7 @@ import org.apache.doris.alter.AlterJobV2.JobType;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.TabletInvertedIndex;
import org.apache.doris.common.Config;
+import org.apache.doris.common.Pair;
import org.apache.doris.common.ThreadPoolManager;
import org.apache.doris.load.EtlJobType;
import org.apache.doris.load.loadv2.JobState;
@@ -41,6 +42,7 @@ import org.apache.doris.transaction.TransactionStatus;
import com.codahale.metrics.Histogram;
import com.codahale.metrics.MetricRegistry;
+import com.google.common.collect.Maps;
import com.google.common.collect.Sets;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@@ -122,6 +124,7 @@ public final class MetricRepo {
public static GaugeMetricImpl<Double> GAUGE_REQUEST_PER_SECOND;
public static GaugeMetricImpl<Double> GAUGE_QUERY_ERR_RATE;
public static GaugeMetricImpl<Long> GAUGE_MAX_TABLET_COMPACTION_SCORE;
+ private static Map<Pair<EtlJobType, JobState>, Long> loadJobNum =
Maps.newHashMap();
private static ScheduledThreadPoolExecutor metricTimer =
ThreadPoolManager.newDaemonScheduledThreadPool(1,
"metric-timer-pool", true);
@@ -134,7 +137,6 @@ public final class MetricRepo {
}
// load jobs
- LoadManager loadManger = Env.getCurrentEnv().getLoadManager();
for (EtlJobType jobType : EtlJobType.values()) {
if (jobType == EtlJobType.UNKNOWN) {
continue;
@@ -146,7 +148,7 @@ public final class MetricRepo {
if (!Env.getCurrentEnv().isMaster()) {
return 0L;
}
- return loadManger.getLoadJobNum(state, jobType);
+ return MetricRepo.getLoadJobNum(jobType, state);
}
};
gauge.addLabel(new MetricLabel("job", "load")).addLabel(new
MetricLabel("type", jobType.name()))
@@ -632,7 +634,11 @@ public final class MetricRepo {
// update the metrics first
updateMetrics();
+ // update load job metrics
+ updateLoadJobMetrics();
+
StringBuilder sb = new StringBuilder();
+
// jvm
JvmService jvmService = new JvmService();
JvmStats jvmStats = jvmService.stats();
@@ -640,11 +646,11 @@ public final class MetricRepo {
visitor.setMetricNumber(DORIS_METRIC_REGISTER.getAllMetricSize());
// doris metrics
- for (Metric metric : DORIS_METRIC_REGISTER.getMetrics()) {
+ for (Metric<?> metric : DORIS_METRIC_REGISTER.getMetrics()) {
visitor.visit(sb, MetricVisitor.FE_PREFIX, metric);
}
// system metric
- for (Metric metric : DORIS_METRIC_REGISTER.getSystemMetrics()) {
+ for (Metric<?> metric : DORIS_METRIC_REGISTER.getSystemMetrics()) {
visitor.visit(sb, MetricVisitor.SYS_PREFIX, metric);
}
@@ -677,4 +683,13 @@ public final class MetricRepo {
public static synchronized List<Metric> getMetricsByName(String name) {
return DORIS_METRIC_REGISTER.getMetricsByName(name);
}
+
+ private static void updateLoadJobMetrics() {
+ LoadManager loadManager = Env.getCurrentEnv().getLoadManager();
+ MetricRepo.loadJobNum = loadManager.getLoadJobNum();
+ }
+
+ private static long getLoadJobNum(EtlJobType jobType, JobState jobState) {
+ return MetricRepo.loadJobNum.getOrDefault(Pair.of(jobType, jobState),
0L);
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]