Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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 {
Expand All @@ -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);
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
}
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<StageId>
{
Expand Down
Loading