You are viewing an old version of this page. View the current version.

Compare with Current View Page History

« Previous Version 5 Next »

1 Motivation


pipeline/向量化引擎上存在的一些问题:

  1. 执行并发上,当前Doris执行并发收到两个因素的制约,一个是fe设置的参数,另一个是受存储层bucket数量的限制,这样的静态并发使得执行引擎无法充分利用机器资源。

  2. 执行逻辑上,当前Doris有一些固定的额外开销,例如表达式部分各个instance彼此独立,而instance的初始化参数有很多公共部分,所以需要额外进行很多重复的初始化步骤。

  3. 调度逻辑上,当前pipeline的调度器会把阻塞task全部放入一个阻塞队列中,由一个线程负责轮询并从阻塞队列中取出可执行task放入runnable队列,所以在有查询执行的过程中,会固定有一个核的资源作为调度的开销。

  4. profile方面,目前pipeline无法为用户提供简单易懂的指标。

2 Goals

  1. 执行并发上,依赖local exchange使pipelinex充分并发,可以让数据被均匀分布到不同的task中,尽可能减少数据倾斜,此外,pipelineX也将不再受存储层tablet数量的制约。

  2. 执行逻辑上,多个pipeline task共享同一个pipeline的全部共享状态,例如表达式和一些const变量,消除了额外的初始化开销。

  3. 调度逻辑上,所有pipeline task的阻塞条件都使用Dependency进行了封装,通过外部事件(例如rpc完成)触发task的执行逻辑进入runnable队列,从而消除了阻塞轮询线程的开销。

  4. profile:为用户提供简单易懂的指标。

3 详细设计

3.1 核心数据结构抽象

3.1.1 继承关系



图1 operator继承关系

图2 LocalState继承关系

3.1.2 组合关系



图3 核心类组合关系

3.2 执行模型

3.2.1 聚合查询


在2BE的集群上运行,执行模型如下:

图4 多BE执行模型(聚合查询)


核心改造:

  1. 执行线程(thread 1, thread 2)执行各自的pipeline task,而pipeline task仅持有一些运行时状态(即local state)。全局信息则由多个task共享的同一个pipeline对象持有(即global state)

  2. 数据分发在单个be上由local shuffle节点完成,由local shuffle来保证多个pipeline task之间的数据均衡。单be上执行模型如下:



图5 单BE执行模型(聚合查询)

3.2.2 带join的查询



图6 单BE执行模型(join查询)

3.2.3 预期收益


引入local shuffle主要是为了解决单机的并发能力,

  1. 可以减少部分情况下的数据倾斜 (详细设计见文末参考文档1。)

  2. 执行并发度不再受存储层tablet数量的制约 (详细设计见文末参考文档4。)

  3. 可以在运行时进行动态并发

3.3 执行流程


pipeline和pipelinex的执行流程对比如下:



图7 pipeline/pipelinex执行流程对比


prepare阶段,pipeline会并发启动多个线程去做不同instance的状态初始化,而由于pipelinex对共享的状态做了复用,也就是把pipeline执行流程中的第3步拆分成了pipelineX执行流程中的第3步和第5步,对比较重的global state只做一次,对更轻量的local state进行串行初始化。

3.3.1 预期收益


执行流程的改造主要是为了降低了初始化的额外开销,

  1. pipeline会启动多个线程同时对多个instance进行初始化的开销

  2. 全局const变量初始化在多个instance中重复初始化的开销

3.4 调度模型


pipeline调度中,就绪task保存在就绪队列中等待调度,阻塞task保存在阻塞队列中等待满足执行条件。而在Doris pipeline当前额外需要一个CPU core去轮询阻塞队列,如果task满足执行条件则保存在就绪队列中。而在pipelineX中,阻塞条件都使用dependency进行了封装,task的阻塞/就绪全部依赖事件通知来解决(详细设计见文末参考文档2。)。例如rpc数据到达将会触发ExchangeSourceOperator满足执行条件进入就绪队列。

图8 pipeline/pipelinex调度模型对比

3.4.1 预期收益


消除了轮询线程的额外开销。

3.5 Profile改造


对于operator的profile,pipelineX做了整理,包括删除不合理的metrics并且添加了必要的metrics。除此之外,得益于调度模型改造,pipelineX中所有阻塞都被dependency封装,所以我们将所有dependency的就绪时间加入profile,通过wait for dependency时间我们可以直观看出每个地方的时间开销。
举几个例子。
Scan operator:

OLAP_SCAN_OPERATOR (id=504):(Active: 17.606ms, % non-child: 0.00%)

    - WaitForDependencyTime: 0ns

        - WaitForData: 14.311ms

        - WaitForEos: 0ns

        - WaitForScannerDone: 0ns

    - WaitForPendingFinishDependency: 16.748ms

其中,OLAP_SCAN_OPERATOR的active总时间是17.606ms(包括了等待scanner读数据的时间和执行的时间),其中因为等待scanner扫描数据阻塞了14.311ms,而OLAP_SCAN_OPERATOR因为PendingFinish阻塞了16.748ms(不包含在active时间中)
Exchange source operator:

EXCHANGE_OPERATOR (id=517):(Active: 575.846us, % non-child: 0.00%)

    - WaitForDependencyTime:

        - WaitForData: 56.122ms

    - WaitForPendingFinishDependency: 0ns

EXCHANGE_OPERATOR的Active时间为575.846us,等待上游数据的时间为56.122ms

4 User Interface


Add 3 SessionVariable knobs.
enable_pipeline_x_engine

enable_local_shuffle

ignore_storage_data_distribution
Ignore storage data distribution or not. If turn on, execution concurrency will not be restricted by tablet num. Please refer to Further Reading 4 for more details.

A new http api:http://{host}:{web_server_port}/api/running_pipeline_tasks
Please refer to Further Reading 3 for more details.

Further Reading

  1. LocalExchanger Rules

  2. PipelineX Dependency Details
  3. pipelineX Debug Manual
  4. PipelineX Parallel Execution Design
  • No labels