2018-05-23 20:19:33 +00:00
|
|
|
#pragma once
|
|
|
|
|
2018-05-24 02:39:22 +00:00
|
|
|
#include <vector>
|
2018-05-23 20:19:33 +00:00
|
|
|
#include <set>
|
|
|
|
#include <mutex>
|
|
|
|
#include <Processors/IProcessor.h>
|
|
|
|
|
2019-02-05 13:01:40 +00:00
|
|
|
template <typename>
|
|
|
|
class ThreadPoolImpl;
|
|
|
|
class ThreadFromGlobalPool;
|
|
|
|
using ThreadPool = ThreadPoolImpl<ThreadFromGlobalPool>;
|
2018-05-23 20:19:33 +00:00
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
/** Wraps pipeline in a single processor.
|
|
|
|
* This processor has no inputs and outputs and just executes the pipeline,
|
|
|
|
* performing all synchronous work within a threadpool.
|
|
|
|
*/
|
|
|
|
class ParallelPipelineExecutor : public IProcessor
|
|
|
|
{
|
|
|
|
private:
|
2018-05-24 02:39:22 +00:00
|
|
|
Processors processors;
|
2018-05-23 20:19:33 +00:00
|
|
|
ThreadPool & pool;
|
|
|
|
|
|
|
|
std::set<IProcessor *> active_processors;
|
|
|
|
std::mutex mutex;
|
|
|
|
|
|
|
|
IProcessor * current_processor = nullptr;
|
|
|
|
Status current_status;
|
|
|
|
|
|
|
|
public:
|
2018-05-24 02:39:22 +00:00
|
|
|
ParallelPipelineExecutor(const Processors & processors, ThreadPool & pool);
|
2018-05-23 20:19:33 +00:00
|
|
|
|
|
|
|
String getName() const override { return "ParallelPipelineExecutor"; }
|
|
|
|
|
|
|
|
Status prepare() override;
|
|
|
|
void schedule(EventCounter & watch) override;
|
|
|
|
};
|
|
|
|
|
|
|
|
}
|