PhysicalOperator是duckdb所有物理算子的基类。
字段
class PhysicalOperator {
public:
//! 子算子列表
ArenaLinkedList<reference<PhysicalOperator>> children;
//! 标识算子的具体类型
PhysicalOperatorType type;
//! 算子返回结果的列类型列表
vector<LogicalType> types;
//! 优化器估算的算子输出行数,用于选择执行策略
idx_t estimated_cardinality;
//! sink算子的全局状态。用于汇总所有输入数据后才能产生输出的算子(如排序、聚合、插入)
unique_ptr<GlobalSinkState> sink_state;
//! 普通算子的全局状态。用于需要在并行执行中共享数据的算子(如哈希表构建)
unique_ptr<GlobalOperatorState> op_state;
//! 多线程保护sink_state, op_state的访问修改安全
mutex lock;
}duckdb是pipeline架构,其中的算子有三种角色
- Source:Pipeline的起点,负责产生数据。通常是SCAN等扫描类算子
- Sink:Pipeline的终点,负责汇集所有数据并产生最终结果
- Operator:处于Pipeline中间,负责处理数据。
这三种解决是执行阶段的概念,并非算子的静态分类属性,根据当前执行的阶段,一个算子可能处于不同的角色。下面按照角色分组看下有哪些接口
source interface
class PhysicalOperator {
public:
// Source interface
virtual unique_ptr<LocalSourceState> GetLocalSourceState(ExecutionContext &context,
GlobalSourceState &gstate) const;
virtual unique_ptr<GlobalSourceState> GetGlobalSourceState(ClientContext &context) const;
protected:
virtual SourceResultType GetDataInternal(ExecutionContext &context, DataChunk &chunk,
OperatorSourceInput &input) const;
public:
SourceResultType GetData(ExecutionContext &context, DataChunk &chunk, OperatorSourceInput &input) const;
virtual OperatorPartitionData GetPartitionData(ExecutionContext &context, DataChunk &chunk,
GlobalSourceState &gstate, LocalSourceState &lstate,
const OperatorPartitionInfo &partition_info) const;
virtual bool IsSource() const {
return false;
}
virtual bool ParallelSource() const {
return false;
}
virtual bool SupportsPartitioning(const OperatorPartitionInfo &partition_info) const {
if (partition_info.AnyRequired()) {
return false;
}
return true;
}
//! The type of order emitted by the operator (as a source)
virtual OrderPreservationType SourceOrder() const {
return OrderPreservationType::INSERTION_ORDER;
}
//! Returns the current progress percentage, or a negative value if progress bars are not supported
virtual ProgressData GetProgress(ClientContext &context, GlobalSourceState &gstate) const;
//! Returns the current progress percentage, or a negative value if progress bars are not supported
virtual ProgressData GetSinkProgress(ClientContext &context, GlobalSinkState &gstate,
const ProgressData source_progress) const {
return source_progress;
}
virtual InsertionOrderPreservingMap<string> ExtraSourceParams(GlobalSourceState &gstate,
LocalSourceState &lstate) const {
return InsertionOrderPreservingMap<string>();
}
}GetGlobalSourceState/GetLocalSourceState:创建全局源状态和局部源状态。前者是所有线程共享的状态,例如扫描文件时的全局文件句柄、全局索引偏移量等。后者时每个工作线程私有的状态,例如线程内部的缓冲区、本地游标等。duckdb采用的时morsel-driver parallelism,多个线程会并发调用GetData获取数据块,为了线程安全和减少锁竞争,将全局数据和线程局部数据分离成为两个状态。GetData/GetDataInternal:前者时对外的公共接口,内部调用GetDataInetrnal。这是source接口最核心的方法,负责将实际数据推给下游GetPartitionData:获取当前数据块的分区键信息,例如在并行哈希连接中,右表构建阶段需要将数据按照哈希值分区,让不同的线程处理不同分区。这个方法可以让source提前告知下游当前chunk应该路由到哪个分区。这个用于支持高级并行特性的。IsSource,ParallelSource,SupportsPartitioning,SourceOrder:这一组方法是用来声明能力的IsSource:返回该算子是否能够作为Source。在构建pipeline的时候用来决定是否能够作为pipeline的起点ParallelSource:该算子是否支持多线程并行生产数据。并不是所有数据源都支持安全地并行扫描,这种情况就只能单线程执行。SupportsPartitioning:该source是否支持按照给定的分区信息进行数据分区输出。上层算子如HashJoin构建的时候可能希望Source按照特定键分区,以减少数据重排开销SourceOrder:返回该Source输出的数据顺序保证类型。默认返回INSERTION_ORDER(即数据写入顺序)。优化器知道数据顺序,就可以决定是否省略后续的SORT算子了
GetProgress:返回当前Source的执行进度。GetSinkProgress:对于同时作为Sink和Source的算子,需要综合上游Sink进度和当前Source阶段进度来计算总进度。ExtraSourceParams:返回一个键值对映射,包含Source算子执行的额外参数。比如使用的索引名称、并行度等信息。
operator interface
class PhysicalOperator {
public:
// Operator interface
virtual unique_ptr<OperatorState> GetOperatorState(ExecutionContext &context) const;
virtual unique_ptr<GlobalOperatorState> GetGlobalOperatorState(ClientContext &context) const;
virtual OperatorResultType Execute(ExecutionContext &context, DataChunk &input, DataChunk &chunk,
GlobalOperatorState &gstate, OperatorState &state) const;
virtual OperatorFinalizeResultType FinalExecute(ExecutionContext &context, DataChunk &chunk,
GlobalOperatorState &gstate, OperatorState &state) const;
virtual OperatorFinalResultType OperatorFinalize(Pipeline &pipeline, Event &event, ClientContext &context,
OperatorFinalizeInput &input) const;
virtual bool ParallelOperator() const {
return false;
}
virtual bool RequiresFinalExecute() const {
return false;
}
virtual bool RequiresOperatorFinalize() const {
return false;
}
//! The influence the operator has on order (insertion order means no influence)
virtual OrderPreservationType OperatorOrder() const {
return OrderPreservationType::INSERTION_ORDER;
}
}GetOperatorState/GetGlobalOperatorState:局部算子状态和全局算子状态。和source接口的两个状态作用类似,不过是operator阶段用的Execute:核心执行函数,接收一个输入chunk,处理后输出一个输出chunk,并更新局部/全局状态FinalExecute:在所有数据都处理完毕后,允许算子产生最后一批输出。如Limit,Distinct这些算子在输入结束之后可能仍然有缓存的输出需要吐出,FinalExecute提供了清理并输出尾部数据的机会,避免丢失数据。OperatorFinalize:在整个pipeline执行结束后调用,用于执行全局汇总或资源释放ParallelOperator:标识该算子是否支持多线程并行执行RequiresFinalExecute:表示该算子是否需要调用FinalExecute。若算子不需要调度器就可以跳过这一阶段,减少不必要的开销。OperatorOrder:描述该算子对数据数据顺序的影响。和SourceOrder类似,同样让优化器决定是否可以消除后续的SORT操作。
sink interface
class PhysicalOperator {
public:
// Sink interface
//! The sink method is called constantly with new input, as long as new input is available. Note that this method
//! CAN be called in parallel, proper locking is needed when accessing dat //! a inside the GlobalSinkState.
virtual SinkResultType Sink(ExecutionContext &context, DataChunk &chunk, OperatorSinkInput &input) const;
//! The combine is called when a single thread has completed execution of its part of the pipeline, it is the final
//! time that a specific LocalSinkState is accessible. This method can be called in parallel while other Sink() or //! Combine() calls are active on the same GlobalSinkState.
virtual SinkCombineResultType Combine(ExecutionContext &context, OperatorSinkCombineInput &input) const;
//! (optional) function that will be called before Finalize
//! For now, its only use is to to communicate memory usage in multi-join pipelines through TemporaryMemoryManager
virtual void PrepareFinalize(ClientContext &context, GlobalSinkState &sink_state) const;
//! The finalize is called when ALL threads are finished execution. It is called only once per pipeline, and is
//! entirely single threaded. //! If Finalize returns SinkResultType::Finished, the sink is marked as finished
virtual SinkFinalizeType Finalize(Pipeline &pipeline, Event &event, ClientContext &context,
OperatorSinkFinalizeInput &input) const;
//! For sinks with RequiresBatchIndex set to true, when a new batch starts being processed this method is called
//! This allows flushing of the current batch (e.g. to disk)
virtual SinkNextBatchType NextBatch(ExecutionContext &context, OperatorSinkNextBatchInput &input) const;
virtual unique_ptr<LocalSinkState> GetLocalSinkState(ExecutionContext &context) const;
virtual unique_ptr<GlobalSinkState> GetGlobalSinkState(ClientContext &context) const;
//! The maximum amount of memory the operator should use per thread.
static idx_t GetMaxThreadMemory(ClientContext &context);
//! Whether operator caching is allowed in the current execution context
static OperatorCachingMode SelectOperatorCachingMode(ExecutionContext &context);
virtual bool IsSink() const {
return false;
}
virtual bool ParallelSink() const {
return false;
}
virtual OperatorPartitionInfo RequiredPartitionInfo() const {
return OperatorPartitionInfo::NoPartitionInfo();
}
//! Whether or not the sink operator depends on the order of the input chunks
//! If this is set to true, we cannot do things like caching intermediate vectors
virtual bool SinkOrderDependent() const {
return false;
}
}GetLocalSinkState/GetGlobalSinkState:与source和operator接口类似,不过是管理sink的局部和全局状态的。Sink:核心的数据接收函数。算是Sink算子消费数据的入口,将输入的chunk处理并合并到局部/全局状态当中。Combine:当一个线程完成其所有输入后调用,用于将该线程的局部状态合并到全局状态,可以并行调用。Combine提供了一个可控的合并点,允许算子高效地将局部哈希表或排序块合并到全局结构中,同时支持并发合并。PrepareFinalize:在Finalize执行调用的准备函数。某些复杂的Sink(例如多路连接)需要在最终处理前通过临时内存管理器当前的内存占用,方便整体协调内存限制。这个钩子函数提供了这样的时机。Finalize:最终处理函数。负责生成最终的输出结果,或者执行提交操作NextBatch:仅当Sink算子设置了RequireBatchIndex时被调用,在开始处理新的一批输入时调用,允许算子刷新当前批次的数据到磁盘等操作。当内存压力大时,NextBatch提供一个检查点,算子可以将当前批次的数据spill到磁盘,清空内存,继续下一批次处理。IsSink:该算子是否可以充当Sink角色。ParallelSink:表示该算子是否支持多线程并行调用Sink。RequiredPartitionInfo:返回该算子对输入数据的分区要求,比如哈希聚合可能希望输入按分组键哈希分区,以减少合并开销。优化器可以据此决定是否在Sink之前插入一个分区算子(Exchange),将数据按特定方式重分布,以提高Sink内部的合并效率。SinkOrderDependent:表示Sink算子的结果是否依赖于输入数据的顺序。某些算子对输入顺序敏感,该标志防止执行引擎进行可能改变顺序的优化。
pipeline interface
这是用于构建pipeline的接口
// Pipeline construction
virtual vector<const_reference<PhysicalOperator>> GetSources() const;
bool AllSourcesSupportBatchIndex() const;
virtual void BuildPipelines(Pipeline ¤t, MetaPipeline &meta_pipeline);GetSources:返回当前算子的源算子列表AllSourcesSupportBatchIndex:检查GetSources返回的源算子是否都支持批次索引BuildPipelines:递归构建Pipeline的核心方法。