PlanningDiagnostics.java
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
package org.apache.doris.common.profile;
import org.apache.doris.common.Config;
import org.apache.doris.common.util.DebugUtil;
import org.apache.doris.qe.ConnectContext;
import org.apache.doris.thrift.TUniqueId;
import com.google.common.annotations.VisibleForTesting;
import com.google.gson.JsonObject;
import org.apache.logging.log4j.CloseableThreadContext;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.util.ArrayList;
import java.util.EnumMap;
import java.util.List;
import java.util.Locale;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/**
* One planning pass. The planning thread owns timers and the bounded slow-operation buffer.
* The connection checker reads only immutable, volatile-published steps, never planner/table locks.
* Timings are inclusive: nested phases and operations must not be added to their enclosing phase.
*/
public final class PlanningDiagnostics {
private static final Logger LOG = LogManager.getLogger(PlanningDiagnostics.class);
private static final int MAX_SLOW_OPERATIONS = 16;
// Internal MV planning can replace ConnectContext while the calling planner still holds table locks.
private static final ThreadLocal<PlanningDiagnostics> ACTIVE = new ThreadLocal<>();
public enum Phase {
PREPROCESS, COLLECT_TABLES, PRELOAD_METADATA, WAIT_CHANGE_VISIBLE, LOCK_TABLES,
ANALYZE, REWRITE, PRE_REWRITE_MV, OPTIMIZE, CHOOSE_PLAN, POST_PROCESS,
TRANSLATE, DISTRIBUTE, RELEASE_RESOURCES;
public String key() {
return name().toLowerCase(Locale.ROOT);
}
}
private final ConnectContext context;
private final SummaryProfile summary;
private final long passId;
private final PlanningDiagnostics parent;
private final PlanningDiagnostics root;
private final boolean inheritsQueryIdentity;
private final TUniqueId queryId;
private final long statementId;
private final String threadName;
private volatile long timeoutMs;
private final LongSupplier clock;
private final long started;
private final EnumMap<Phase, Timing> timings = new EnumMap<>(Phase.class);
private final List<Event> slowOperations = new ArrayList<>();
private final AtomicLong lastReport;
private volatile Step current;
private volatile Class<?> currentJob;
private volatile HeldLocks heldLocks = new HeldLocks(0, 0, "");
private volatile boolean finished;
// Only the root publishes this pointer, so its registered connection can inspect an internal pass.
private volatile PlanningDiagnostics activePass;
private String failedStep = "";
private Throwable stepFailure;
private long omittedOperations;
private PlanningDiagnostics(ConnectContext context) {
this(context, System::nanoTime);
}
@VisibleForTesting
PlanningDiagnostics(ConnectContext context, LongSupplier clock) {
this.context = context;
this.parent = ACTIVE.get();
this.root = parent == null ? this : parent.root;
this.lastReport = parent == null ? new AtomicLong(Long.MIN_VALUE) : parent.lastReport;
SummaryProfile contextSummary = SummaryProfile.getSummaryProfile(context);
this.summary = contextSummary != null ? contextSummary : parent == null ? null : parent.summary;
this.passId = summary == null ? 0 : summary.nextPlanningPassId();
TUniqueId contextQueryId = context.queryId();
this.inheritsQueryIdentity = contextQueryId == null && parent != null;
this.queryId = inheritsQueryIdentity ? parent.queryId
: contextQueryId == null ? null : contextQueryId.deepCopy();
this.statementId = inheritsQueryIdentity ? parent.statementId : context.getStmtId();
this.threadName = Thread.currentThread().getName();
this.timeoutMs = inheritsQueryIdentity ? parent.timeoutMs : context.getExecTimeoutS() * 1000L;
this.clock = clock;
this.started = clock.getAsLong();
}
/** Includes cleanup and restores an outer pass even when internal planning uses another connection. */
public static <T> T plan(ConnectContext context, Supplier<T> action) {
return execute(new PlanningDiagnostics(context), action);
}
@VisibleForTesting
static <T> T execute(PlanningDiagnostics diagnostics, Supplier<T> action) {
ConnectContext context = diagnostics.context;
PlanningDiagnostics previous = context.getPlanningDiagnostics();
ACTIVE.set(diagnostics);
context.setPlanningDiagnostics(diagnostics);
diagnostics.root.activePass = diagnostics;
boolean success = false;
try {
T result = action.get();
success = true;
return result;
} catch (RuntimeException | Error failure) {
if (diagnostics.stepFailure != failure) {
diagnostics.failedStep = "between_phases";
}
throw failure;
} finally {
diagnostics.finished = true;
context.setPlanningDiagnostics(previous);
diagnostics.root.activePass = diagnostics.parent;
if (diagnostics.parent == null) {
ACTIVE.remove();
} else {
ACTIVE.set(diagnostics.parent);
}
diagnostics.finish(success);
}
}
public static void runPhase(ConnectContext context, Phase phase, Runnable action) {
phase(context, phase, () -> {
action.run();
return null;
});
}
public static <T> T phase(ConnectContext context, Phase phase, Supplier<T> action) {
PlanningDiagnostics diagnostics = current(context);
return diagnostics == null ? action.get() : diagnostics.measure(phase, "", "", -1, action);
}
/** A negative operation budget means that timeout handling belongs to the connector. */
public static <T> T operation(ConnectContext context, String operation, String target,
long budgetMs, Supplier<T> action) {
return operation(context, operation, () -> target, budgetMs, action);
}
/** Resolve the target only for an active planning pass, on the calling thread. */
public static <T> T operation(ConnectContext context, String operation, Supplier<String> target,
long budgetMs, Supplier<T> action) {
PlanningDiagnostics diagnostics = current(context);
if (diagnostics == null) {
return action.get();
}
Step step = diagnostics.current;
return diagnostics.measure(step == null ? Phase.PREPROCESS : step.phase,
operation, target.get(), budgetMs, action);
}
public static PlanningDiagnostics current(ConnectContext context) {
return context == null ? null : context.getPlanningDiagnostics();
}
private <T> T measure(Phase phase, String operation, String target, long budgetMs, Supplier<T> action) {
Step previous = current;
Step step = new Step(phase, operation, target, budgetMs, clock.getAsLong(),
operation.isEmpty() || previous == null ? -1 : previous.phaseStarted);
current = step;
boolean success = false;
try {
T result = action.get();
success = true;
return result;
} catch (RuntimeException | Error failure) {
// Preserve the innermost failing step of this exception, not an earlier recovered MV failure.
if (stepFailure != failure) {
stepFailure = failure;
failedStep = phase.key() + (operation.isEmpty() ? "" : "/" + operation);
}
throw failure;
} finally {
long elapsed = clock.getAsLong() - step.started;
if (operation.isEmpty()) {
Timing timing = timings.computeIfAbsent(phase, key -> new Timing());
timing.nanos += elapsed;
timing.calls++;
timing.failures += success ? 0 : 1;
if (phase == Phase.PREPROCESS && success) {
// SET_VAR hints are applied by preprocessing, after this pass was created.
timeoutMs = inheritsQueryIdentity ? parent.timeoutMs : context.getExecTimeoutS() * 1000L;
}
} else {
recordOperation(step, elapsed, success);
}
current = previous;
}
}
/** Avoid allocating task descriptions for every optimizer job; format the class only on a slow report. */
public Class<?> setCurrentJob(Class<?> job) {
Class<?> previous = currentJob;
currentJob = job;
return previous;
}
public LockHold acquiredLock(String table) {
long now = clock.getAsLong();
HeldLocks previous = heldLocks;
heldLocks = new HeldLocks(previous.count + 1, previous.count == 0 ? now : previous.started,
previous.count == 0 ? table : previous.oldestTable);
return new LockHold(table, now);
}
/** Closed after unlocking; planner resources release in reverse acquisition order. */
public final class LockHold implements AutoCloseable {
private final String table;
private final long acquired;
private LockHold(String table, long acquired) {
this.table = table;
this.acquired = acquired;
}
@Override
public void close() {
HeldLocks previous = heldLocks;
heldLocks = new HeldLocks(previous.count - 1, previous.started, previous.oldestTable);
recordOperation(new Step(Phase.RELEASE_RESOURCES, "read_lock_hold", table, -1, acquired, -1),
clock.getAsLong() - acquired, true);
}
}
private void recordOperation(Step step, long nanos, boolean success) {
long threshold = Config.nereids_planning_operation_log_threshold_ms;
if (success && (threshold <= 0 || millis(nanos) < threshold)) {
return;
}
if (slowOperations.size() == MAX_SLOW_OPERATIONS) {
omittedOperations++;
return;
}
JsonObject event = step.toJson();
event.addProperty("elapsed_ms", millis(nanos));
event.addProperty("status", success ? "completed" : "failed");
slowOperations.add(new Event(this, "Planning operation", event));
}
/** Invoked by the existing connection timeout checker, including while a planner job is blocked. */
public void reportIfSlow() {
long threshold = Config.nereids_planning_log_threshold_ms;
PlanningDiagnostics active = root.activePass;
if (root.finished || active == null || threshold <= 0) {
return;
}
// Snapshot published starts before reading the clock. A newer step/lock sampled after the clock
// could otherwise appear to have started in the future when the checker was descheduled.
Step step = active.current;
Class<?> job = active.currentJob;
HeldLocks oldest = active.heldLocks;
int lockCount = oldest.count;
for (PlanningDiagnostics outer = active.parent; outer != null; outer = outer.parent) {
HeldLocks locks = outer.heldLocks;
lockCount += locks.count;
if (locks.count > 0 && (oldest.count == 0 || locks.started < oldest.started)) {
oldest = locks;
}
}
long now = root.clock.getAsLong();
if (millis(now - root.started) < threshold) {
return;
}
long previous = root.lastReport.get();
// A minimum interval protects the checker from an accidentally zero/negative dynamic setting.
long interval = Math.max(1000, Config.nereids_planning_log_interval_ms);
if ((previous != Long.MIN_VALUE && millis(now - previous) < interval)
|| !root.lastReport.compareAndSet(previous, now)) {
return;
}
JsonObject event = step == null ? new JsonObject() : step.toJson();
event.addProperty("status", "running");
event.addProperty("elapsed_ms", millis(now - active.started));
event.addProperty("root_pass_id", root.passId);
event.addProperty("root_elapsed_ms", millis(now - root.started));
event.addProperty("phase_elapsed_ms", step == null ? 0 : millis(now - step.phaseStarted));
event.addProperty("operation_elapsed_ms", step == null || step.operation.isEmpty()
? 0 : millis(now - step.started));
event.addProperty("job", job == null ? "" : job.getSimpleName());
event.addProperty("held_locks", lockCount);
event.addProperty("oldest_lock", lockCount == 0 ? "" : bounded(oldest.oldestTable));
event.addProperty("oldest_lock_hold_ms", lockCount == 0 ? 0 : millis(now - oldest.started));
active.log("Slow planning", event);
}
private void finish(boolean success) {
JsonObject result = new JsonObject();
result.addProperty("pass_id", passId);
result.addProperty("status", success ? "completed" : "failed");
result.addProperty("elapsed_ms", millis(clock.getAsLong() - started));
result.addProperty("failed_step", success ? "" : failedStep);
JsonObject phases = new JsonObject();
for (Phase phase : Phase.values()) {
Timing timing = timings.get(phase);
JsonObject value = new JsonObject();
value.addProperty("status", timing == null ? "not_run" : timing.failures > 0 ? "failed" : "completed");
value.addProperty("elapsed_ms", timing == null ? 0 : millis(timing.nanos));
value.addProperty("calls", timing == null ? 0 : timing.calls);
phases.add(phase.key(), value);
}
result.add("phases", phases);
result.addProperty("omitted_operations", omittedOperations);
if (summary != null) {
summary.addPlanningPass(result);
}
long threshold = Config.nereids_planning_log_threshold_ms;
boolean report = !success || !slowOperations.isEmpty()
|| (threshold > 0 && result.get("elapsed_ms").getAsLong() >= threshold);
// An enclosing MV planning pass may still hold locks. Defer its children's logs as well.
if (parent != null) {
for (Event event : slowOperations) {
parent.defer(event);
}
parent.omittedOperations += omittedOperations;
if (report) {
parent.defer(new Event(this, "Planning finished", result));
}
} else {
for (Event event : slowOperations) {
event.owner.log(event.message, event.json);
}
if (report) {
log("Planning finished", result);
}
}
}
private void defer(Event event) {
if (slowOperations.size() == MAX_SLOW_OPERATIONS) {
omittedOperations++;
} else {
slowOperations.add(event);
}
}
private void log(String message, JsonObject event) {
String formattedId = queryId == null || (queryId.hi == 0 && queryId.lo == 0)
? "" : DebugUtil.printId(queryId);
event.addProperty("query_id", formattedId);
event.addProperty("pass_id", passId);
event.addProperty("statement_id", statementId);
event.addProperty("planner_thread", threadName);
event.addProperty("query_timeout_ms", timeoutMs);
// Bind only this event; the checker or a deferred event may run under another query's context.
try (CloseableThreadContext.Instance ignored = CloseableThreadContext.put("query_id", formattedId)) {
LOG.warn("{}: {}", message, event);
}
}
private static long millis(long nanos) {
return TimeUnit.NANOSECONDS.toMillis(nanos);
}
private static String bounded(String text) {
return text.substring(0, Math.min(text.length(), 256)).replace('\n', ' ').replace('\r', ' ');
}
private static final class Event {
private final PlanningDiagnostics owner;
private final String message;
private final JsonObject json;
private Event(PlanningDiagnostics owner, String message, JsonObject json) {
this.owner = owner;
this.message = message;
this.json = json;
}
}
private static final class Timing {
private long nanos;
private long calls;
private long failures;
}
private static final class HeldLocks {
private final int count;
private final long started;
private final String oldestTable;
private HeldLocks(int count, long started, String oldestTable) {
this.count = count;
this.started = started;
this.oldestTable = oldestTable;
}
}
private static final class Step {
private final Phase phase;
private final String operation;
private final String target;
private final long budgetMs;
private final long started;
private final long phaseStarted;
private Step(Phase phase, String operation, String target, long budgetMs, long started, long phaseStarted) {
this.phase = phase;
this.operation = operation;
this.target = target;
this.budgetMs = budgetMs;
this.started = started;
this.phaseStarted = phaseStarted == -1 ? started : phaseStarted;
}
private JsonObject toJson() {
JsonObject json = new JsonObject();
json.addProperty("phase", phase.key());
json.addProperty("operation", operation);
json.addProperty("target", bounded(target));
json.addProperty("operation_timeout_ms", budgetMs);
return json;
}
}
}