From a39a060752c55f8ab17a92512c38c21998ab9550 Mon Sep 17 00:00:00 2001 From: nanjeshramesh Date: Thu, 8 Oct 2026 13:00:35 -0700 Subject: [PATCH] chore: restore query ID verification in MSQ workers --- .../org/apache/druid/msq/exec/WorkerImpl.java | 17 +++++++++++++---- .../org/apache/druid/msq/kernel/StageId.java | 3 --- 2 files changed, 13 insertions(+), 7 deletions(-) diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/WorkerImpl.java b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/WorkerImpl.java index fa7d5c232ce6..6b2859a20915 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/exec/WorkerImpl.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/exec/WorkerImpl.java @@ -1135,7 +1135,7 @@ public static KernelHolders create(final WorkerContext workerContext, final Clos */ public void addKernel(final WorkerStageKernel kernel) { - final StageId stageId = kernel.getWorkOrder().getStageDefinition().getId(); + final StageId stageId = verifyQueryId(kernel.getWorkOrder().getStageDefinition().getId()); if (holderMap.putIfAbsent(stageId.getStageNumber(), new KernelHolder(kernel)) != null) { // Already added. Do nothing. @@ -1151,7 +1151,7 @@ public void addKernel(final WorkerStageKernel kernel) */ public void finishProcessing(final StageId stageId) { - final KernelHolder kernel = holderMap.get(stageId.getStageNumber()); + final KernelHolder kernel = holderMap.get(verifyQueryId(stageId).getStageNumber()); if (kernel != null) { try { @@ -1172,7 +1172,7 @@ public void finishProcessing(final StageId stageId) */ public void removeKernel(final StageId stageId) { - final KernelHolder removed = holderMap.remove(stageId.getStageNumber()); + final KernelHolder removed = holderMap.remove(verifyQueryId(stageId).getStageNumber()); if (removed == null) { throw new ISE("No kernel for stage[%s]", stageId); @@ -1226,7 +1226,7 @@ public int runningKernelCount() @Nullable public WorkerStageKernel getKernelFor(final StageId stageId) { - final KernelHolder holder = holderMap.get(stageId.getStageNumber()); + final KernelHolder holder = holderMap.get(verifyQueryId(stageId).getStageNumber()); if (holder != null) { return holder.kernel; } else { @@ -1275,6 +1275,15 @@ public void setDone() { this.done = true; } + + private StageId verifyQueryId(final StageId stageId) + { + if (!stageId.getQueryId().equals(workerContext.queryId())) { + throw new ISE("Unexpected queryId[%s], expected queryId[%s]", stageId.getQueryId(), workerContext.queryId()); + } + + return stageId; + } } /** diff --git a/multi-stage-query/src/main/java/org/apache/druid/msq/kernel/StageId.java b/multi-stage-query/src/main/java/org/apache/druid/msq/kernel/StageId.java index a928e5834fc4..5b98eed0da95 100644 --- a/multi-stage-query/src/main/java/org/apache/druid/msq/kernel/StageId.java +++ b/multi-stage-query/src/main/java/org/apache/druid/msq/kernel/StageId.java @@ -31,9 +31,6 @@ /** * Globally unique stage identifier: query ID plus stage number. - * - * Note: Versions till Druid 30 had a bug in the QueryKits which populated the {@link #queryId} field with random - * UUIDs. Therefore, all usage of the field must be vetted instead of assuming that it will be the expected query id */ public class StageId implements Comparable {