This is an automated email from the ASF dual-hosted git repository. FrankChen021 pushed a commit to branch codex/fix-dart-report-publication in repository https://gitbox.apache.org/repos/asf/druid.git
commit b29c66d8582b8f666ea2a0347723922774f067c1 Author: Frank Chen <[email protected]> AuthorDate: Sat Sep 12 23:14:36 2026 +0800 fix: separate live and final MSQ reports --- .../druid/msq/exec/CaptureReportQueryListener.java | 42 +++++++------------ .../java/org/apache/druid/msq/exec/Controller.java | 11 +++++ .../apache/druid/msq/exec/ControllerHolder.java | 47 +++++++++------------- .../org/apache/druid/msq/exec/ControllerImpl.java | 22 ++++++++-- .../druid/msq/exec/ControllerHolderTest.java | 14 +++++++ 5 files changed, 78 insertions(+), 58 deletions(-) diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/CaptureReportQueryListener.java b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/CaptureReportQueryListener.java index 7632ca2ecf2..74a9c214cc4 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/CaptureReportQueryListener.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/CaptureReportQueryListener.java @@ -19,46 +19,32 @@ package org.apache.druid.msq.exec; -import org.apache.druid.error.DruidException; import org.apache.druid.frame.read.FrameReader; +import org.apache.druid.indexer.report.TaskReport; +import org.apache.druid.msq.indexing.report.MSQTaskReport; import org.apache.druid.msq.indexing.report.MSQTaskReportPayload; import org.apache.druid.query.rowsandcols.RowsAndColumns; -import javax.annotation.Nullable; +import java.util.concurrent.atomic.AtomicReference; /** - * A {@link QueryListener} wrapper that captures the report from {@link #onQueryComplete(MSQTaskReportPayload)}. + * A {@link QueryListener} wrapper that captures the final report before delegating query completion. */ public class CaptureReportQueryListener implements QueryListener { private final QueryListener delegate; + private final String queryId; + private final AtomicReference<TaskReport.ReportMap> reportReference; - @Nullable - private volatile MSQTaskReportPayload report; - - public CaptureReportQueryListener(final QueryListener delegate) + public CaptureReportQueryListener( + final QueryListener delegate, + final String queryId, + final AtomicReference<TaskReport.ReportMap> reportReference + ) { this.delegate = delegate; - } - - /** - * Whether this listener has captured a report. Will be true if the query has completed, false otherwise. - */ - public boolean hasReport() - { - return report != null; - } - - /** - * Retrieves the report. Can only be called once the query is complete. - */ - public MSQTaskReportPayload getReport() - { - if (report == null) { - throw DruidException.defensive("Query not complete, cannot call getReport()"); - } - - return report; + this.queryId = queryId; + this.reportReference = reportReference; } @Override @@ -88,7 +74,7 @@ public class CaptureReportQueryListener implements QueryListener @Override public void onQueryComplete(final MSQTaskReportPayload report) { - this.report = report; + reportReference.set(TaskReport.buildTaskReports(new MSQTaskReport(queryId, report))); delegate.onQueryComplete(report); } } diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/Controller.java b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/Controller.java index 75b47e7e262..816f28c4eeb 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/Controller.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/Controller.java @@ -26,6 +26,7 @@ import org.apache.druid.msq.indexing.MSQControllerTask; import org.apache.druid.msq.indexing.client.ControllerChatHandler; import org.apache.druid.msq.indexing.error.CancellationReason; import org.apache.druid.msq.indexing.error.MSQErrorReport; +import org.apache.druid.msq.indexing.report.MSQTaskReportPayload; import org.apache.druid.msq.kernel.StageId; import org.apache.druid.msq.statistics.PartialKeyStatisticsInformation; import org.apache.druid.query.QueryContext; @@ -132,9 +133,19 @@ public interface Controller */ boolean hasWorker(String workerId); + /** + * Returns the current live report snapshot, or null if query execution has not started. + */ @Nullable TaskReport.ReportMap liveReports(); + /** + * Returns the final report once the query is complete, or null while the query is running. + * The final report is published before {@link QueryListener#onQueryComplete(MSQTaskReportPayload)} is called. + */ + @Nullable + TaskReport.ReportMap finalReport(); + ControllerContext getControllerContext(); QueryContext getQueryContext(); diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerHolder.java b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerHolder.java index 465f6e5c74e..09a9b57ba8d 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerHolder.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerHolder.java @@ -67,10 +67,6 @@ public class ControllerHolder private final DateTime startTime; - // Published before notifying the result listener, so report readers see final counters once results end. - @Nullable - private volatile TaskReport.ReportMap finalReport; - @GuardedBy("this") private State state = State.ACCEPTED; @@ -111,7 +107,7 @@ public class ControllerHolder @Nullable public TaskReport.ReportMap getReports() { - final TaskReport.ReportMap report = finalReport; + final TaskReport.ReportMap report = controller.finalReport(); return report == null ? controller.liveReports() : report; } @@ -181,20 +177,12 @@ public class ControllerHolder Thread.currentThread().setName(makeThreadName()); try { - final CaptureReportQueryListener reportListener = new CaptureReportQueryListener(listener) - { - @Override - public void onQueryComplete(final MSQTaskReportPayload report) - { - finalReport = TaskReport.buildTaskReports(new MSQTaskReport(controller.queryId(), report)); - super.onQueryComplete(report); - } - }; + TaskReport.ReportMap reportMap = null; try { if (transitionToRunning()) { try { - controller.run(reportListener); + controller.run(listener); } finally { synchronized (this) { @@ -205,11 +193,22 @@ public class ControllerHolder } } - updateStateOnQueryComplete(reportListener.getReport()); + reportMap = controller.finalReport(); + if (reportMap != null) { + final TaskReport taskReport = reportMap.get(MSQTaskReport.REPORT_KEY); + if (taskReport instanceof MSQTaskReport) { + final MSQTaskReportPayload report = ((MSQTaskReport) taskReport).getPayload(); + if (report != null) { + updateStateOnQueryComplete(report); + } + } + } } else { // Canceled before running. + final MSQTaskReportPayload canceledReport = makeCanceledReport(cancelReason); + reportMap = TaskReport.buildTaskReports(new MSQTaskReport(controller.queryId(), canceledReport)); synchronized (this) { - reportListener.onQueryComplete(makeCanceledReport(cancelReason)); + listener.onQueryComplete(canceledReport); } } } @@ -223,18 +222,12 @@ public class ControllerHolder ); } finally { - // Build report and then call "deregister". - final MSQTaskReport taskReport; - - if (reportListener.hasReport()) { - taskReport = new MSQTaskReport(controller.queryId(), reportListener.getReport()); - } else { - taskReport = null; + if (reportMap == null) { + // ControllerImpl publishes its final report before invoking the completion listener, including when + // controller.run() exits with an exception. + reportMap = controller.finalReport(); } - final TaskReport.ReportMap reportMap = new TaskReport.ReportMap(); - reportMap.put(MSQTaskReport.REPORT_KEY, taskReport); - if (controllerRegistry != null) { controllerRegistry.deregister(this, reportMap); } diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerImpl.java b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerImpl.java index 17a729f2348..030af55c758 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerImpl.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/ControllerImpl.java @@ -274,6 +274,9 @@ public class ControllerImpl implements Controller // For live reports. Written by the main controller thread, read by HTTP threads. private final AtomicReference<QueryDefinition> queryDefRef = new AtomicReference<>(); + // Final report. Written before notifying the query listener, read by HTTP threads. + private final AtomicReference<TaskReport.ReportMap> finalReport = new AtomicReference<>(); + // Last reported CounterSnapshots per stage per worker // For live reports. Written by the main controller thread, read by HTTP threads. private final CounterSnapshotsTree taskCountersForLiveReports = new CounterSnapshotsTree(); @@ -362,17 +365,23 @@ public class ControllerImpl implements Controller @Override public void run(final QueryListener queryListener) throws Exception { + final QueryListener reportPublishingListener = new CaptureReportQueryListener( + queryListener, + queryId(), + finalReport + ); + final MSQTaskReportPayload reportPayload; try (final Closer closer = Closer.create()) { - reportPayload = runInternal(queryListener, closer); + reportPayload = runInternal(reportPublishingListener, closer); } catch (Throwable e) { log.error(e, "Controller internal execution encountered exception."); - queryListener.onQueryComplete(makeStatusReportForException(e)); + reportPublishingListener.onQueryComplete(makeStatusReportForException(e)); throw e; } // Call onQueryComplete after Closer is fully closed, ensuring no controller-related processing is ongoing. - queryListener.onQueryComplete(reportPayload); + reportPublishingListener.onQueryComplete(reportPayload); } @@ -1062,6 +1071,13 @@ public class ControllerImpl implements Controller ); } + @Override + @Nullable + public TaskReport.ReportMap finalReport() + { + return finalReport.get(); + } + @Override @Nullable public TaskReport.ReportMap liveReports() diff --git a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/ControllerHolderTest.java b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/ControllerHolderTest.java index 1b4dfcbcc90..38e53adbf09 100644 --- a/multi-stage-query/src/test/java/org/apache/druid/msq/exec/ControllerHolderTest.java +++ b/multi-stage-query/src/test/java/org/apache/druid/msq/exec/ControllerHolderTest.java @@ -384,6 +384,7 @@ public class ControllerHolderTest public void testFinalReportVisibleBeforeCompletionListenerReturns() throws Exception { final MSQTaskReportPayload finalPayload = makeSuccessReport(); + final AtomicReference<TaskReport.ReportMap> finalReports = new AtomicReference<>(); final CountDownLatch listenerCalled = new CountDownLatch(1); final CountDownLatch releaseListener = new CountDownLatch(1); final DartControllerRegistry registry = new DartControllerRegistry(new DartControllerConfig()); @@ -392,8 +393,15 @@ public class ControllerHolderTest @Override public void run(final QueryListener listener) { + finalReports.set(TaskReport.buildTaskReports(new MSQTaskReport("test-query", finalPayload))); listener.onQueryComplete(finalPayload); } + + @Override + public TaskReport.ReportMap finalReport() + { + return finalReports.get(); + } }; final ControllerHolder holder = new ControllerHolder( controller, @@ -569,6 +577,12 @@ public class ControllerHolderTest return new TaskReport.ReportMap(); } + @Override + public TaskReport.ReportMap finalReport() + { + return TaskReport.buildTaskReports(new MSQTaskReport(queryId, makeSuccessReport())); + } + @Override public ControllerContext getControllerContext() { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
