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]

Reply via email to