首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >​工业分析实战: 分区并行与分布式聚合

​工业分析实战: 分区并行与分布式聚合

原创
作者头像
User_芊芊君子
发布2026-08-21 10:35:06
发布2026-08-21 10:35:06
410
举报

摘要

做工业 IoT 平台的人,大概都碰过这样的场景:设备数据攒了半年,某天领导要看一份全量的季度统计——每台设备每天的特征指标、趋势、异常天数。SQL 写好了,一跑,40 分钟。加节点?加内存?都不便宜。但很多时候,慢的根因不是算力不够,而是这条查询压根没吃满集群的并行能力——6 个节点的集群,查询实际只调动了其中一小部分资源,剩下的在旁边看戏。

这是典型的"单机思维写分布式 SQL":语法没错、结果没错,但执行方式是串行的。分布式数据库的性能上限,不取决于你写了多复杂的函数,而取决于查询被拆成了多少个可以同时跑的任务

本文从一个 40 分钟的季度统计出发,拆解 DolphinDB 批量查询的并行执行机制:分区怎么变成并行任务、聚合怎么在分区上两阶段合并、**context by** 的并行边界在哪里、什么样的写法会悄悄杀死并行,最后给出一份"让 SQL 吃满集群"的实践清单。

img
img

一、慢在哪:一个 40 分钟的季度统计

先把问题摆出来。场景:振动测点数据攒了半年,约 20 亿行,要做一份季度统计——每台设备每天算一组特征指标(均值、峰值、超标时长):

代码语言:SQL
复制
select
    date(ts) as day, deviceId,
    avg(rms)              as rms_mean,
    max(rms)              as rms_peak,
    sum(iif(rms > 5.0, 1, 0)) \ count(*) as over_ratio
from loadTable("dfs://iot", "vibration")
where ts between 2026.01.01 00:00:00 : 2026.03.31 23:59:59
group by date(ts), deviceId

SQL 很朴素。在测试集群(6 个数据节点)上跑,40 分钟。第一反应是数据量大——20 亿行确实不小。但看了执行情况才发现问题:这条查询只被拆成了 6 个子任务,每个节点一个,每个任务扛着几亿行在串行扫描。集群另外五个节点的 CPU 大部分时间在空转。

换句话说:20 亿行的计算量被切成了 6 大块,每块内部是单线程硬啃。并行度 = 6,再多的节点也帮不上忙。

img
img

问题的根源不在 SQL 写法,在分区设计——这张表当时是按月分区的,一个季度只有 3 个月,跨 3 个月再乘设备的 HASH 子分区,能调度的任务数就是那么几个。分区是并行的单位,分区数不够,并行度就上不去。

二、分区是并行调度的单位

这是理解 DolphinDB(以及大多数 MPP 风格系统)批量查询性能的第一原则:查询并行度的天花板,由命中的分区数决定。

2.1 执行模型:一条查询怎么被拆开

当一条批量查询命中多个分区时,执行过程大致是:

代码语言:Plaintext
复制
查询: select ... from 大表 where 时间范围 group by ...
        │
        ▼
1. 分区裁剪 → 根据过滤条件筛掉无关分区
        │
        ▼
2. 任务切分 → 每个命中的分区变成一个独立的子任务
        │
        ▼
3. 并行调度 → 子任务被分发到各节点的执行器上同时跑
        │
        ▼
4. 结果合并 → 各分区的局部结果汇总成最终结果
img
img

关键在第二步:分区是切分粒度。一个分区内部的扫描和计算,对一个子任务来说基本是串行的。所以:

  • 命中 6 个分区 → 最多 6 路并行
  • 命中 180 个分区 → 最多 180 路并行(受集群总核数限制)

2.2 分区数的两头陷阱

分区数决定并行上限,但不是越多越好。两头都有陷阱:

陷阱一:分区太大,并行度低

按月分区存高频测点,一个月动辄几亿行。查一个季度只命中 3 个月 × 设备 HASH 桶数个分区,任务数少,单任务数据量巨大——这就是第一节 40 分钟的根因。

陷阱二:分区太小,调度开销吃掉收益

反过来按小时分区,一天 24 个分区,一年 8000+。查一个月命中 720 个分区,看着并行度很高,但每个分区只有几万行,任务创建、调度、结果合并的固定开销占比急剧上升,还可能触发系统对单查询分区数的保护性限制。最终快是快了,但离线性加速差得远,且集群并发能力被一条大查询占满。

工程上的平衡点:让"常规查询命中的分区数"和"集群能同时调度的任务数"在同一量级,且单分区数据量在可接受的扫描粒度(经验上,单分区几百万到几千万行是个常见区间)。高频测点按天分区、低频数据按月分区,多数场景能落在合理区间——具体粒度要结合自己的查询模式压一压。

img
img

分区设计是同一枚硬币的两面:治理视角关心分区对 TTL、降采样的影响,并行视角关心分区对任务切分的影响。建表前两边都要想。

2.3 回到那个 40 分钟的查询

把测点表从按月分区改为按天分区后,同样的季度统计:

代码语言:Plaintext
复制
按月分区: 命中 3 个月 × N 个 HASH 桶   → 任务数少, 单任务几亿行  → 40 分钟
按天分区: 命中 90 天 × N 个 HASH 桶    → 任务数 ×30, 单任务千万级 → 几分钟

SQL 一个字没改,只是分区粒度变了,查询就能吃满集群。这是分布式调优里性价比很高的一招:先看分区数,再看别的。

三、分布式聚合:两阶段合并的语义

分区并行解决了"扫得快",这一节讲"算得对、算得省"——聚合在分区上是怎么执行的。

3.1 标准聚合的两阶段执行

group by 聚合在多分区数据上天然是两阶段的:

代码语言:Plaintext
复制
map 阶段(各分区并行): 每个分区先算局部聚合
    分区 A: (deviceId=PUMP-007, day=D1, sum=350, count=100)
    分区 B: (deviceId=PUMP-007, day=D1, sum=280, count=80)
    分区 C: (deviceId=PUMP-007, day=D1, sum=410, count=120)
              │
              ▼
reduce 阶段: 按分组键合并局部结果
    sum = 350+280+410 = 1040
    count = 100+80+120 = 300
    avg = 1040 / 300

sum / count / max / min 这类聚合的局部结果可以直接合并,两阶段执行没有语义问题,引擎自动完成。

3.2 自定义指标要想着"能不能合并"

两阶段模型对标准聚合是透明的,但当你写自定义聚合逻辑时,就要自己考虑合并语义了。一个典型陷阱:想在 group by 里算一个"合并后才能算"的指标。

比如"每台设备每天的中位数"——中位数不能像 sum 那样把各分区的中位数再取中位数合并。这类指标的执行要么由引擎退化为收集原始数据再算(代价是数据移动),要么需要你自己设计可合并的中间量(比如用 t-digest 这类可合并的近似结构的思想)。

实践上的建议:

  • 能用标准聚合表达的,用标准聚合sum + count 拆开算再相除,比硬凑一个"一步到位"的自定义逻辑更利于两阶段执行
  • iif + sum 是个好用的可合并模式。第一节例子里的 sum(iif(rms > 5.0, 1, 0)) 就是:局部可算、合并就是加法,两阶段友好
  • 分位数类指标接受近似。工程上报表用 P95,通常 pctile 就够了,不必强求精确中位数的完美并行

四、context by 的并行边界

context by 是 DolphinDB 处理时序的利器(前面多期都在用),这一节从执行角度讲它的并行边界。

4.1 组间天然并行,组内依赖顺序

context by deviceId 的语义是"按设备分组,组内按顺序算窗口函数"。从执行视角看:

  • 组间完全独立——PUMP-007 的滑动均值和 PUMP-008 的没有任何关系,可以分配到不同任务并行算
  • 组内有序依赖——mavg / prev / mstd 这类函数依赖组内前后文,必须组内按时间顺序串行推进

所以一条 context by 查询的理想并行度 ≈ 分组数。设备上千台,理论并行度上千,天然吃满集群——这也是为什么特征工程(十二期)用 context by 滚动算特征很少成为瓶颈。

代码语言:SQL
复制
// 组间并行的典型形态:每台设备独立算滚动特征
select deviceId, ts,
       mavg(rms, 600)  as rms_ma,
       mstd(rms, 600)  as rms_std
from loadTable("dfs://iot", "vibration")
context by deviceId

4.2 全局状态是并行杀手

会破坏这个并行性的写法,是引入跨组全局状态。最常见的两种:

img
img

全局排序后再分组

代码语言:SQL
复制
// 反例:先全局 order by,再 context by
select ... from (
    select * from loadTable("dfs://iot", "vibration") order by ts
) context by deviceId

内层的全局 order by 需要把所有分区的数据拉到一起排一次,这一步就是单点串行。而实际上 context by deviceId 本身就隐含"组内排序"的语义(配合 csort ts),根本不需要先做全局排序:

代码语言:SQL
复制
// 正解:context by 分组 + csort 组内排序,各组并行
select deviceId, ts, mavg(rms, 600) as rms_ma
from loadTable("dfs://iot", "vibration")
context by deviceId csort ts

先聚合再窗口:如果先把数据按时间聚合成一条全局序列再算窗口(比如把全部设备混在一起按分钟求和再 mavg),分组独立性就没了。需要"总体指标"时,明确它本来就是全局的,接受这部分串行;需要"设备级指标"时,守住 context by deviceId 的边界。

一句话:context by 的并行红利来自分组独立,写 SQL 时守住组边界,别让全局操作横插进来。

五、把并行写进 SQL:一份实践清单

前几节是原理,这一节收成可以直接照着做的清单。

5.1 过滤先行,且过滤条件要"够得着分区"

分区裁剪发生在过滤阶段。两个要求:

  1. where 里带上分区列。表按天分区,where 里就要有时间条件;按设备 HASH 分区,带上 deviceId 等值条件能进一步裁剪。分区列不出现在过滤条件里,裁剪就无从谈起
  2. 条件写得可下推where date(ts) = 2026.03.15 不如 where ts between 2026.03.15 00:00:00 : 2026.03.15 23:59:59——前者对列套了函数,裁剪判断更难;后者是裸列上的范围条件,裁剪直接生效

5.2 少读列

列存的基本功,但在写宽表查询时最容易忘:select * 会把二十几个列全读出来,而统计往往只用三四个。显式列出需要的列,I/O 直接降一个量级。这不是新知识,但在"先跑通再说"的临时分析里,select * 是惯性写法——批量任务上线前值得全部过一遍。

5.3 大任务拆批

一次性跑半年的全量统计,不如按天/周拆成小批次循环跑,每批结果增量写入结果表。好处有三:

  • 每批命中的分区数与集群并行度匹配,资源利用率平稳
  • 失败重跑只损失一批,不用从头再来
  • 中间进度可见,长任务不再是个黑盒
代码语言:SQL
复制
// 按天拆批的季度统计骨架
for(d in 2026.01.01..2026.03.31) {
    result = select deviceId,
                    avg(rms) as rms_mean, max(rms) as rms_peak
             from loadTable("dfs://iot", "vibration")
             where ts between datetime(d) : datetime(d) + 1d - 1s
             group by deviceId
    // 增量写入结果表
    loadTable("dfs://iot", "quarterly_stats").tableInsert(result)
}

5.4 排查一个慢查询的顺序

最后给一个诊断顺序,遇到慢查询按这个次序过:

  1. 命中的分区数——是不是分区粒度导致任务太少?(第二节)
  2. 过滤条件——分区裁剪有没有生效?where 里有没有裸分区列?(5.1)
  3. 读的列——是不是 select *?(5.2)
  4. 全局操作——有没有全局排序/全局聚合横在 context by 前面?(第四节)
  5. 都排除了,再看集群配置层面:单查询可并发调度的分区任务数上限、节点执行器并发数等参数(以所用版本的官方文档为准),以及是否有别的任务在抢资源
img
img

十次慢查询里,前两步就能解释七八次。

六、手写 MapReduce 的边界

前面的内容都在"让 SQL 自动并行"。DolphinDB 也提供了 mr / imrsqlDS 这套 MapReduce API,允许完全自定义分布式计算逻辑——二期的多范式评测里已经展示过用它做分布式回归的例子,语言能力层面不再重复,这里只从执行角度补一句什么时候真的需要它

  • 标准 SQL + 标准聚合能表达的——不需要。SQL 自动两阶段执行,已经吃满分区并行
  • 自定义聚合且可设计出可合并中间量的——可以考虑。sqlDS 把查询转成数据源,mr 的 map 阶段逐分区跑你的逻辑,reduce 合并
  • 迭代型计算(每轮用上一轮的结果,比如分布式优化)——imr 的场景

大多数 IoT 批量分析停在前两类。把 MapReduce 当成"SQL 撑不住时的逃生通道",而不是性能优化的第一步——先修分区、修过滤、修写法,再考虑换工具

七、写在最后

这篇补上了第三块拼图——批量查询在集群上怎么并行执行。

回到那句主线:分布式数据库里,查询的上限不是函数多强,而是被拆成了多少个能同时跑的任务。

  • 分区是并行的单位——分区数够,并行度才够(第二节)
  • 标准聚合自动两阶段合并——写自定义指标时想着"能不能合并"(第三节)
  • context by 的并行红利来自分组独立——守住组边界,别让全局操作横插(第四节)
  • 清单化排查——分区数、裁剪、读列、全局操作,按序过一遍(第五节)

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 摘要
  • 一、慢在哪:一个 40 分钟的季度统计
  • 二、分区是并行调度的单位
    • 2.1 执行模型:一条查询怎么被拆开
    • 2.2 分区数的两头陷阱
    • 2.3 回到那个 40 分钟的查询
  • 三、分布式聚合:两阶段合并的语义
    • 3.1 标准聚合的两阶段执行
    • 3.2 自定义指标要想着"能不能合并"
  • 四、context by 的并行边界
    • 4.1 组间天然并行,组内依赖顺序
    • 4.2 全局状态是并行杀手
  • 五、把并行写进 SQL:一份实践清单
    • 5.1 过滤先行,且过滤条件要"够得着分区"
    • 5.2 少读列
    • 5.3 大任务拆批
    • 5.4 排查一个慢查询的顺序
  • 六、手写 MapReduce 的边界
  • 七、写在最后
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档