Skip to content
Draft
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
5 changes: 5 additions & 0 deletions be/src/pipeline/pipeline_task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,11 @@ void PipelineTask::set_task_queue(TaskQueue* task_queue) {
_task_queue = task_queue;
}

std::shared_ptr<PipelineTask> PipelineTask::get_task_holder() {
auto context_holder = _fragment_context->shared_from_this();
return std::shared_ptr<PipelineTask>(context_holder, this);
}

Status PipelineTask::execute(bool* eos) {
SCOPED_TIMER(_task_profile->total_time_counter());
SCOPED_TIMER(_exec_timer);
Expand Down
7 changes: 6 additions & 1 deletion be/src/pipeline/pipeline_task.h
Original file line number Diff line number Diff line change
Expand Up @@ -259,7 +259,12 @@ class PipelineTask {
virtual bool is_pipelineX() const { return false; }

bool is_running() { return _running.load(); }
void set_running(bool running) { _running = running; }
// Return the previous state so a scheduler can atomically claim the task.
bool set_running(bool running) { return _running.exchange(running); }

// Pipeline tasks are owned by their fragment context. The aliasing shared pointer keeps the
// context (and therefore this task) alive while the task is waiting in an asynchronous queue.
std::shared_ptr<PipelineTask> get_task_holder();

bool is_exceed_debug_timeout() {
if (_has_exceed_timeout) {
Expand Down
13 changes: 8 additions & 5 deletions be/src/pipeline/pipeline_x/dependency.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -44,17 +44,18 @@ Dependency* BasicSharedState::create_sink_dependency(int dest_id, int node_id, s
}

void Dependency::_add_block_task(PipelineXTask* task) {
DCHECK(_blocked_task.empty() || _blocked_task[_blocked_task.size() - 1] != task)
auto task_holder = std::static_pointer_cast<PipelineXTask>(task->get_task_holder());
DCHECK(_blocked_task.empty() || _blocked_task.back().lock().get() != task)
<< "Duplicate task: " << task->debug_string();
_blocked_task.push_back(task);
_blocked_task.emplace_back(task_holder);
}

void Dependency::set_ready() {
if (_ready) {
return;
}
_watcher.stop();
std::vector<PipelineXTask*> local_block_task {};
std::vector<std::weak_ptr<PipelineXTask>> local_block_task {};
{
std::unique_lock<std::mutex> lc(_task_lock);
if (_ready) {
Expand All @@ -63,8 +64,10 @@ void Dependency::set_ready() {
_ready = true;
local_block_task.swap(_blocked_task);
}
for (auto* task : local_block_task) {
task->wake_up();
for (auto& task : local_block_task) {
if (auto task_holder = task.lock()) {
task_holder->wake_up();
}
}
}

Expand Down
2 changes: 1 addition & 1 deletion be/src/pipeline/pipeline_x/dependency.h
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ class Dependency : public std::enable_shared_from_this<Dependency> {
MonotonicStopWatch _watcher;

std::mutex _task_lock;
std::vector<PipelineXTask*> _blocked_task;
std::vector<std::weak_ptr<PipelineXTask>> _blocked_task;

// If `_always_ready` is true, `block()` will never block tasks.
std::atomic<bool> _always_ready = false;
Expand Down
3 changes: 3 additions & 0 deletions be/src/pipeline/pipeline_x/pipeline_x_task.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -547,6 +547,9 @@ std::string PipelineXTask::debug_string() {

void PipelineXTask::wake_up() {
// call by dependency
if (is_finished()) {
return;
}
static_cast<void>(get_task_queue()->push_back(this));
}
} // namespace doris::pipeline
54 changes: 36 additions & 18 deletions be/src/pipeline/task_queue.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,11 +29,11 @@ namespace pipeline {

TaskQueue::~TaskQueue() = default;

PipelineTask* SubTaskQueue::try_take(bool is_steal) {
PipelineTaskSPtr SubTaskQueue::try_take(bool is_steal) {
if (_queue.empty()) {
return nullptr;
}
auto task = _queue.front();
auto task = std::move(_queue.front());
_queue.pop();
return task;
}
Expand All @@ -49,13 +49,30 @@ PriorityTaskQueue::PriorityTaskQueue() : _closed(false) {
}

void PriorityTaskQueue::close() {
std::unique_lock<std::mutex> lock(_work_size_mutex);
_closed = true;
_wait_task.notify_all();
DorisMetrics::instance()->pipeline_task_queue_size->increment(-_total_task_size);
std::vector<PipelineTaskSPtr> pending_tasks;
{
std::unique_lock<std::mutex> lock(_work_size_mutex);
if (_closed) {
return;
}
_closed = true;
_wait_task.notify_all();
const auto pending_task_size = _total_task_size.exchange(0);
DorisMetrics::instance()->pipeline_task_queue_size->increment(
-static_cast<int64_t>(pending_task_size));
pending_tasks.reserve(pending_task_size);
for (auto& queue : _sub_queues) {
queue.drain(&pending_tasks);
}
}
// Releasing a task may destroy its fragment context, so do it outside the queue lock.
for (const auto& task : pending_tasks) {
task->pop_out_runnable_queue();
}
pending_tasks.clear();
}

PipelineTask* PriorityTaskQueue::_try_take_unprotected(bool is_steal) {
PipelineTaskSPtr PriorityTaskQueue::_try_take_unprotected(bool is_steal) {
if (_total_task_size == 0 || _closed) {
return nullptr;
}
Expand Down Expand Up @@ -92,13 +109,13 @@ int PriorityTaskQueue::_compute_level(uint64_t runtime) {
return SUB_QUEUE_LEVEL - 1;
}

PipelineTask* PriorityTaskQueue::try_take(bool is_steal) {
PipelineTaskSPtr PriorityTaskQueue::try_take(bool is_steal) {
// TODO other efficient lock? e.g. if get lock fail, return null_ptr
std::unique_lock<std::mutex> lock(_work_size_mutex);
return _try_take_unprotected(is_steal);
}

PipelineTask* PriorityTaskQueue::take(uint32_t timeout_ms) {
PipelineTaskSPtr PriorityTaskQueue::take(uint32_t timeout_ms) {
std::unique_lock<std::mutex> lock(_work_size_mutex);
auto task = _try_take_unprotected(false);
if (task) {
Expand All @@ -113,20 +130,21 @@ PipelineTask* PriorityTaskQueue::take(uint32_t timeout_ms) {
}
}

Status PriorityTaskQueue::push(PipelineTask* task) {
Status PriorityTaskQueue::push(PipelineTaskSPtr task) {
auto level = _compute_level(task->get_runtime_ns());
std::unique_lock<std::mutex> lock(_work_size_mutex);
if (_closed) {
return Status::InternalError("WorkTaskQueue closed");
}
auto level = _compute_level(task->get_runtime_ns());
std::unique_lock<std::mutex> lock(_work_size_mutex);

// update empty queue's runtime, to avoid too high priority
if (_sub_queues[level].empty() &&
_queue_level_min_vruntime > _sub_queues[level].get_vruntime()) {
_sub_queues[level].adjust_runtime(_queue_level_min_vruntime);
}

_sub_queues[level].push_back(task);
task->put_in_runnable_queue();
_sub_queues[level].push_back(std::move(task));
_total_task_size++;
DorisMetrics::instance()->pipeline_task_queue_size->increment(1);
_wait_task.notify_one();
Expand All @@ -148,8 +166,8 @@ void MultiCoreTaskQueue::close() {
[](auto& prio_task_queue) { prio_task_queue.close(); });
}

PipelineTask* MultiCoreTaskQueue::take(int core_id) {
PipelineTask* task = nullptr;
PipelineTaskSPtr MultiCoreTaskQueue::take(int core_id) {
PipelineTaskSPtr task = nullptr;
while (!_closed) {
DCHECK(_prio_task_queue_list.size() > core_id)
<< " list size: " << _prio_task_queue_list.size() << " core_id: " << core_id
Expand All @@ -175,7 +193,7 @@ PipelineTask* MultiCoreTaskQueue::take(int core_id) {
return task;
}

PipelineTask* MultiCoreTaskQueue::_steal_take(int core_id) {
PipelineTaskSPtr MultiCoreTaskQueue::_steal_take(int core_id) {
DCHECK(core_id < _core_size);
int next_id = core_id;
for (int i = 1; i < _core_size; ++i) {
Expand Down Expand Up @@ -203,8 +221,8 @@ Status MultiCoreTaskQueue::push_back(PipelineTask* task) {

Status MultiCoreTaskQueue::push_back(PipelineTask* task, int core_id) {
DCHECK(core_id < _core_size);
task->put_in_runnable_queue();
return _prio_task_queue_list[core_id].push(task);
auto task_holder = task->get_task_holder();
return _prio_task_queue_list[core_id].push(std::move(task_holder));
}

} // namespace pipeline
Expand Down
31 changes: 21 additions & 10 deletions be/src/pipeline/task_queue.h
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@
#include <ostream>
#include <queue>
#include <set>
#include <utility>
#include <vector>

#include "common/status.h"
#include "pipeline_task.h"
Expand All @@ -35,14 +37,16 @@
namespace doris {
namespace pipeline {

using PipelineTaskSPtr = std::shared_ptr<PipelineTask>;

class TaskQueue {
public:
TaskQueue(int core_size) : _core_size(core_size) {}
virtual ~TaskQueue();
virtual void close() = 0;
// Get the task by core id.
// TODO: To think the logic is useful?
virtual PipelineTask* take(int core_id) = 0;
virtual PipelineTaskSPtr take(int core_id) = 0;

// push from scheduler
virtual Status push_back(PipelineTask* task) = 0;
Expand All @@ -63,9 +67,16 @@ class SubTaskQueue {
friend class PriorityTaskQueue;

public:
void push_back(PipelineTask* task) { _queue.emplace(task); }
void push_back(PipelineTaskSPtr task) { _queue.emplace(std::move(task)); }

PipelineTaskSPtr try_take(bool is_steal);

PipelineTask* try_take(bool is_steal);
void drain(std::vector<PipelineTaskSPtr>* tasks) {
while (!_queue.empty()) {
tasks->emplace_back(std::move(_queue.front()));
_queue.pop();
}
}

void set_level_factor(double level_factor) { _level_factor = level_factor; }

Expand All @@ -81,7 +92,7 @@ class SubTaskQueue {
bool empty() { return _queue.empty(); }

private:
std::queue<PipelineTask*> _queue;
std::queue<PipelineTaskSPtr> _queue;
// depends on LEVEL_QUEUE_TIME_FACTOR
double _level_factor = 1;

Expand All @@ -95,18 +106,18 @@ class PriorityTaskQueue {

void close();

PipelineTask* try_take(bool is_steal);
PipelineTaskSPtr try_take(bool is_steal);

PipelineTask* take(uint32_t timeout_ms = 0);
PipelineTaskSPtr take(uint32_t timeout_ms = 0);

Status push(PipelineTask* task);
Status push(PipelineTaskSPtr task);

void inc_sub_queue_runtime(int level, uint64_t runtime) {
_sub_queues[level].inc_runtime(runtime);
}

private:
PipelineTask* _try_take_unprotected(bool is_steal);
PipelineTaskSPtr _try_take_unprotected(bool is_steal);
static constexpr auto LEVEL_QUEUE_TIME_FACTOR = 2;
static constexpr size_t SUB_QUEUE_LEVEL = 6;
SubTaskQueue _sub_queues[SUB_QUEUE_LEVEL];
Expand Down Expand Up @@ -135,7 +146,7 @@ class MultiCoreTaskQueue : public TaskQueue {
void close() override;

// Get the task by core id.
PipelineTask* take(int core_id) override;
PipelineTaskSPtr take(int core_id) override;

// TODO combine these methods to `push_back(task, core_id = -1)`
Status push_back(PipelineTask* task) override;
Expand All @@ -149,7 +160,7 @@ class MultiCoreTaskQueue : public TaskQueue {
}

private:
PipelineTask* _steal_take(int core_id);
PipelineTaskSPtr _steal_take(int core_id);

std::vector<PriorityTaskQueue> _prio_task_queue_list;
std::atomic<uint32_t> _next_core = 0;
Expand Down
15 changes: 11 additions & 4 deletions be/src/pipeline/task_scheduler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -264,16 +264,23 @@ void _close_task(PipelineTask* task, PipelineTaskState state, Status exec_status
void TaskScheduler::_do_work(size_t index) {
const auto& marker = _markers[index];
while (*marker) {
auto* task = _task_queue->take(index);
if (!task) {
auto task_holder = _task_queue->take(index);
if (!task_holder) {
continue;
}
if (task->is_pipelineX() && task->is_running()) {
auto* task = task_holder.get();
if (task->is_pipelineX() && task->set_running(true)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Claim or coalesce duplicates before touching task accounting

The atomic claim happens only after take() has updated this task's _queue_level/_core_id and stopped its wait-worker watcher. When this branch loses, it immediately pushes the same token again, so an idle worker can hot-loop pop/push for the owner's full execution slice (or longer close) while push() reads _runtime, increments _schedule_time, and restarts the same plain stopwatch. Meanwhile the owner calls update_statistics(), which writes _runtime and reads _core_id/_queue_level; per-core queue locks do not synchronize these task fields, especially after steals. This leaves C++ data races plus corrupted priority/profile accounting even though execute() is serialized. Please claim/coalesce before per-task dequeue accounting and add a latched multi-worker duplicate test.

static_cast<void>(_task_queue->push_back(task, index));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Do not recycle a wake token across task-state epochs

Requeuing the losing entry preserves it without tying it to the dependency epoch that created it. A delayed duplicate can sit in the queue while the owner runs; if the owner then reaches EOS, is_pending_finish() registers an unready finish dependency and the owner publishes PENDING_FINISH before clearing _running. This stale entry can now claim the task, and the PENDING_FINISH branch calls _close_task() without rechecking finish dependencies for PipelineX, releasing async-writer/exchange state before the callback makes the dependency ready. If the owner instead enters an ordinary blocked state, the stale entry re-enters the same unready dependency and trips _add_block_task's duplicate DCHECK (or appends duplicate waiters in release builds). Please coalesce or generation-tag pending wakes so an old token cannot satisfy a later wait epoch, and add a latched regression covering both transitions.

continue;
}
if (task->is_finished()) {
task->set_running(false);
continue;
}
task->log_detail_if_need();
task->set_running(true);
if (!task->is_pipelineX()) {
task->set_running(true);
}
task->set_task_queue(_task_queue.get());
auto* fragment_ctx = task->fragment_context();
bool canceled = fragment_ctx->is_canceled();
Expand Down
Loading
Loading