ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Doris:数据湖加速与实时数仓的统一查询中枢

Apache Doris:数据湖加速与实时数仓的统一查询中枢 先聊一个这两年很多团队都会遇到的场景数据湖、实时数仓、统一的 SQL 查询入口这三件事业务都想要但往往不是同一套引擎能搞定的。数据湖便宜能装海量数据可查询性能总差口气实时数仓延迟低但存不下所有历史明细统一查询层听起来完美搞到最后却变成在 BI 工具里写一堆不同引擎的方言。我最近在几个项目里反复折腾 Apache Doris慢慢发现它其实很适合站在中间当那个“高性能中枢”——对外接数据湖和消息流对内提供实时数仓能力往上又是统一的查询入口。这篇文章就把我实际的用法、踩过的坑、以及为什么在多个引擎之间选它来做加速层的思考都整理出来给正在纠结数据架构选型的朋友一个参考。1. 先把问题说清楚为什么现代数据架构里总缺一个“中枢”1.1 数据湖、实时数仓、统一查询层各管什么数据湖首先要解决的是“什么都存得下”的问题。业务原始日志、订单流水、用户行为事件、甚至图片元数据都可以用很低成本的存储先堆在湖里比如 HDFS、S3、OSS 上的 Parquet/ORC 文件再通过 Hive、Iceberg、Hudi、Paimon 这类表格式管理元数据和事务。它的优点是存储便宜、格式开放、生态广缺点是查询性能做不到很低延迟尤其是要频繁跑交互式报表或者服务化查询时光扫描文件元数据和列裁剪都够折腾。实时数仓要解决的是“指标不能等到明天”的问题。线上业务要看今日实时 GMV、实时新增用户、实时转化率数据从业务库或消息队列产生到出指标需要控制在秒级到分钟级。传统做法是 Kafka 进实时计算引擎再落地到可查询的 OLAP 引擎。数仓能保证低延迟和高并发查询但通常不会把所有原始日志都存进去成本太高。统一查询层要解决的是“用户不想关心数据存在哪儿”的问题。业务分析师、数据产品、运营后台只需要一个 SQL 入口最好能跨数据源 join能同时读到湖里的历史数据和数仓里的实时数据。过去很多团队用 Presto/Trino 做联邦查询但它自身不存数实时更新和写入能力薄弱遇到高并发服务化查询往往要搭配额外的缓存层。1.2 多引擎各自为战的乱象很多公司的真实情况是数据湖放一套 Presto 集群跑 Ad-hoc实时数仓用 ClickHouse 或 DorisBI 报表再直连一个 MySQL 汇总库。听起来分工明确实际用起来问题不少。首先是数据口径不统一。同样的“用户数”在湖上用 SQL 算一遍和数仓里算一遍因为去重逻辑、时区处理、事件过滤条件稍微不同结果就对不上。其次是权限体系割裂湖上有一套 Ranger/Sentry数仓有自己的用户权限BI 还要单独配账号。再一个是运维成本高每套引擎都有各自的监控、调优、扩容节奏。我们当时就在想与其继续加引擎不如选一个能“骑在”数据湖上、又能自己做实时数仓的引擎作为中枢统一提供查询服务。1.3 为什么是 Apache Doris 站到这个位置上Doris 本身是一个 MPP 架构的分布式 SQL 数据库擅长高并发点查和高性能聚合查询这点很多人知道。但真正让我注意到它的是近两个版本里把数据湖接入做得越来越“实在”。它支持多 Catalog 方式对接 Hive、Iceberg、Hudi、Paimon、Delta Lake还能通过外部表直接查 MySQL、PostgreSQL、Elasticsearch查询时既能做联邦分析也能通过物化视图把湖上常用结果落到 Doris 内部存储形成“湖上原始数据 Doris 加速层”的组合。加上它自带的实时导入链路Stream Load、Routine Load、Flink Doris Connector一个引擎就能同时覆盖实时数仓和统一查询层两个角色。对我来说最关键的一点是Doris 不是要替代数据湖而是给数据湖加一个高性能的“前端”。湖还是那个湖历史数据继续低成本放着但常用查询路径会经过 Doris 的物化视图和缓存因此用户感知到的查询速度是数仓级的。2. 拆解 Doris 的三个核心定位它到底做了什么2.1 数据湖加速不搬数据但让查询变快数据湖加速这个说法听起来有点虚第一个要搞明白的是Doris 去查湖上数据到底比 Trino 或者 Spark SQL 快在哪。Doris 在 FE 节点维护了数据湖表的元数据缓存BE 节点则负责实际扫描外部文件。查询时Doris 的查询规划器会把过滤条件下推到文件格式层比如 Parquet 的 row-group 统计信息、ORC 的 stripe 统计信息利用 min/max 索引跳过大量不满足条件的文件。这点很多引擎都支持。Doris 更独特的是文件缓存File Cache能力首次从远端对象存储或 HDFS 读取的数据会按照一定策略缓存在本地 BE 磁盘上后续相同数据块的查询直接读本地网络和带宽开销大幅下降。我在实际场景里测过同样的 Iceberg 表查询最近 7 天用户行为明细第一次查询耗时 15 秒左右开启文件缓存后第二次查询降到 2 秒多。如果再把高频聚合结果做成 Doris 内部的物化视图查询耗时能压到几百毫秒。这种体验让数据湖不再是“只适合跑离线批任务”的存储而是可以直接支撑 BI 看板的实时分析。2.2 实时数仓从明细到汇总一条链路打通Doris 作为实时数仓底层是列式存储、向量化执行、MPP 并行查询单表聚合能力很强。重要的是它支持三类数据模型Duplicate Key、Aggregate Key、Unique Key。其中 Unique Key 模型配合 Merge-on-Write 实现主键更新可以用流式导入完成业务库的 CDC 同步也可以承接 Kafka 里的实时明细数据。典型的做法是业务库 MySQL 用 Flink CDC 同步到 Doris 的 ODS 层Kafka 实时日志用 Routine Load 或 Flink Doris Connector 写入 DWD 层再通过 Doris 的 Insert Into Select 做实时或准实时的汇总到 ADS 层。这个过程不需要额外部署一套 OLAP 引擎也不需要维护复杂的 Lambda 架构一套 Doris 就能做到 OD S 到 ADS 的分层建模。我到现在还记得第一次用 Doris 的 Unique Key 模型做精确去重 UV 的场景。以前用 Kafka Flink Redis 维护实时 UV状态后端一出问题就要回追换到 Doris 后直接把明细写入 Unique Key 表用 Bitmap 函数做精确去重既不需要外部状态存储查询还能并行跑凌晨回刷数据也容易得多。2.3 统一查询层让用户只认一个 SQL 入口统一查询层并不是简单地把 SQL 转发到后端引擎而是要解决“一个 SQL 能同时查多个源”的问题。Doris 的多 Catalog 机制非常关键。你可以创建一个名为 hive_catalog 的 Catalog映射 Hive Metastore创建一个 iceberg_catalog映射 Iceberg REST Catalog再创建一个 mysql_catalog映射业务库。用户查询时只需要写hive_catalog.db.table或者internal.db.tableDoris 会把这些源当成统一的关系表来处理。更重要的是Doris 支持跨 Catalog 的 Join比如把数据湖里的历史用户画像和实时数仓里的今日行为 join 起来不需要先把数据导出再合并。我在做统一查询层时还会在 Doris 上建一层逻辑视图View把业务口径封装好。比如“核心交易指标宽表”底层可能是湖里的订单历史、实时数仓里的支付流水、MySQL 里的商品维表但在用户眼里就是一个大宽表直接 select 就行。这层视图也方便权限管控不需要给分析师开放底层 Catalog 的权限只开放视图即可。3. 关键功能的使用细节哪些才是真正的提效神器3.1 Multi-Catalog 配置的三个要点Multi-Catalog 是 Doris 接入外部数据源的入口但很多人配置完发现查询很慢问题往往出在细节上。我总结三个最容易踩坑的点。第一个是元数据同步方式。Doris 不会实时去拉取 Hive Metastore 的所有分区默认情况下首次查询一个新建的外部表Doris 会主动拉取元数据。如果你的湖表是个超大分区表首次查询会比较慢。建议手动执行REFRESH CATALOG hive_catalog或者设置 FE 参数enable_metacache相关配置让元数据常驻缓存。对于分区频繁变化的表可以设置定时刷新避免每次查询都触发全量拉取。第二个是文件系统配置。连接 HDFS 时需要把hdfs-site.xml、core-site.xml放到 FE 和 BE 的conf目录下或者通过 Catalog 属性传dfs.nameservices、dfs.ha.namenodes等参数。连接 S3 或 OSS 时需要配好 endpoint、access key、secret key。这个听起来简单但集群一多很容易漏 BE 节点的配置导致 FE 能查到元数据BE 扫描文件时却报权限或路径错误。第三个是谓词下推的真实效果。虽然 Doris 支持下推过滤条件但有些函数会阻断下推比如在过滤条件里对列做函数运算、使用非确定性函数等。我建议在关键查询上打开 Profile查看扫描行数和实际返回行数来判断下推是否生效。如果发现扫描了大量不该扫的分区优先检查查询条件是否命中了分区字段。3.2 物化视图加速数据湖的杀手锏如果只是通过外部表直接查数据湖性能始终受限于文件扫描。Doris 的异步物化视图允许你把数据湖上的常用聚合结果“物化”到 Doris 内部表存储并且支持透明改写。这意味着业务 SQL 不需要改Doris 会自动判断能否命中物化视图如果命中直接从加速表返回结果。实际使用时我一般会为三类场景建物化视图第一类是高基数维度的明细汇总例如按天、按省份、按商品维度的订单聚合第二类是跨 Catalog 的 Join 结果例如数据湖用户画像和实时数仓行为表 join 后的宽表第三类是周期性指标快照例如每天凌晨计算一次累计指标。物化视图刷新可以配置定时任务也可以手动触发。我遇到过一个问题物化视图的刷新 SQL 如果太重会占用系统资源影响实时查询。后来我把它单独设置资源组限制刷新任务的并发和内存避免高峰期互相干扰。3.3 实时导入链路怎么设计才稳Doris 支持多种实时导入方式选型要看数据源和延迟要求。Stream Load适合应用程序直接推送数据比如 Java 服务批量上报日志HTTP 接口导入秒级延迟。Routine Load适合消费 Kafka 数据Doris 常驻任务持续拉取消息可做过滤、列映射、时间戳转换是最常用的实时写入方式。Flink Doris Connector适合 Flink 做实时 ETL 后写入支持两阶段提交能保证精确一次语义。Insert Into Select适合 Doris 内部表之间的 ETL 和分层加工。链路设计上最容易忽略的是导入并发和副本数的关系。如果 BE 节点数量不多但导入任务并发开得过高会导致磁盘 IO 打满、版本堆积。我一般建议控制导入并发数为 BE 数量的 2 到 3 倍并设置routine_load_parallelism合理取值。另一个经验是实时写入尽量追加写入避免高频删除因为高频更新会触发 Compaction影响查询性能。4. 实操过程一个“湖仓一体”场景的落地全流程4.1 场景设定和表结构规划为了把上面的思路讲透我用一个比较典型的例子电商平台的用户行为分析。数据分布如下历史订单数据存放在数据湖的 Iceberg 表中存储在对象存储上分区字段是日期数据量大约几十 TB。商品维表存放在 MySQL 业务库中需要实时同步。用户实时点击流日志进入 Kafkatopic 为user_click_log。最终需要提供一个统一的查询入口让分析师在 BI 工具里既能查历史订单又能看实时点击流同时关联商品信息。整体架构就是数据湖Iceberg MySQL Kafka全部接入 Apache Doris对外提供 MySQL 协议查询。4.2 创建 Catalog 接入数据湖首先在 Doris 中创建 Iceberg Catalog。我使用的是 Iceberg REST CatalogCREATE CATALOG iceberg_catalog PROPERTIES ( type iceberg, iceberg.catalog.type rest, uri http://iceberg-rest:8181, warehouse s3://my-bucket/warehouse, s3.endpoint https://s3.amazonaws.com, s3.access_key AK..., s3.secret_key ... );创建后执行SHOW DATABASES FROM iceberg_catalog;能看到 Iceberg 里的所有库表。然后测试一个查询SELECT order_date, sum(amount) FROM iceberg_catalog.analytics.orders WHERE order_date 2025-01-01 GROUP BY order_date;第一次执行会比较慢因为需要拉取元数据和扫描远端文件。为了加速我会开启文件缓存SET enable_file_cache true;同时为常用的聚合结果创建物化视图。这里要注意物化视图如果要跨 Catalog 访问外部表需要确保引用的外部表元数据稳定。我采用手动刷新方式CREATE MATERIALIZED VIEW mv_order_daily BUILD DEFERRED REFRESH AUTO ON MANUAL AS SELECT order_date, user_id, count(order_id) AS order_cnt, sum(amount) AS amount_sum FROM iceberg_catalog.analytics.orders GROUP BY order_date, user_id;刷新之后后续查询如果没有更细粒度的要求Doris 会直接读取这个物化视图来响应速度从秒级降到毫秒级。4.3 接入 Kafka 实时点击流接下来用 Routine Load 接入 Kafka 的点击流日志。假设 Kafka 消息是 JSON 格式包含字段user_id、product_id、click_time、page。先在 Doris 建内部表CREATE TABLE dwd_click_log ( user_id BIGINT, product_id BIGINT, click_time DATETIME, page VARCHAR(128) ) DUPLICATE KEY(user_id, product_id, click_time) DISTRIBUTED BY HASH(user_id) BUCKETS 16 PROPERTIES ( replication_num 2 );然后创建 Routine Load 任务CREATE ROUTINE LOAD job_click_log ON dwd_click_log COLUMNS(user_id, product_id, click_time, page) PROPERTIES ( desired_concurrent_number 4, max_batch_interval 10, format json ) FROM KAFKA ( kafka_broker_list kafka1:9092,kafka2:9092, kafka_topic user_click_log, kafka_partitions 0,1,2,3, property.group.id doris_group );这里我踩过一个坑一开始没设置max_batch_interval导致小批量文件频繁触发导入BE 上产生大量小版本查询变慢。后来调整到 10 秒一个批次情况明显好转。要注意 Routine Load 的字段顺序要与表结构一致如果原始消息里有多余字段需要显式映射或忽略。4.4 实时同步 MySQL 维表商品维表数据量不大但更新频繁我选择用 Flink CDC 同步到 Doris。Flink 里配置 MySQL CDC source 和 Doris sink通过Flink Doris Connector写入。Doris 端建 Unique Key 模型表CREATE TABLE dim_product ( product_id BIGINT, product_name VARCHAR(256), category_id BIGINT, update_time DATETIME ) UNIQUE KEY(product_id) DISTRIBUTED BY HASH(product_id) BUCKETS 8 PROPERTIES ( replication_num 2, enable_unique_key_merge_on_write true );这里最关键的是enable_unique_key_merge_on_write要开 true这样才能在实时 upsert 场景下保证查询性能。如果不开Doris 会在读取时做 Merge性能会差很多。4.5 建立统一查询视图数据都接入后我会创建一个逻辑视图把历史订单、实时点击流、商品维表关联起来。业务想分析“每个商品今天的点击量、近 30 天销量和毛利”一个视图就能搞定。CREATE VIEW v_product_analysis AS SELECT p.product_id, p.product_name, today_clk.click_cnt AS today_click_cnt, his.order_cnt AS history_order_cnt, his.amount_sum AS history_amount_sum FROM dim_product p LEFT JOIN ( SELECT product_id, count(*) AS click_cnt FROM dwd_click_log WHERE click_time CURDATE() GROUP BY product_id ) today_clk ON p.product_id today_clk.product_id LEFT JOIN ( SELECT product_id, count(order_id) AS order_cnt, sum(amount) AS amount_sum FROM mv_order_daily WHERE order_date DATE_SUB(CURDATE(), INTERVAL 30 DAY) GROUP BY product_id ) his ON p.product_id his.product_id;分析师直接查询视图即可不需要知道底层是 Iceberg 还是 Kafka 数据。你也可以在 Doris 里给不同用户授权视图权限隐藏底层 Catalog 信息。4.6 常见的查询优化配置在真实业务里查询延迟和并发往往一起压过来。我常用的几个优化手段如下开启物化视图透明改写保证高频聚合查询命中物化视图。合理设置 Bucket 数和副本数数据量小时 bucket 过多会导致调度开销大数据量大时 bucket 过少会导致并行度不足。使用 Runtime Filter大表和小表 Join 时开启自动 Runtime Filter 能大幅减少数据传输量。控制外部表查询并发对湖表的查询不要开太高的并发否则对象存储请求限流会导致大面积报错。设置查询超时和资源组把报表查询和 Ad-hoc 查询放到不同资源组避免相互影响。5. 常见问题与排查技巧实录5.1 外表查询慢到底是慢在哪外部表查询慢的原因通常是元数据拉取慢、文件扫描慢、谓词下推失效三者之一。排查顺序建议先看 Profile。通过EXPLAIN或EXPLAIN ANALYZE查看执行计划重点看 Table Scan 节点检查numScanRange、numRowsScanned和filteredRows是否合理。如果扫描行数远大于实际返回行数说明过滤条件下推不彻底。还有一种情况是数据湖表文件数量太多比如几万个 1MB 小文件即使过滤条件再好每个文件都要打开开销依然很大。这种情况下我会在湖上做小文件合并或者在 Doris 中建立物化视图把明细归并成比较大的分段。5.2 文件缓存不生效怎么排查开启 File Cache 后如果第二次查询仍然慢大概率是缓存没有命中。常见原因有查询列太多导致缓存块被淘汰并发查询量太大缓存容量设置过小BE 节点重启后缓存被清理。你可以通过查询 Profile 中的 Cache 命中率指标确认。我一般会把file_cache_max_size_per_disk调大一些比如单盘 200GB 以上并确保查询模式相对固定。如果业务查询模式很发散缓存命中率很难提升这时候就别指望缓存直接上物化视图更靠谱。5.3 Routine Load 延迟越来越高延迟升高通常是两个原因Kafka 分区数少于导入并发设置或者单条消息处理耗时过长比如复杂 JSON 解析。我会先用SHOW ROUTINE LOAD;查看任务状态看lag和pendingTasks指标。如果 pendingTasks 很高说明 Doris 消费速度跟不上生产速度需要增加 BE 资源或者提高desired_concurrent_number同时确认目标表的写入没有成为瓶颈。目标表如果有过多小批量导入建议增大批次时间窗口减少文件版本数。5.4 Join 查询内存爆掉涉及数据湖大表和内部实时表的 Join最容易爆内存。我遇到过多次。解决思路有几种调整 Join 顺序用小表做驱动表开启 Runtime Filter如果是等值 Join考虑将大表按 Join 字段分桶让数据分布一致从而使用 Colocate Join 避免 Shuffle。Doris 的colocate_group配置需要两表的分桶列和桶数一致。若表已经建好可以在建表时指定也可以临时用SET enable_optimizer_reorder true;让优化器自动调整 Join 顺序。5.5 数据口径对不上这个往往不是 Doris 的问题而是同步逻辑的问题。比如实时表和历史表对“当天”的界定有的用 UTC 时间有的用东八区时间。我建议在接入层统一时间和时区字段并在视图或物化视图里明确时间口径。另外实时数据和湖上历史数据可能因为同步延迟存在重叠或缺口我会为实时数仓表加上业务日期字段查询时用 coalesce 或 union all 补全。6. 最后的一些实操心得这个架构落地到现在我最深的体会是Doris 作为数据湖加速、实时数仓、统一查询层的“高性能中枢”并不需要你把所有数据都搬到 Doris 里而是“有意为之”地把高频和实时部分放进来低频和历史部分留在湖里通过查询路由和物化视图让用户感受不到边界。这里面的分寸感很重要一开始我也犯过把所有湖表都同步到 Doris 内部的错误结果存储成本翻番查询却没有显著变快。后来我学会了一个原则明细数据按需物化汇总数据必须物化维表全量同步。如果你打算在自己的团队里引入 Doris我建议先不要追求大而全的架构先从最痛的一个业务场景切入把实时指标查询从原来的跑批改成 Doris 查询再加一个数据湖外部表做历史对比。跑通一条链路之后再逐步扩展 Catalog 和物化视图。Doris 配置灵活但坑也不少尤其是外部表元数据、文件缓存、实时导入并发这些细节每一项都需要生产环境中的真实数据来验证理论和参考架构都只是起点。
返回列表