This is an automated email from the ASF dual-hosted git repository.
cryptoe pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 04bd7f64126 minor: Add method Task.getDefaultPriority and coerce
priority value to integer (#19984)
04bd7f64126 is described below
commit 04bd7f641264efe2acb7f12053dcc6ecd9d25dfb
Author: Kashif Faraz <[email protected]>
AuthorDate: Wed Aug 12 20:16:00 2026 +0530
minor: Add method Task.getDefaultPriority and coerce priority value to
integer (#19984)
---
.../indexing/common/task/AbstractBatchIndexTask.java | 4 ++--
.../apache/druid/indexing/common/task/CompactionTask.java | 4 ++--
.../org/apache/druid/indexing/common/task/NoopTask.java | 4 ++--
.../java/org/apache/druid/indexing/common/task/Task.java | 15 ++++++++++++++-
.../indexing/seekablestream/SeekableStreamIndexTask.java | 4 ++--
.../apache/druid/indexing/common/task/NoopTaskTest.java | 10 ++++++++++
.../org/apache/druid/indexing/overlord/TaskQueueTest.java | 6 ++++--
.../org/apache/druid/msq/indexing/MSQControllerTask.java | 4 ++--
.../java/org/apache/druid/msq/indexing/MSQWorkerTask.java | 4 ++--
9 files changed, 40 insertions(+), 15 deletions(-)
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/AbstractBatchIndexTask.java
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/AbstractBatchIndexTask.java
index b9d22aca592..a555056c995 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/AbstractBatchIndexTask.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/AbstractBatchIndexTask.java
@@ -303,9 +303,9 @@ public abstract class AbstractBatchIndexTask extends
AbstractTask
public abstract Granularity getSegmentGranularity();
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY);
+ return Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY;
}
public TaskLockHelper getTaskLockHelper()
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/CompactionTask.java
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/CompactionTask.java
index 1955ff7e413..d4a8b68d26c 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/CompactionTask.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/CompactionTask.java
@@ -470,9 +470,9 @@ public class CompactionTask extends AbstractBatchIndexTask
implements PendingSeg
}
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_MERGE_TASK_PRIORITY);
+ return Tasks.DEFAULT_MERGE_TASK_PRIORITY;
}
@Override
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/NoopTask.java
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/NoopTask.java
index 99c71562e3c..ed81857a129 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/NoopTask.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/NoopTask.java
@@ -121,9 +121,9 @@ public class NoopTask extends AbstractTask implements
PendingSegmentAllocatingTa
}
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY);
+ return Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY;
}
@Override
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/Task.java
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/Task.java
index a9884438d5b..662c25ae212 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/common/task/Task.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/common/task/Task.java
@@ -42,6 +42,7 @@ import
org.apache.druid.indexing.common.task.batch.parallel.SinglePhaseSubTask;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.UOE;
import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryContexts;
import org.apache.druid.query.QueryRunner;
import org.apache.druid.server.coordination.BroadcastDatasourceLoadingSpec;
import org.apache.druid.server.lookup.cache.LookupLoadingSpec;
@@ -118,7 +119,19 @@ public interface Task
*/
default int getPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY, Tasks.DEFAULT_TASK_PRIORITY);
+ return QueryContexts.getAsInt(
+ Tasks.PRIORITY_KEY,
+ getContextValue(Tasks.PRIORITY_KEY),
+ getDefaultPriority()
+ );
+ }
+
+ /**
+ * Default value for {@link Tasks#PRIORITY_KEY} if not specified in the task
context.
+ */
+ default int getDefaultPriority()
+ {
+ return Tasks.DEFAULT_TASK_PRIORITY;
}
/**
diff --git
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTask.java
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTask.java
index 816cfc1574d..08813e61133 100644
---
a/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTask.java
+++
b/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/SeekableStreamIndexTask.java
@@ -123,9 +123,9 @@ public abstract class
SeekableStreamIndexTask<PartitionIdType, SequenceOffsetTyp
}
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_REALTIME_TASK_PRIORITY);
+ return Tasks.DEFAULT_REALTIME_TASK_PRIORITY;
}
@Override
diff --git
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/NoopTaskTest.java
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/NoopTaskTest.java
index 96bdef75f0f..978188c6f9d 100644
---
a/indexing-service/src/test/java/org/apache/druid/indexing/common/task/NoopTaskTest.java
+++
b/indexing-service/src/test/java/org/apache/druid/indexing/common/task/NoopTaskTest.java
@@ -22,6 +22,8 @@ package org.apache.druid.indexing.common.task;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import java.util.Map;
+
public class NoopTaskTest
{
@Test
@@ -30,4 +32,12 @@ public class NoopTaskTest
NoopTask task = NoopTask.create();
Assertions.assertTrue(task.getInputSourceResources().isEmpty());
}
+
+ @Test
+ public void test_getPriority_coercesStringValuesToInteger()
+ {
+ final NoopTask task = new NoopTask(null, null, null, 0, 0,
Map.of(Tasks.PRIORITY_KEY, "1000"));
+ Assertions.assertEquals(1000, task.getPriority());
+ Assertions.assertEquals(Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY,
task.getDefaultPriority());
+ }
}
diff --git
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueTest.java
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueTest.java
index d64b4003934..c09f623ede0 100644
---
a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueTest.java
+++
b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueTest.java
@@ -751,7 +751,7 @@ public class TaskQueueTest extends IngestionTestBase
}
@Test
- public void testTaskSubmissionToTaskRunnerBasedOnPriority() throws Exception
+ public void testTaskSubmissionToTaskRunnerBasedOnPriority()
{
final RecordingTaskRunner recordingRunner = new
RecordingTaskRunner(serviceEmitter);
final TaskQueue priorityQueue = new TaskQueue(
@@ -772,7 +772,9 @@ public class TaskQueueTest extends IngestionTestBase
final NoopTask lowPriority2 = NoopTask.ofPriority(10);
final NoopTask medPriority = NoopTask.ofPriority(50);
final NoopTask highPriority1 = NoopTask.ofPriority(90);
- final NoopTask highPriority2 = NoopTask.ofPriority(100);
+
+ // Create a task with a String priority value to verify that it gets
coerced as an integer
+ final NoopTask highPriority2 = new NoopTask(null, null, null, 0, 0,
Map.of(Tasks.PRIORITY_KEY, "100"));
priorityQueue.add(lowPriority1);
priorityQueue.add(medPriority);
diff --git
a/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQControllerTask.java
b/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQControllerTask.java
index b452ad1bb53..33648426445 100644
---
a/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQControllerTask.java
+++
b/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQControllerTask.java
@@ -322,9 +322,9 @@ public class MSQControllerTask extends AbstractTask
implements ClientTaskQuery,
}
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY);
+ return Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY;
}
@Nullable
diff --git
a/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQWorkerTask.java
b/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQWorkerTask.java
index a6517d88dc3..01935943c1a 100644
---
a/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQWorkerTask.java
+++
b/multi-stage-query/src/main/java/org/apache/druid/msq/indexing/MSQWorkerTask.java
@@ -175,9 +175,9 @@ public class MSQWorkerTask extends AbstractTask
}
@Override
- public int getPriority()
+ public int getDefaultPriority()
{
- return getContextValue(Tasks.PRIORITY_KEY,
Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY);
+ return Tasks.DEFAULT_BATCH_INDEX_TASK_PRIORITY;
}
@Override
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]