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应该路由到哪个分区。这个用于支持高级并行特性的。
  • IsSourceParallelSource, 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 &current, MetaPipeline &meta_pipeline);
  • GetSources:返回当前算子的源算子列表
  • AllSourcesSupportBatchIndex:检查GetSources返回的源算子是否都支持批次索引
  • BuildPipelines:递归构建Pipeline的核心方法。