Apache Doris 3.0:告别分库分表,构建统一实时数仓的架构实践 📅 2026/8/26 9:58:18 1. 项目概述从“拆”到“合”的架构演进如果你在数据团队待过几年肯定对“分库分表”这四个字有复杂的感情。早期业务跑得快数据库单表撑不住大家的第一反应就是“拆”。按用户ID哈希分128个库每个库再分256张表一套组合拳下来应用层要处理的路由逻辑复杂得像蜘蛛网。这还只是开始等到业务方要一个跨多个分片的全局报表或者想实时分析用户最近30天的行为路径时噩梦才真正降临。数据工程师要么得写复杂的聚合脚本在各个分库间游走抽取再合并要么就得再建一套专门用于分析的数据仓库ETL链路长得让人心焦数据延迟几个小时是家常便饭。这正是过去十年很多互联网公司数据架构的真实写照交易库OLTP和分析库OLAP被硬生生割裂中间隔着一条由定时任务和复杂脚本构成的“数据鸿沟”。而今天要聊的Apache Doris 3.0瞄准的就是这个核心痛点。它不再仅仅是一个传统的MPP分析型数据库而是试图成为一块能够直接承接高并发点查、支持实时数据更新、并能无缝进行复杂分析的“实时数仓基石”。尤其是在AI应用爆发的当下特征工程、模型训练、在线推理都对数据的实时性和一致性提出了前所未有的要求。一个还在跑批量T1的数据平台根本无法支撑起敏捷的AI迭代。Doris 3.0通过其全新的存算分离架构、更强的数据湖分析能力以及对向量化计算引擎的优化正在努力填平这道鸿沟让企业能够用一套系统同时应对“交易”和“分析”的挑战告别那个不停“分拆”和“搬运”数据的时代。2. 核心设计思路为何是“实时数仓基石”要理解Doris 3.0的定位我们得先拆解“实时数仓基石”这个说法。传统的数仓无论是基于Hadoop的离线数仓还是早期的一些MPP数据库其核心流程是“抽取-转换-加载”ETL。数据从业务系统产生经过一段时间的累积可能是分钟、小时甚至天级别被批量同步到数仓中进行处理和分析。这个模式在应对固定报表和历史分析时没问题但面对实时风控、实时推荐、动态定价等场景时就力不从心了。Doris 3.0的设计思路是推动数仓向“ELT”甚至“ETLT”演进。它强调数据的实时接入Ingestion和就地转换Transformation。其基石作用体现在三个层面2.1 统一的服务层入口过去应用开发可能需要面对多个数据端点写事务数据去MySQL写日志去Kafka查用户画像去HBase跑分析报表去Presto/ClickHouse。Doris 3.0通过支持多种数据模型聚合表、唯一表、重复表和多种导入方式Stream Load、Routine Load、Binlog Load试图提供一个统一的SQL入口。对于应用而言无论是写入一条订单记录还是查询一个复杂的用户行为漏斗理论上都可以通过同一个Doris集群和同一种SQL方言来完成。这极大地简化了数据架构的复杂度。2.2 融合的存储与计算范式“分库分表”本质是通过水平切分来分散存储和计算压力但这带来了数据分布的碎片化。Doris作为一款MPP数据库其底层数据存储本身就是分布式的它通过分区Partition和分桶Bucket机制在数据入库时就已经完成了有序的分布。计算时查询可以并行地在所有数据节点上执行由FEFrontend节点协调汇总。这种原生的分布式设计使得用户无需再关心“我的数据到底在哪个物理分片上”这种底层问题而是通过标准的SQL来操作逻辑上的一张表。Doris 3.0的存算分离架构进一步将这种范式深化计算节点可以按需弹性伸缩独立于存储节点为云原生部署和成本优化提供了可能。2.3 面向AI的数据服务能力AI时代的数据需求有两个鲜明特点一是对新鲜数据的渴求模型需要最新的特征二是对复杂数据类型如图、向量和计算模式如相似度搜索的支持。Doris 3.0在这两方面都做了重点增强。其极致的实时数据更新能力通过Unique Key或Merge-on-Write模型可以保证特征库的实时更新。同时社区正在积极集成向量搜索能力未来有望直接支持基于向量的近似最近邻ANN搜索这使得Doris不仅能做传统的结构化数据分析还能直接服务于AI场景下的特征检索和召回。3. 关键技术特性深度解析Doris 3.0并非一个凭空出现的新版本它是在2.0版本坚实基础上的一次重大演进。要理解它如何解决“分库分表”的痛点我们需要深入其几个关键的技术特性。3.1 存算分离架构Compute Node这是3.0版本最引人注目的特性。在传统架构中Doris的BEBackend节点同时承担数据存储和SQL计算双重职责。这带来了资源耦合的问题计算资源紧张时需要扩容但连带存储成本也上升了反之存储空间不足时扩容又会带来计算资源的浪费。3.0版本引入了独立的计算节点Compute Node。计算节点只负责执行查询计算不持久化存储数据。数据仍然存储在原有的BE节点或外部对象存储如S3、OSS中。这个架构带来了根本性的好处极致弹性在面对“双十一”、“618”等流量洪峰时可以快速、低成本地增加大量计算节点来应对突发的查询压力事后再缩容。计算节点的启动速度远快于传统的BE节点。成本优化存储和计算可以独立规划。可以将冷数据存储在更便宜的对象存储上而计算节点按需使用实现了存储成本与计算成本的双重优化。资源共享与隔离不同的业务线或租户可以共享同一份数据存储但使用独立的计算节点集群进行查询实现了资源的有效隔离与复用。注意存算分离的部署和运维模式与一体架构有较大差异。需要仔细规划网络带宽计算节点与存储节点间数据拉取、元数据缓存策略以及计算节点的调度策略避免因网络延迟导致查询性能下降。3.2 数据湖分析能力增强“湖仓一体”是当前大数据领域的主流方向。Doris 3.0极大地强化了对外部数据湖如Apache Hudi、Iceberg、Delta Lake的表格式支持。用户可以通过创建“多源目录”Multi-Catalog直接使用SQL查询存储在HDFS或对象存储上的数据湖表无需数据导入。这个特性直接击中了“数据孤岛”和“数据冗余”的痛点。很多企业的原始数据或轻度汇总数据已经存在于数据湖中传统做法是需要将这些数据再次导入到数仓才能进行高效分析。现在Doris可以直接对这些外部数据执行高性能查询。这意味着免搬迁分析可以直接对湖中数据做即席查询探索性分析无需等待漫长的ETL。统一元数据Doris可以同步HMSHive Metastore或其它Catalog的元数据提供统一的库表视图。增量与时间旅行查询对于支持快照的湖格式如Iceberg可以直接查询历史快照数据方便进行数据回溯和审计。3.3 向量化执行引擎优化查询性能是分析型数据库的立身之本。Doris很早就实现了向量化执行引擎而3.0版本在此基础上有持续优化。向量化执行的核心思想是摒弃传统的“一次处理一行数据”的模式改为“一次处理一批数据一个向量”充分利用现代CPU的SIMD单指令多数据流指令集进行并行计算。在3.0中优化可能深入到更具体的算子和场景例如聚合算子优化对于高基数分组聚合优化哈希表实现减少内存占用和CPU缓存失效。连接Join优化完善Colocate Join、Bucket Shuffle Join等分布式Join策略减少节点间数据Shuffle的网络开销。这对于替代需要跨分库分表做Join的业务场景至关重要。谓词下推与存储层协同更智能地将过滤条件下推到存储层或扫描层在读取数据时就尽可能过滤掉无关数据块减少后续计算的数据量。3.4 实时数据更新与部分更新在需要替代“分库分表”的OLTP场景时数据的实时更新能力是关键。Doris提供了两种模型来支持更新Unique Key模型定义主键对于相同主键的数据后写入的会自动覆盖先写入的Replace。这非常适合需要实时同步MySQL等业务库binlog的场景实现“实时数仓”的镜像。Merge-on-Write (MoW) 表这是2.0版本引入、在3.0中持续稳定的特性。它同样支持定义主键和实时UPSERT插入/更新。与Unique Key模型在查询时合并不同MoW表在数据写入时即完成合并牺牲一部分写入性能但换来了极致的查询性能无需运行时合并。这对于点查频繁的场景如用户画像实时查询非常友好。此外部分列更新功能允许只更新一行中的某几列而不是整行替换。这在更新宽表、且每次只更新少量字段的场景下能大幅减少IO和网络传输是面向OLTP场景的一个非常实用的优化。4. 从分库分表迁移到Doris的实操路径假设我们有一个用户订单系统最初基于MySQL并已按user_id进行了分库分表。现在希望迁移到Doris以解决跨分片查询复杂、分析能力弱的问题。以下是详细的迁移与实现路径。4.1 数据模型与表结构设计这是最关键的一步设计好坏直接影响后续性能和易用性。选择数据模型如果订单数据一旦生成不再变更如日志或仅追加选择Duplicate Key模型指定排序列如order_time写入性能最高。如果订单状态会更新如从“待支付”变为“已支付”且需要实时查询最新状态选择Unique Key模型主键为order_id或Merge-on-Write模型。通常读多写少、点查多的场景用MoW写频繁、对写入延迟敏感的场景用Unique Key。设计分区与分桶分区Partition按时间范围分区是最常见的做法例如按order_time字段每月一个分区。这可以方便地进行历史数据冷热分离删除旧分区或迁移到对象存储也利于查询时快速裁剪数据。分桶Bucket分桶是数据在分区内进一步水平切分和分布的方式。分桶列的选择至关重要应选择高频查询条件或Join条件的列。对于订单表user_id和order_id都是候选。如果查询经常按user_id查询其所有订单那么用user_id作为分桶列相同用户的订单会落在同一个桶内查询效率高。如果查询经常按order_id进行点查那么用order_id作为分桶列更合适。可以使用多个列的组合作为分桶列。分桶数量需要合理规划通常建议每个分桶的数据量在100MB-1GB之间。太少不利于并行太多会增加元数据开销。创建表示例CREATE TABLE order_analysis ( order_id BIGINT, user_id BIGINT, product_id INT, order_amount DECIMAL(12,2), order_status VARCHAR(20), province_code VARCHAR(10), order_time DATETIME, update_time DATETIME ) UNIQUE KEY(order_id) -- 使用Unique Key模型主键为order_id DISTRIBUTED BY HASH(user_id) BUCKETS 32 -- 按user_id哈希分桶32个桶 PARTITION BY RANGE(order_time) -- 按订单时间范围分区 ( PARTITION p202401 VALUES [(2024-01-01), (2024-02-01)), PARTITION p202402 VALUES [(2024-02-01), (2024-03-01)), PARTITION p202403 VALUES [(2024-03-01), (2024-04-01)) ) PROPERTIES ( replication_num 3, -- 副本数高可用 enable_unique_key_merge_on_write true -- 启用Merge-on-Write以优化点查性能 );4.2 数据实时同步方案如何将分散在数十个甚至上百个MySQL分片中的数据实时同步到Doris的一张表里使用 Canal / Debezium Routine Load在MySQL每个分片实例上部署Canal或Debezium来捕获Binlog。将Binlog消息发送到统一的Kafka Topic中。这里需要注意为了避免主键冲突最好在Kafka消息的Key或Value中包含源分片信息或者确保全局主键如order_id本身是全局唯一的。在Doris中为这个Kafka Topic创建一个Routine Load任务持续消费数据并导入到order_analysis表中。由于Doris表的主键是order_id来自不同分片的相同订单更新会在Doris内部自动合并。使用 Flink CDC这是目前更流行和强大的方案。使用Flink SQL直接创建MySQL CDC源表捕获所有分片然后通过Flink进行简单的ETL如字段转换、过滤后用Flink-Doris-Connector写入Doris。优势在于Flink提供了强大的流处理能力可以在入库前进行更复杂的数据清洗、聚合或打宽。同时Flink CDC能更好地处理Schema变更。双写与灰度迁移在迁移初期可以采用应用双写策略应用同时写入旧的MySQL分片和新的Doris表。这可以用于数据比对和验证。逐步将读流量切到Doris先切离线分析查询再切简单的点查最后切复杂的OLAP查询。整个过程需要密切监控Doris集群的负载和查询延迟。4.3 查询服务改造与优化数据同步完成后业务查询需要从面向分库分表的复杂逻辑改为面向Doris的统一SQL。重构数据访问层DAL将原来包含分片路由逻辑的代码替换为标准的JDBC或ORM方式连接Doris。查询语句变得异常简单例如查某个用户最近一个月的订单SELECT * FROM order_analysis WHERE user_id 123456 AND order_time 2024-03-01;利用物化视图预聚合对于频繁查询的聚合指标如“每个省份的日销售总额”可以创建物化视图来加速。CREATE MATERIALIZED VIEW province_daily_sales_mv AS SELECT province_code, DATE(order_time) as dt, SUM(order_amount) as total_sales FROM order_analysis GROUP BY province_code, DATE(order_time);当查询命中该物化视图时Doris会自动路由并从预聚合的结果中读取速度极快。复杂查询实践过去需要业务代码拼凑的跨分片复杂查询现在可以直接用SQL完成。例如查询购买过某热门商品且订单金额大于1000元的用户分布SELECT province_code, COUNT(DISTINCT user_id) as user_count FROM order_analysis WHERE product_id 8888 AND order_amount 1000 GROUP BY province_code ORDER BY user_count DESC;5. 生产环境部署与运维核心要点将Doris用于核心生产环境尤其是在替代原有数据库的场景下稳定的运维至关重要。5.1 集群规划与硬件选型FE节点负责元数据管理、查询规划、集群调度。属于控制面对CPU和内存要求较高但IO压力小。建议使用高配虚拟机或物理机至少3个节点组成高可用1个Leader2个Follower。生产环境务必部署3个及以上Follower。BE节点在存算一体模式下承担数据存储和计算。需要均衡的CPU、内存、磁盘和网络资源。CPU建议现代多核处理器主频越高向量化计算收益越大。内存越大越好用于查询计算、排序、聚合等。建议至少64GB起步。磁盘使用SSDNVMe最佳。BE数据目录建议使用多块盘做RAID0或直接配置为多数据目录以提升IO吞吐。切忌使用机械硬盘会严重成为性能瓶颈。网络万兆网卡是标配节点间数据传输和Shuffle对带宽和延迟非常敏感。计算节点CN如果采用存算分离架构CN节点需要强大的CPU和充足的内存磁盘容量要求不高主要用于临时数据。可以选用计算优化型实例。5.2 关键的配置调优max_query_memory_limit限制单个查询在BE上能使用的最大内存。防止大查询拖垮整个节点。query_timeout和load_timeout设置合理的超时时间避免慢查询或导入任务长期占用资源。压缩算法表数据的压缩算法在表属性中设置compression推荐使用lz4它在压缩比和解压速度之间取得了很好的平衡对查询性能友好。索引Doris内置了智能的前缀索引基于排序列。在CREATE TABLE时将最常作为查询条件的列放在DUPLICATE KEY或UNIQUE KEY的前面能有效加速查询。对于非前缀的常用过滤列可以考虑使用Bloom Filter索引。5.3 监控与告警体系必须建立完善的监控核心指标包括集群健康度FE、BE、CN节点存活状态。资源使用率CPU、内存、磁盘使用率、网络IO。特别是BE节点的磁盘使用率需设置水位线告警如85%。查询性能query_latency查询延迟、qps每秒查询数、慢查询数量。分析慢查询日志是性能优化的关键入口。导入性能load_rpc_rate导入RPC速率、导入任务队列长度、失败率。副本状态表副本的健康状态是否有副本缺失或损坏。推荐使用Prometheus Grafana来搭建监控看板社区提供了标准的导出器metrics exporter。6. 常见问题与故障排查实录在实际使用中你肯定会遇到各种问题。以下是一些典型场景及排查思路。6.1 写入速度变慢或失败现象Stream Load或Routine Load任务耗时变长甚至失败。排查检查BE节点磁盘IO使用iostat -x 1查看磁盘使用率%util和响应时间await。如果持续接近100%说明磁盘已达瓶颈需要考虑扩容或增加BE节点。检查内存导入任务会消耗内存。查看BE日志是否有Memory limit exceeded相关错误。可以适当调大load_process_max_memory_limit_percent参数需谨慎。检查副本同步如果写入涉及多个副本网络延迟或某个BE节点繁忙会导致副本同步慢。观察tablet_commit_count等指标。小文件问题频繁的小批量导入会产生大量小文件影响后续查询和合并Compaction性能。应适当调大导入的批量大小或调整Compaction策略如cumulative_compaction_num_threads_per_disk。6.2 查询延迟高或不稳定现象同一个SQL有时快有时慢或整体比预期慢。排查首先查看查询计划在Doris Web UI或通过EXPLAIN命令查看SQL的执行计划。重点关注分区裁剪是否有效过滤了分区PartitionRange是否缩小到了目标范围。索引命中PREDICATES部分是否下推了谓词是否使用了合适的索引。Join类型是Broadcast Join、Shuffle Join还是Colocate/Bucket Shuffle Join不合理的Join类型会导致巨大的网络开销。确保Join键是分桶列以实现本地Join。检查集群负载查询期间是否有其他重计算任务如Compaction、其他大查询在运行导致资源争抢。监控CPU和内存使用率。检查统计信息Doris的CBO基于成本的优化器依赖统计信息。如果表数据量发生重大变化如增长10倍但统计信息未更新优化器可能会生成次优计划。定期执行ANALYZE TABLE更新统计信息。热点问题如果数据分布严重倾斜例如某个分桶的数据量是其他的几十倍会导致处理该分桶的BE节点成为瓶颈。需要重新审视分桶列的选择和分桶数。6.3 内存不足OOM问题现象查询失败BE日志出现Memory exceed limit或BE进程崩溃。排查与解决分析查询通常是涉及大表全表扫描、高基数分组聚合DISTINCT COUNT、或非等值Join如导致中间结果集膨胀。尝试优化SQL增加过滤条件或分阶段查询。调整参数全局层面适当增加mem_limitBE总内存限制和query_mem_limit单个查询内存限制但不要超过物理内存。会话层面对于确定需要大量内存的查询可以在Session中临时设置set exec_mem_limitxxxxx;。启用Spill to Disk对于排序、聚合等操作如果内存不足可以启用落盘功能。设置spill_mode为auto或force并指定spill_storage_root路径。这用磁盘空间换取了内存避免OOM但会降低查询速度。6.4 数据不一致或丢失现象查询结果与源数据对不上或部分数据查不到。排查检查导入任务状态确认Routine Load或Stream Load任务状态是否为RUNNING且没有持续的ERROR。检查错误信息常见原因是数据格式不匹配或主键冲突。检查副本一致性使用ADMIN SHOW REPLICA STATUS命令查看表的副本状态确认所有副本的Version是否一致。不一致可能需要触发副本修复。检查数据版本Doris的多版本并发控制MVCC机制下过时的查询可能读到旧版本数据。确认查询是否使用了正确的时间范围。核对数据源如果是CDC同步核对Kafka中消息的偏移量是否正常Flink作业是否有异常重启导致数据重复或丢失。迁移到Doris是一个系统工程从模型设计、数据同步到查询改造和运维每一步都需要仔细考量。它带来的收益是巨大的简化到极致的架构、统一的数据视图、强大的实时分析能力。但与此同时你也需要接受一种新的思维模式——从“如何分拆数据”转向“如何更好地组织数据以利用分布式并行计算”。这个过程可能会有阵痛但当你看到那些曾经需要几个小时才能跑出来的跨分片报表现在秒级返回时你会觉得这一切都是值得的。