简介:《Doris在数仓中的实践》是一份面向大数据工程师与数仓架构师的技术 PDF,围绕 Doris 这个 MPP 架构 OLAP 引擎,系统梳理其在企业数仓中的选型依据与落地经验。内容先交代业务背景与旧方案性能差、维护成本高等痛点,再依次说明技术选型、系统架构,以及 Kafka、Flume、Spark、Doris 协作的实时查询链路;同时深入 Aggregate 模型、RollUp 预聚合、Base 表与 Bitmap 存储,并解释 Doris on ES 的高性能查询设计。资源为单个 PDF 文件,大小 1.77MB,便于阅读收藏。目前已有 1323 人学习,适合正在调研实时 OLAP 引擎、规划数仓实时化改造的读者。透过 BI 报表、PV/UV、教研工作台等案例,可掌握 Doris 的建模手段与调优方向,也能理解其相比 Presto on ES、Druid 等方案的选型优势。整体内容包含真实业务场景下的架构拆解与性能对比,能帮助读者快速形成实时数仓建设思路,避开常见设计坑点。
1. 从 ES 裸查到 Doris on ES:作业帮实时数仓查询层的重构路径
流量分析、教研工作台、BI 报表,这些场景在作业帮的数仓体系里有着截然不同的查询特征:前者要秒级返回 PV/UV 聚合,后者要在一节课内对千万级明细做任意列过滤和分组统计。过去这些需求分别压在 Druid、ES、Kafka 接口和一堆定制 API 上,业务线每接一个需求就要重新走一遍「数据清洗 → 建索引 → 写接口」的流程,交付周期按周计算,ES 裸查在千万级数据上跑一个聚合甚至要十个小时以上。2020 年上半年团队用 Doris 逐步替换掉这套组合拳,半年时间接入 7 条业务线、近 1T 数据,查询延迟从分钟级降到秒级,且没有出现 P2 及以上事故。这篇文章把这套方案里的选型依据、Doris 数据模型设计、Doris on ES 的执行原理以及元数据治理细节拆开讲清楚,适合正在做实时数仓选型、或者已经在用 Doris/ES 但查询性能不达标的工程师参考。
2. 查询引擎选型与 Doris 核心设计:为什么是它
2.1 旧架构的痛点:重复建设与性能瓶颈
作业帮过去支撑业务查询的组件包括 Druid、ES、Kafka、Spark 以及大量手写 API。这套体系的问题从两个维度暴露出来。在建设成本上,每个业务线都是 case by case 地开发,从 Kafka 取数、Spark 清洗、写入 ES、再封装接口,链路重复建设且无法复用,非标接口需要独立维护,业务侧还要裸查 ES——学习成本高、SQL 不完备、稳定性差。从查询能力上看,Druid 只支持聚合、明细数据丢失,Presto on ES 延迟 99 分位约 25 秒且不支持 DDL,ES 自带的 SQL 方言(6.3 版本)语法不完备,不支持 join 和多列 group by。这些组件各自能解决一部分问题,但没有任何一个能同时覆盖「明细 + 聚合」两类查询需求。
对比之下,Doris 本身的特性恰好补齐了这些短板。它是 MPP 架构的 OLAP 引擎,FE 负责解析和元数据管理,BE 负责执行和存储;兼容 MySQL 协议和标准 SQL;同时支持离线批量导入和实时流式导入;支持 Rollup 表和 Base 表的智能路由;支持 Schema 在线变更。这套特性意味着业务侧可以直接用 MySQL 客户端连上来写 SQL,不用再学一套查询方言,也不用依赖独立的接口层做转发。Doris on ES 则通过外表(External Table)的方式把 ES 的索引映射成 Doris 表,可以在 Doris 里用完整 SQL 语法查询 ES 中的数据。两者结合后,聚合类流量分析走 Doris 原生表,明细类教研工作台走 Doris on ES,统一了实时查询的入口。
2.2 MPP 架构与查询执行模型的匹配度
选择 Doris 而不是 Presto 或 ClickHouse,需要结合业务特征来看。流量分析场景的查询以 PV/UV 为主,UV 计算依赖精确去重。Doris 的 Aggregate 模型配合 BITMAP 类型,可以把 UV 预聚合到分钟甚至小时粒度,查询时直接对预聚合结果做 BITMAP 并集计算,避免了扫描明细。教研工作台的查询特征是「给定 lesson_id 和 teacher_id,统计出勤学生数」,本质上是一个带过滤条件的 group by,并且需要实时写入——学生出勤数据是持续产生的。Presto 在 ES 上的表现不佳,主要原因在于它无法把 limit、过滤条件下推到 ES,导致大量数据跨节点传输;而 Doris on ES 支持谓词下推和分片级并发扫描,在架构上更适合这类场景。
Doris 的另一个关键设计是 Rollup 预聚合。Base 表存储明细或原始聚合数据,Rollup 表存储按更粗粒度预聚合的结果。查询时 FE 会根据查询的维度和聚合函数自动选择最优的 Rollup 表,而不是让用户手动指定。这个「智能路由」能力直接影响查询性能:同样一份 UV 数据,按天聚合的 Rollup 表在查询「某天活跃用户数」时只需要扫描一行,而从 Base 表算则需要扫描全表明细。实际使用中,Rollup 表不能盲目建多,每增加一张 Rollup 表都会带来数据导入时的额外聚合开销,通常只针对高频查询维度组合建 2-4 张。
3. Doris 数据模型设计与实时写入链路
3.1 Aggregate 模型:UV 场景下的 Rollup 设计
流量分析场景最典型的查询是「作业帮主 App 某天活跃用户数」和「某个小时段各版本下的活跃用户数」。前者只需要按天做 UV 去重,后者需要按小时、版本维度做 UV 去重。这里对应两类查询频率:天级聚合每天被报表任务大量调用,小时级的版本维度分析则用于运营排查问题,频率相对较低。如果用一张 Base 表存全量明细,每次查询都扫描明细数据,在千万级 UV 的体量下延迟无法接受。
用 Aggregate 模型建表示例如下:
CREATE TABLE app_uv_agg ( dt DATE, hour INT, app_version VARCHAR(32), uv BITMAP BITMAP_UNION, pv BIGINT SUM ) AGGREGATE KEY (dt, hour, app_version) DISTRIBUTED BY HASH(dt) BUCKETS 16 PROPERTIES ("replication_num" = "3");建表完成后,创建按天聚合的 Rollup 表:
ALTER TABLE app_uv_agg ADD ROLLUP rollup_dt_uv (dt, uv, pv);这里有几个要点。第一,uv字段使用BITMAP类型,配合BITMAP_UNION聚合函数,Doris 会在导入时自动对相同 Key 的 bitmap 做合并,而不需要业务侧先去重。第二,AGGREGATE KEY只能包含维度列,指标列必须在 Key 之外;查询时如果group by的维度是 Key 的子集,Doris 就能命中 Rollup 表。第三,DISTRIBUTED BY HASH(dt)决定了数据分布方式,UV 场景的查询几乎总是带时间范围,按天 hash 分桶可以保证分桶裁剪生效,避免全表扫描。如果业务侧高频按app_version过滤,可以考虑把app_version加入分桶键,但这样会导致桶数膨胀,需要根据实际查询特征权衡。
3.2 明细写入:Flink SQL 实时导入
实时流量数据通过 Kafka 接入,然后用 Flink SQL 写入 Doris。Flink SQL 写 Doris 的常见方式是通过 Doris 的 Stream Load 接口,Flink 官方连接器封装了这部分逻辑。示例 DDL 如下:
CREATE TABLE doris_sink ( dt DATE, hour INT, app_version STRING, uv BITMAP, pv BIGINT ) WITH ( 'connector' = 'doris', 'fenodes' = 'fe01:8030,fe02:8030', 'table.identifier' = 'dwd.app_uv_agg', 'username' = 'writer', 'password' = '******', 'sink.label-prefix' = 'doris_uv_20240115', 'sink.properties.format' = 'json', 'sink.properties.columns' = 'dt,hour,app_version,uv,pv', 'sink.enable.batch-mode' = 'true' );使用 Flink SQL 写入 Doris 时,sink.label-prefix必须每条作业唯一,Doris 的 Stream Load 依赖 label 实现幂等写入,label 重复会导致作业报错。sink.enable.batch-mode开启后,写入会攒批提交,减少小文件数量,对提升导入性能和降低 BE 压力有明显帮助。另外需要注意,Doris 表如果是 Aggregate 模型,Flink SQL 写入的数据会被 Doris 按照聚合 Key 和聚合函数自动合并,写明细的dup表则不需要有聚合语义的 DDL,直接映射字段即可。
提示:Flink SQL 中如果
uv字段是从明细数据实时计算出来的 bitmap,需要在 Flink 侧先用BITMAP_UDF_TO_BITMAP或类似 UDF 把数值转成 bitmap 再写入,否则 Druid/SQL 端到端链路里这一步最容易出现类型不匹配。
3.3 写入链路参数调优
实际生产环境里实时写入的稳定性往往比查询更重要。Stream Load 的批量大小、并发数、BE 磁盘类型等都会影响端到端延迟。常见的参数配置思路如下:
| 参数 | 推荐值 | 说明 |
|---|---|---|
sink.buffer-count | 10-20 | 攒批缓冲数量,太小容易频繁 flush |
sink.buffer-flush.max-rows | 50000-100000 | 写入行数阈值,根据单行大小调整 |
sink.buffer-flush.max-bytes | 100MB | 批量字节数阈值,避免单次导入过大 |
sink.buffer-flush.interval | 2-5s | 时间阈值,满足实时性要求即可 |
sink.max-retries | 3 | 写入失败重试次数,重试过多会造成延迟堆积 |
这些参数需要根据业务容忍的延迟和 Doris BE 的处理能力平衡。如果追求秒级可见性,sink.buffer-flush.interval可以设到 1 秒,但导入频率升高会对 FE 产生更多轮询请求。如果业务容忍分钟级延迟,5 秒甚至 10 秒的攒批间隔对 BE 更友好。
4. 基于 Doris 的实时查询系统架构与 Doris on ES 执行原理
4.1 架构总览:数据摄入、清洗、查询三层
整个实时查询系统的架构分三层。数据摄入层用 Kafka 承接业务侧日志和 MySQL Binlog,Flume 做链路传输或备份;数据清洗层用 Spark 和 Flink SQL 做净化、维表关联、格式转换;存储查询层用 Doris 承接聚合查询,用 Doris on ES 承接明细查询。业务侧通过 OpenAPI 或直连 Doris 的 MySQL 协议端口访问数据,前端工作台、BI 报表等直接面向 Doris 发 SQL。
这条链路和过去的方案相比,核心变化是「查询入口统一」。过去一个报表需求要串联 Kafka、Spark、ES、自研 API 四个组件,现在业务侧只需要面向 Doris 写 SQL——能命中 Rollup 的就走 Doris 原生表,需要明细检索的就走 Doris on ES 外表。Doris on ES 的映射关系在 FE 侧完成,用户无感知。这也意味着稳定性问题从「多个组件各自排查」收敛到「Doris 与 ES 两端排查」,运维复杂度显著降低。
4.2 Doris on ES 的查询改写与两阶段取数
Doris on ES 之所以比裸查 ES 快,核心在于执行策略的差异。ES 的常规搜索走 Query-Then-Fetch 两阶段:先通过分片和排序逻辑拿到 Top ID 列表,再根据 ID 集合去获取完整文档。Doris on ES 在默认分析场景下使用 Query-And-Fetch 模式,直接在分片上过滤并返回最终结果,减少了一次跨节点取数的轮次。
Doris on ES 的具体执行流程可以概括为四步。第一步,FE 将 SQL 改写成 ES DSL,下推到 ES;第二步,BE 节点并发访问 ES 各分片,每个分片只扫描自身数据,实现分片级并发;第三步,对 ES 返回的文档执行列裁剪、谓词过滤、聚合等计算;第四步,如果 SQL 中带 limit 或 first-scroll 语义,会触发提前终止。整个流程中,Doris 的 BE 不做全量数据拉取,而是把计算尽量推向 ES 分片,传输层只保留需要的列。这种「计算找数据」的策略比把数据拉回来再算要高效得多。
4.3 谓词下推与扫描优化细节
Doris on ES 的谓词下推包括两个层面:一是 Doris 将 SQL 中 WHERE 条件下推到 ES,让 ES 在 Lucene 层面过滤;二是 source/path filter 减少 ES 返回的数据量。实际使用中需要特别注意类型一致性——ES 侧字段如果是keyword类型,Doris 外表建表时对应字段不能定义为bigint,否则下推的查询可能执行时报错或结果异常。类型对齐这件事需要纳入建表规范,而不是等问题暴露再修数据。
扫描速度优化方面,Doris on ES 支持列存优先原则,即 BE 从 ES 拉取数据时优先选择列存格式读取,对宽表场景能明显降低传输量。分片级并发数可以通过外表属性es.nodes和 BE 数量间接控制,一般来说 ES 分片数最好是 BE 节点的整数倍,这样每个 BE 能均匀消费分片任务,避免某个 BE 空闲、某个 BE 过载。
提示:ES 索引的
index.mapping.total_fields.limit如果设置过小,Doris on ES 查询时容易字段映射失败。建议对映射到 Doris 的索引提前检查字段数上限,并确保不需要检索的字段关闭doc_values或index: false,降低存储开销同时提升扫描性能。
4.4 环境配置与常见故障排查
Doris on ES 初期接入时最容易踩的坑集中在三块。第一,网络连通性:BE 节点必须能访问 ES 的 transport 端口(默认 9300),否则查询长时间超时。第二,ES 集群安全认证:如果 ES 开启 x-pack 认证,在 Doris 建外表时需要在 properties 中配置"es.username"和"es.password",早期版本不支持加密传输,要用 HTTP。第三,查询超时设置:Doris 默认查询超时时间受query_timeout限制,外表查询涉及 ES 扫描时耗时通常高于原生表,需要按需调大。用 MySQL 客户端执行:
SET query_timeout = 120;这条命令只在当前会话生效,适合调优时使用。如果要在全局生效,可以修改 FE 配置项query_timeout的默认值。查询超时日志在 FE 节点fe.log中表现为query timeout关键字,排障时先确认查询类型——是 ES 扫描慢还是 BE 聚合慢。
5. 元数据管理与 Schema 一致性:Doris on ES 稳定性的保证
Doris on ES 的使用体验虽然统一到了 SQL 层,但本质上是两个存储引擎的协作。ES 索引和 Doris 外表必须保持字段名、字段类型的严格一致,否则会出现三种典型问题:类型不一致导致的查询报错;字段缺失导致的数据同步质量不可控;新增字段后需要两侧同步修改建表语句。这些都是线上事故的高发源头,因此团队把元数据管理作为一个独立的治理模块来设计。
元数据管理覆盖的核心对象包括:ES 索引的 mapping,Doris 外表的 Schema,Doris 原生表的 Rollup 定义,以及 Flink SQL 写入端的 DDL。维护目标是「一处定义、多处复用」。ES 建索引时,需要把字段清单、类型、是否开启 doc_values、是否索引等属性统一维护到元数据中心;Doris 侧建外表时从元数据中心读取字段定义,自动生成建表语句;Flink SQL 写入时根据元数据中心生成目标表 DDL,确保上游写入和下游查询字段语义一致。
这里给出一个简易的元数据核对脚本思路,定期巡检 Doris 外表和 ES index mapping 是否一致:
import pymysql from elasticsearch import Elasticsearch # 读取 Doris 外表定义 conn = pymysql.connect(host='fe_host', user='user', password='pass', port=9030) cur = conn.cursor() cur.execute("SHOW CREATE TABLE ads_lesson_attend_detailed") doris_schema = cur.fetchone()[0] # 读取 ES 索引 mapping es = Elasticsearch(['es_host:9200']) es_mapping = es.indices.get_mapping(index='lesson_attend_detailed') # 对比字段名集合 doris_fields = set(parse_doris_schema(doris_schema)) # 伪代码 es_fields = set(es_mapping['lesson_attend_detailed']['mappings']['properties'].keys()) print("缺失字段:", es_fields - doris_fields)这段脚本的逻辑是:从 Doris 的SHOW CREATE TABLE结果中解析出外表字段集合,再读 ES 的 mapping 拿到索引字段集合,做差集对比。生产环境可以做成定时巡检任务,每次发布前执行一次,避免 Schema 不一致导致线上查询报错。需要注意的是,Doris on ES 外表不支持自动感知 ES 新增字段——ES mapping 加字段后,Doris 侧也需手动ALTER TABLE添加或重建外表。
6. 查询性能验证方法:从 SQL 执行计划到 Rollup 命中率
Doris 查询性能的验证不能只靠肉眼感受,需要从执行计划层面确认是否命中了 Rollup 表、是否走对了索引和外表路由。Doris 提供了EXPLAIN命令来查看 SQL 的执行计划,这一步是排查慢查询的第一道工序。
对一个典型的 UV 查询,执行计划如下:
EXPLAIN SELECT dt, COUNT(DISTINCT uv) FROM app_uv_agg WHERE dt = '2024-01-15' GROUP BY dt;执行计划中重点关注两部分:一是TABLE节点显示的表名,如果命中了 Rollup 表,会显示app_uv_agg下的rollup_dt_uv而非 Base 表;二是AGGREGATE节点的聚合方式,如果看到BITMAP_UNION,说明 UV 聚合在预聚合阶段已完成,扫描的数据量远小于 Base 表。如果TABLE节点直接显示 Base 表且扫描行数超过百万,说明 Rollup 没有命中,需要检查查询维度和 Rollup 定义是否匹配。
Doris on ES 的查询验证则要看EXPLAIN中的SCAN节点,确认谓词是否下推到 ES。正常执行计划中应该能看到ES_SCAN节点,且包含下推的过滤条件。如果WHERE条件没有出现在ES_SCAN节点的PREDICATES中,说明谓词下推失败,查询会把 ES 数据全量拉回 BE 再过滤,性能急剧下降。这种情况下优先检查外表字段类型和 ES 索引字段类型是否一致——这是导致下推失败的最常见原因。
对于线上慢查询,还可以通过 Profile 机制做进一步定位。Doris 支持开启查询 Profile,在 FE 节点执行:
SET enable_profile = true;执行慢查询后,在http://fe_host:8030/query_profile页面查看对应查询的 Profile。重点看SCAN节点的rows_read和bytes_read,以及EXCHANGE节点的网络传输量。rows_read远大于预期值时说明扫描范围过大,优先调整分桶键或 Rollup 设计;bytes_read大且耗时集中在EXCHANGE阶段,说明跨节点数据传输是瓶颈,需要通过谓词下推或列裁剪来压缩传输量。
最后一个常用的验证维度是前端工具的可视化监控。Doris 的 Grafana 监控面板中核心指标包括 BE 的scan行数、scan耗时、query耗时分位数,以及 Stream Load 的导入成功率。日常巡检以扫描行数和查询分位延迟为主线,如果 P99 查询延迟持续升高而扫描行数没有明显变化,大概率是 BE 机器 CPU 或磁盘 IO 到达瓶颈,需要扩容或优化分桶布局。
这套方法论可以用在任何 Doris 集群的巡检和调优过程中,不依赖具体业务——先看执行计划是否按预期走 Rollup/外表路由,再看 Profile 中扫描和传输的量是否合理,最后针对性调整 Schema 或查询逻辑。Doris 上手快是因为 SQL 语法友好,但要做到「把查询延迟稳定控制在秒级」,关键是对执行计划、Rollup 命中、谓词下推这几层有明确把握,而不是等线上出了问题再逐条排查。
本文还有配套的精品资源,点击获取