从ES裸查到Doris on ES:作业帮实时数仓查询层重构实践

发布时间:2026/9/19 15:35:13
从ES裸查到Doris on ES:作业帮实时数仓查询层重构实践 简介《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 秒且不支持 DDLES 自带的 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-count10-20攒批缓冲数量太小容易频繁 flushsink.buffer-flush.max-rows50000-100000写入行数阈值根据单行大小调整sink.buffer-flush.max-bytes100MB批量字节数阈值避免单次导入过大sink.buffer-flush.interval2-5s时间阈值满足实时性要求即可sink.max-retries3写入失败重试次数重试过多会造成延迟堆积这些参数需要根据业务容忍的延迟和 Doris BE 的处理能力平衡。如果追求秒级可见性sink.buffer-flush.interval可以设到 1 秒但导入频率升高会对 FE 产生更多轮询请求。如果业务容忍分钟级延迟5 秒甚至 10 秒的攒批间隔对 BE 更友好。4. 基于 Doris 的实时查询系统架构与 Doris on ES 执行原理4.1 架构总览数据摄入、清洗、查询三层整个实时查询系统的架构分三层。数据摄入层用 Kafka 承接业务侧日志和 MySQL BinlogFlume 做链路传输或备份数据清洗层用 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 索引的 mappingDoris 外表的 SchemaDoris 原生表的 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(hostfe_host, useruser, passwordpass, port9030) 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(indexlesson_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 命中、谓词下推这几层有明确把握而不是等线上出了问题再逐条排查。本文还有配套的精品资源点击获取