尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

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

Apache Doris + SelectDB:构建AI时代实时分析平台的三大核心范式

Apache Doris + SelectDB:构建AI时代实时分析平台的三大核心范式 1. 项目概述当实时分析遇上AI我们如何重新定义范式最近几年数据领域最火的两个词一个是“实时”另一个就是“AI”。从业务侧看老板们不再满足于T1的报表他们想知道“此刻”发生了什么从技术侧看大模型和AI应用井喷对数据的实时处理、低延迟查询和向量检索能力提出了前所未有的要求。传统的离线数仓和批处理架构在应对这些新场景时常常显得力不从心。正是在这样的背景下“Apache Doris SelectDB”这个组合开始频繁出现在技术社区的讨论中它被很多人视为构建下一代实时分析平台的关键拼图。简单来说Apache Doris是一个开源的、高性能的MPP分析型数据库以其极致的查询速度和易用性著称。而SelectDB则是基于Apache Doris内核的商业化产品与服务它在提供企业级功能、稳定性保障和云原生部署体验的同时也深度参与了Doris社区的建设与创新。这个组合的核心价值就在于它试图为AI时代的实时数据分析提供一套全新的、高效的“玩法”或曰“范式”。我接触Doris和SelectDB有段时间了从早期的单表查询优化到后来支撑起公司核心的实时数仓和用户行为分析平台再到最近开始尝试对接AI应用做向量检索整个过程踩过不少坑也积累了一些心得。今天我就想结合自己的实践经验来聊聊我所理解的、由“Apache Doris SelectDB”所定义的AI时代实时分析的三大核心范式。这不仅仅是工具介绍更是一套从架构设计到应用落地的完整思路。2. 范式一流批一体的统一数据服务层第一个范式也是最基础、最广泛应用的范式就是构建一个流批一体的统一数据服务层。它的核心目标是打破实时数据流与历史批量数据之间的壁垒让业务用一个接口、一套语法就能同时查询最新的实时数据和沉淀的历史数据且保证毫秒到秒级的响应。2.1 为什么需要流批一体在传统架构里实时数据和离线数据往往是两套系统。实时数据可能走Kafka - Flink - 实时OLAP库如Druid, ClickHouse的链路用于监控和实时大盘离线数据则走HDFS - Hive/Spark - 离线OLAP库或数据服务层的链路用于复杂的报表和深度分析。这种架构带来的问题显而易见数据口径不一致两套计算逻辑极易导致同一个指标在实时和离线报表中数值对不上引发信任危机。开发运维成本高需要维护两套技术栈开发人员要写两套代码运维要管两套集群。使用体验割裂分析师需要知道数据在哪套系统用不同的查询方式无法进行跨实时与历史数据的关联分析。Apache Doris通过其独特的“明细模型”、“聚合模型”和“更新模型”配合强大的实时数据接入能力为流批一体提供了优雅的解决方案。2.2 核心实现实时数据接入与数据模型设计实时数据接入这是流批一体的入口。Doris支持多种实时数据摄入方式最常用的是通过Routine Load任务从Kafka持续消费数据。我通常会这样创建一个任务CREATE ROUTINE LOAD db1.kafka_load ON table1 COLUMNS TERMINATED BY “,”, COLUMNS (k1, k2, v1, v2) PROPERTIES ( “desired_concurrent_number”“3” “max_batch_interval”“20” “max_batch_rows”“300000” “max_batch_size”“209715200” ) FROM KAFKA ( “kafka_broker_list”“broker1:9092,broker2:9092” “kafka_topic”“my_topic” “property.group.id”“doris_consumer_group” );这里有几个关键参数的经验之谈desired_concurrent_number并发数通常设置为Kafka topic的分区数以实现并行消费提升吞吐。max_batch_interval和max_batch_rows/size控制微批处理的节奏。间隔太短会频繁提交产生小文件影响查询性能间隔太长则数据延迟高。在生产中我们根据数据流量通常将间隔设置在10-30秒行数在20-50万之间做平衡。严格模式strict_mode对于数据质量要求高的场景建议开启。它会在遇到类型转换错误或空值违反非空约束时暂停任务防止脏数据污染。数据模型设计这是决定查询效率和资源消耗的关键。对于流批一体场景明细模型Duplicate适用于需要保留原始明细、进行任意维度下钻分析的场景如用户行为日志、交易流水。它是流数据最自然的落地方式。聚合模型Aggregate适用于指标汇总场景如PV、UV、GMV。Doris会在数据摄入时进行预聚合极大提升sum count min max等查询的速度。这是实现“实时汇总”能力的核心。更新模型Unique适用于状态表、维度表需要按主键更新。比如用户属性表最新的用户画像数据会覆盖旧数据。实操心得模型选择与分区键设计不要一味追求聚合模型。如果业务需要频繁查询明细或者聚合维度经常变化明细模型加物化视图可能是更好的选择。分区键Partition Key和分桶键Bucket Key的设计至关重要。我们通常按时间如dt字段做分区实现数据生命周期管理轻松删除旧分区按查询最频繁的维度如user_id,product_id做分桶并结合Bloom Filter索引能极大提升点查和范围查询的效率。2.3 统一查询SQL的力量当实时流数据通过Routine Load持续写入历史批量数据也可以通过Broker Load或Spark-Doris-Connector一次性或周期性导入后业务侧看到的就是一张完整的表。无论是BI工具、数据API还是即席查询只需要使用标准的SQL无需关心数据是刚刚从Kafka进来还是昨天从HDFS导入的。例如一个查询最近7天实时销量并同比上周的复杂SQL可以一气呵成。这种范式将数据架构极大简化从原来的“Lambda架构”进化到“Kappa架构”的简化版降低了至少50%的开发和运维复杂度是我们团队构建数据中台实时数据服务的基石。3. 范式二高并发点查与极速宽表分析第二个范式聚焦于具体的查询性能尤其是在两种典型但需求迥异的场景下面向C端应用的高并发、低延迟的点查询和面向内部分析的极速、复杂的宽表关联查询。Doris通过不同的技术特性同时在这两个战场表现出色。3.1 高并发点查应对千万级QPS的挑战点查即基于主键如订单ID、用户ID快速查询一行或少量行数据。这在电商订单详情、用户中心、实时风控等场景下需求巨大要求毫秒级响应且能承受极高的并发。Doris满足点查的核心依靠两点前缀索引和内存化优化。前缀索引Short Key IndexDoris在每间隔若干行数据默认1024行生成一个索引项存储的是该行数据前36字节的哈希值。当查询条件包含这些前缀列时能快速定位到数据块避免全表扫描。内存化优化PageCacheDoris利用操作系统的Page Cache缓存数据块热点数据查询速度极快。Tablet本地缓存对于更新不频繁的维度表可以将其设置为“内存表”通过storage_medium MEMORY全量驻留内存实现微秒级响应。配置示例与压测经验 为了支撑高并发点查我们在表设计上和集群配置上做了大量优化-- 建表示例突出点查优化 CREATE TABLE user_profile ( user_id BIGINT NOT NULL name VARCHAR(50) age INT city VARCHAR(20) last_login DATETIME ) UNIQUE KEY(user_id) -- 使用更新模型user_id为主键 DISTRIBUTED BY HASH(user_id) BUCKETS 32 -- 按user_id分桶保证点查落在一个bucket内 PROPERTIES ( “replication_num” “3” “storage_medium” “SSD” -- 对于极热表可以尝试 “storage_medium” “MEMORY” “light_schema_change” “true” -- 开启轻量Schema变更避免在频繁加列时重导数据 );在压测时我们使用sysbench或自定义脚本模拟并发请求。一个关键的发现是连接池的管理至关重要。应用端需要使用高效的连接池如HikariCP并合理设置最大连接数和超时时间避免连接数暴涨拖垮FEFrontend。同时通过部署多个FE节点并配置负载均衡可以有效分散连接压力。3.2 极速宽表分析告别“大宽表”与“星型模型”的纠结传统上为了应对复杂的多表关联查询如星型模型要么提前进行大量的ETL加工成“大宽表”牺牲灵活性和实时性要么让OLAP引擎执行昂贵的运行时Join速度堪忧。Doris提出了一个更优解预物化视图和Colocate Join。预物化视图你可以针对常用的、固定的关联查询模式创建物化视图。Doris会在底层自动维护这个预计算好的结果集。查询时优化器会自动路由到物化视图速度极快。这相当于按需、自动构建的“宽表”比手动维护ETL任务灵活得多。Colocate Join对于无法预知的关联查询Doris的Colocate Join功能可以在建表时指定相关的表采用相同的分桶方式和副本分布。这样关联计算时数据无需在网络间Shuffle直接在本地完成性能提升一个数量级。场景对比 假设有订单事实表和用户维度表。老方法大宽表每天凌晨ETL任务跑一小时把用户属性拼接到订单表生成一张巨大的宽表。分析师无法查询到最近一小时的订单且用户属性更新有延迟。Doris方法对“按用户城市分析订单总额”这个固定报表创建一个聚合物化视图。对即席查询在建表时让订单表按user_id分桶和用户表同样按user_id分桶设置为Colocate Group。 结果是固定报表亚秒级响应即席的多表关联查询也比传统方式快5-10倍。我们团队的一个核心宽表分析场景在切换到Colocate Join后查询耗时从平均20秒降到了2秒以内。这个范式让数据分析师和业务系统都能“鱼与熊掌兼得”既享有宽表的查询速度又保留了模型的灵活性和数据的实时性。4. 范式三AI-Native的数据智能底座第三个范式是最前沿、也最能体现“AI时代”特征的即让实时分析数据库本身具备AI原生的能力成为AI应用的高效数据底座。这主要体现在两个方面高性能向量检索和内置机器学习函数。4.1 向量检索让数据库“理解”非结构化数据大模型应用如智能问答、推荐系统、图像检索的核心需求之一是根据文本、图像等非结构化数据生成的“向量”Embedding快速找到最相似的条目。这要求数据库能高效处理向量数据的存储和近似最近邻ANN搜索。Apache Doris从2.0版本开始正式支持了向量化索引。其核心流程是数据写入应用端将原始内容如商品描述、文章通过模型如OpenAI的text-embedding-ada-002转化为固定维度的浮点数向量连同原始数据一起写入Doris。索引构建Doris支持在向量列上创建INDEX目前主要支持FLAT暴力计算精度100%和IVF_FLAT倒排文件查询更快精度略有损失等索引类型。相似度查询使用DOT_PRODUCT点积或COSINE余弦相似度等函数进行相似度计算并结合索引快速返回Top-K结果。实操示例-- 1. 创建包含向量列的表 CREATE TABLE article_embeddings ( article_id BIGINT title VARCHAR(500) content TEXT embedding ARRAYFLOAT -- 假设是1536维的向量 ) DUPLICATE KEY(article_id) DISTRIBUTED BY HASH(article_id) BUCKETS 16 PROPERTIES (“replication_num” “1”); -- 2. 在向量列上创建IVF_FLAT索引假设有100万数据建256个聚类中心 CREATE INDEX embedding_idx ON article_embeddings(embedding) USING IVF_FLAT PROPERTIES(“nlist” “256”); -- 3. 进行相似度查询找到与给定向量最相似的10篇文章 SELECT article_id title DOT_PRODUCT(embedding [0.1 0.2 ... 0.5]) AS similarity FROM article_embeddings ORDER BY similarity DESC LIMIT 10;注意事项向量索引的调优nlist参数是IVF索引的核心它定义了聚类中心的数量。值越大搜索精度越高但建索引和搜索的成本也越高。通常建议在sqrt(N)N为总数据量附近调整。对于千万级数据nlist2048或4096是常见的起点。向量维度不宜过高通常1536维OpenAI ada模型或768维常用Sentence-BERT模型是平衡点。维度越高存储和计算成本呈线性增长。目前Doris的向量索引在批量导入数据后需要手动触发BUILD INDEX或等待后台合并时构建对于实时写入的场景需要关注索引构建的延迟。4.2 内置机器学习函数与联邦查询除了向量检索Doris还在向更广泛的AI能力演进。例如通过内置的机器学习函数用户可以直接用SQL调用一些简单的模型进行预测比如使用LOGISTIC_REGRESSION_PREDICT函数进行在线推理无需将数据导出到专门的ML系统。更强大的能力在于联邦查询Federation。Doris可以通过Multi-Catalog功能直接对接外部数据源如Hive、Iceberg、Elasticsearch、MySQL甚至通过JDBC连接其他数据库。这意味着你可以写一条SQL同时关联查询Doris本地表里的实时指标、Hive里的历史数据、ES里的日志文本并将结果统一返回。这对于构建跨数据源的AI特征平台至关重要数据科学家无需在不同系统间搬运数据可以一站式完成特征抽取和样本构建。踩坑实录向量检索的性能陷阱我们早期在测试向量检索时直接对百万级数据做全表DOT_PRODUCT计算查询耗时超过10秒完全不可用。创建了IVF索引后性能提升到200毫秒以内。但另一个坑是内存高并发进行向量检索时索引和数据会被加载到内存如果并发数过高容易导致BEBackend节点内存溢出。我们的解决方案是合理规划BE节点内存预留足够空间给向量检索。在应用层做查询限流和队列管理。考虑将向量数据单独放在专用的、内存更大的节点组Tag中通过Doris的资源隔离功能实现物理隔离。这个范式正在将Doris从一个单纯的OLAP数据库转变为一个支持AI工作负载的“智能数据平台”这也是SelectDB商业版重点发力的方向提供了更稳定的向量索引、云原生的弹性伸缩以及企业级的运维支持。5. 实战从零构建一个实时AI推荐系统数据层理论说了这么多我们来看一个综合性的实战案例如何利用上述三大范式构建一个支持实时AI推荐系统的数据层。这个系统需要1实时处理用户点击流2存储物品和用户的向量化特征3支持低延迟的特征检索和模型推理。5.1 架构设计与数据流我们的架构如下图所示文字描述数据源用户行为日志Kafka、物品元数据MySQL、离线训练好的用户/物品向量HDFS。实时接入层使用Doris的Routine Load从Kafka消费用户点击、搜索行为写入user_behavior明细表范式一。同时通过DataX或Flink CDC将MySQL的物品元数据同步到item_meta表。向量存储与检索层将离线训练好的物品向量通过Spark-Doris-Connector批量导入item_vectors表并创建IVF索引范式三。实时用户向量可以由实时模型推理后写入或通过联邦查询从模型服务中获取。统一服务层推荐引擎RecSys通过JDBC连接Doris。当需要为某个用户做推荐时首先从user_behavior表中实时获取该用户最近的行为范式一点查。然后获取该用户的向量在item_vectors表中进行近似最近邻搜索找到最相似的N个候选物品范式三。最后将这些候选物品的ID与item_meta表进行关联获取详细信息可能还会关联实时聚合表如物品实时热度item_hot_agg由聚合模型实现进行打分排序范式二Colocate Join确保关联性能。结果推荐引擎综合各种特征和分数生成最终的推荐列表。整个数据查询过程在百毫秒内完成。5.2 关键DDL与配置片段-- 1. 用户行为明细表 (范式一流批一体入口) CREATE TABLE dwd.user_behavior ( user_id BIGINT item_id BIGINT behavior_type VARCHAR(10) -- ‘click’ ‘like’ ‘buy’ ts DATETIME dt DATE -- 用于分区 ) DUPLICATE KEY(user_id item_id ts) PARTITION BY RANGE(dt)() -- 动态分区后续添加 DISTRIBUTED BY HASH(user_id) BUCKETS 32 PROPERTIES ( “replication_num” “3” “dynamic_partition.enable” “true” “dynamic_partition.time_unit” “DAY” “dynamic_partition.start” “-7” -- 保留最近7天 “dynamic_partition.end” “3” ); -- 2. 物品向量表 (范式三AI原生) CREATE TABLE rec.item_vectors ( item_id BIGINT embedding ARRAYFLOAT update_time DATETIME ) UNIQUE KEY(item_id) DISTRIBUTED BY HASH(item_id) BUCKETS 16 PROPERTIES (“replication_num” “3”); -- 创建向量索引 CREATE INDEX idx_vec ON rec.item_vectors(embedding) USING IVF_FLAT PROPERTIES(“nlist” “1024”); -- 3. 物品实时热度聚合表 (范式二极速聚合) CREATE TABLE dws.item_hot_agg ( item_id BIGINT dt DATE hour DATETIME click_count BIGINT SUM -- 聚合模型自动求和 unique_user_count HLL_UNION HLL -- 使用HLL进行UV近似统计 ) AGGREGATE KEY(item_id dt hour) PARTITION BY RANGE(dt)() DISTRIBUTED BY HASH(item_id) BUCKETS 16 PROPERTIES ( “replication_num” “3” -- 设置数据过期时间例如保留7天详细数据 “storage_cooldown_time” “9999-12-31 23:59:59” -- 实际需根据分区设置 ); -- 此表的数据由Flink或Doris的物化视图从user_behavior表实时聚合而来。 -- 4. 为item_vectors和item_hot_agg设置Colocate Group加速关联 (范式二) ALTER TABLE rec.item_vectors SET (“colocate_with” “hot_group”); ALTER TABLE dws.item_hot_agg SET (“colocate_with” “hot_group”);5.3 性能优化与踩坑点在这个项目中我们遇到了几个典型问题向量索引构建慢初始导入千万级向量时构建索引耗时数小时。解决方案是调大build_heap_memory_limit参数增加索引构建的内存并采用分批导入、分批建索引的方式。实时聚合资源消耗大item_hot_agg表如果按分钟聚合数据量会爆炸。我们最终按小时聚合并通过ROLLUP物化视图预先计算天级别的聚合数据空间换时间。联邦查询效率尝试过用Doris直接联邦查询TensorFlow Serving的接口获取实时用户向量延迟不稳定。后来改为由推荐引擎自行调用模型服务获取向量再将向量作为参数传给Doris做物品检索架构更清晰稳定性更好。这个实战案例充分展示了三大范式如何协同工作将一个复杂的、多数据源、高实时性要求的AI推荐系统数据层变得清晰、高效和可维护。6. 常见问题与排查技巧实录在实际运维和开发过程中总会遇到各种问题。下面是我总结的一些常见问题及其排查思路希望能帮你少走弯路。6.1 数据写入问题问题现象可能原因排查步骤与解决方案Routine Load任务持续PAUSED1. 数据格式错误如列数不匹配。2. 严格模式strict_mode下遇到空值或类型错误。3. Kafka分区偏移量问题。1.SHOW ROUTINE LOAD\G查看错误信息。2. 检查错误样例数据调整COLUMNS映射或关闭strict_mode。3. 检查Kafka集群和Topic状态尝试重置offset。数据写入延迟高1. BE节点写入压力大CPU/IO高。2. 单批次数据量过大或过小。3. 副本同步慢。1. 监控BE节点资源使用率考虑扩容或均衡负载。2. 调整Routine Load的max_batch_interval和max_batch_rows找到吞吐和延迟的平衡点。3. 检查网络和磁盘性能SHOW PROC ‘/backends’查看副本健康状况。Broker Load导入HDFS数据失败1. HDFS文件路径或权限错误。2. Broker节点网络不通或配置错误。3. 文件格式不匹配。1. 在Broker节点上用hdfs dfs -ls命令手动测试路径。2. 检查broker_load的WITH BROKER配置确保Broker名称正确且节点存活。3. 确认FORMAT AS指定的格式如Parquet ORC与文件实际格式一致。6.2 查询性能问题问题现象可能原因排查步骤与解决方案简单点查变慢1. 前缀索引未命中导致全表扫描。2. 数据分布严重倾斜导致某个Bucket压力过大。3. FE或BE节点负载过高。1. 使用EXPLAIN查看执行计划确认是否使用了索引。2. 检查表的分桶键选择是否合理SHOW DATA SKEW查看数据倾斜情况。3. 监控集群负载考虑增加节点或优化查询并发。关联查询Join慢1. 未使用Colocate Join产生了网络Shuffle。2. 右表过大导致Broadcast Join负担重。3. Join条件上缺乏索引或分区裁剪。1. 使用EXPLAIN查看Join类型确认是否为ColocateJoin。2. 对于大表关联尝试调整exec_mem_limit增加内存或使用SHUFFLEHint强制分桶Join。3. 确保关联键是分桶键或建有索引并利用分区条件减少数据量。内存不足Out of Memory1. 单个查询处理的数据量过大。2. 并发查询过多总量超出限制。3. 向量检索等内存密集型操作并发高。1. 优化SQL增加过滤条件避免SELECT *。2. 通过SET exec_mem_limit限制单查询内存通过资源标签Resource Tag隔离不同业务负载。3. 对于向量检索限制并发数或使用专用高内存节点组。6.3 运维与集群问题磁盘空间告警Doris的数据删除是标记删除需要通过ALTER TABLE … COMPACT进行压缩Compaction才能真正释放空间。定期检查表的数据版本对版本数过多SHOW TABLET中VersionCount过大的表进行手动Compaction。同时合理设置动态分区过期策略自动清理旧数据。FE元数据故障FE负责元数据管理和查询规划是高可用关键。务必部署至少3个Follower组成集群并定期备份元数据mysqldump备份fe/meta目录对应的数据库。如果Leader FE宕机集群会自动选举但需确保客户端配置了多个FE地址实现重连。BE节点扩容后数据不均衡新增BE节点后旧表的数据不会自动迁移。需要使用ADMIN REPAIR TABLE … REBALANCE命令手动触发表级别的副本均衡或者使用ADMIN SET REPLICA STATUS命令进行更精细的调整。一个宝贵的调试技巧当遇到难以理解的慢查询时除了EXPLAIN一定要用EXPLAIN ANALYZE。它会实际执行查询所以最好在测试环境进行并输出每个执行环节的详细耗时能精准定位到是扫描数据慢、网络传输慢还是聚合计算慢是性能调优的终极利器。7. 总结与个人体会回顾这三大范式本质上是从不同维度解决了AI时代数据应用的痛点流批一体解决了数据“时效”与“统一”的矛盾让实时与历史数据无缝融合高并发点查与极速分析解决了数据“速度”与“复杂度”的矛盾让系统既能扛住洪峰流量又能进行深度洞察AI-Native能力则解决了数据“形态”与“智能”的矛盾让非结构化的向量数据也能被高效管理和检索。从我个人的使用体验来看Apache Doris SelectDB这个组合最大的优势在于“一体化”和“简洁性”。它试图用一个系统覆盖从数据接入、存储、加工到服务、分析乃至AI集成的全链路极大地简化了技术架构。对于中小型团队或需要快速迭代的业务来说这意味着更少的运维负担、更低的开发成本和更快的需求响应速度。当然它并非银弹。在超大规模数据PB级以上的纯离线复杂ETL场景可能还是Spark/Flink更专业在超高性能的键值查询场景专门的KV数据库仍有优势。但在实时分析、湖仓一体、数据服务以及新兴的AI数据底座这些交汇地带Doris展现出了强大的竞争力和生命力。最后分享一个小心得学习Doris一定要动手实践。它的很多特性比如物化视图的自动路由、Colocate Join的性能提升、向量索引的效果光看文档是很难有深刻体会的。不妨就从搭建一个单机测试环境开始用自己业务的一小部分真实数据跑一跑感受一下它如何重新定义你处理数据的方式。毕竟在数据领域没有什么比“跑起来”更有说服力了。
返回列表