ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Doris与SelectDB:AI时代实时分析的三大范式与实践

Apache Doris与SelectDB:AI时代实时分析的三大范式与实践 1. 从“实时数仓”到“AI时代分析”一个被重新定义的战场如果你在过去几年里关注过大数据技术栈的演进那么对“实时数仓”这个词一定不陌生。从早期的T1报表到后来的小时级、分钟级延迟再到秒级甚至毫秒级的实时分析这几乎是所有数据团队追求的目标。Apache Doris作为一款开源的MPP分析型数据库正是在这个浪潮中凭借其极致的查询性能和易用性成为了许多企业实时数仓的核心引擎。而SelectDB作为基于Apache Doris内核的商业化公司则进一步推动了其在企业级市场的应用。但今天当“AI时代”成为所有技术讨论的默认前缀时我们谈论的“分析”已经发生了根本性的变化。它不再仅仅是给业务人员看一张实时更新的销售大屏或者让风控系统在几百毫秒内拦截一笔可疑交易。AI驱动的分析意味着数据系统需要直接与模型训练、推理、Agent决策等环节深度耦合。数据流的终点从一个BI图表变成了一个AI模型的参数或者一个自主智能体的“思考”依据。这就是“Apache Doris SelectDB定义 AI 时代实时分析的三大范式”这个标题背后真正的张力所在。它宣告了一个转变我们熟悉的那个“实时数仓”Doris正在进化成一个面向AI的“实时数据智能平台”。这三大范式不是简单的功能叠加而是从架构理念到应用场景的体系化重构。接下来我将结合一线的观察和实践为你拆解这三大范式究竟是什么以及它们如何在实际项目中落地重新定义我们处理数据价值的方式。2. 范式一从“数据服务”到“模型特征”的实时供给在传统的机器学习流水线中特征工程是一个离线的、批处理的过程。数据工程师每天定时跑ETL任务生成特征表然后由算法工程师抽取到模型中进行训练或推理。这个过程的延迟通常是小时级甚至天级。在推荐、风控、实时定价等场景下这种延迟是无法接受的。用户点击一个商品后推荐系统需要在毫秒级内根据用户最新的行为例如刚刚浏览了哪些同类商品、停留了多久更新推荐列表这就要求特征必须是实时更新的。这就是第一个范式要解决的核心问题实现高吞吐、低延迟的实时特征供给。Apache Doris在这一范式中扮演的角色从一个批处理后的数据查询引擎转变为一个在线的特征存储与计算中心。2.1 架构变革流批一体与主键模型要实现实时特征首先需要改变数据入库的方式。过去我们可能通过Flink等流处理引擎进行一些聚合后再批量写入Doris。现在更优的路径是利用Doris本身对高频数据写入的支持能力。SelectDB在Apache Doris内核基础上强化了其Unique Key主键模型和Partial Update部分列更新能力。对于用户画像、商品特征这类以唯一ID为主键的表我们可以直接以流的方式通过Routine Load对接Kafka或通过Flink-Doris-Connector持续写入数据。当同一主键的新数据到来时Doris会自动进行更新。这意味着特征表始终保持着最新状态。例如一个用户实时行为特征表表结构为主键是user_id列包括last_click_time、last_viewed_item、session_duration等。用户每次点击一条包含user_id和last_click_time可能还有其他字段的记录就会实时写入。Doris会根据user_id找到原有行并只更新last_click_time这个字段其他字段保持不变。这个过程是毫秒级的。实操心得在设计实时特征表时务必合理定义主键。主键的粒度决定了更新的粒度。例如如果你需要更新的是“用户-商品”维度的特征那么主键就应该是(user_id, item_id)的组合。同时要谨慎选择允许更新的字段避免不必要的写放大。2.2 查询优化应对高并发点查特征供给的另一个挑战是查询模式。模型推理服务通常需要在一瞬间拉取成千上万个user_id对应的上百个特征列这是一种典型的高并发、低延迟的点查场景。这与Doris擅长的海量数据复杂分析OLAP场景有所不同。为此SelectDB优化了Doris的查询引擎和存储引擎以更好地支持这种“大宽表点查”索引优化除了主键索引对经常作为查询条件的字段如city,age_group建立前缀索引或Bloom Filter索引加速过滤。行存缓存对于主键查询Doris可以高效地定位到具体的数据行Rowset。SelectDB进一步引入了行级缓存机制将热点用户的特征数据缓存在内存中后续查询可以直接从内存获取将延迟降低到亚毫秒级。向量化查询与结果集优化即使一次查询上千个用户Doris的向量化执行引擎也能高效处理。同时优化网络传输协议减少结果集序列化/反序列化的开销。在实际部署中我们通常会将特征服务封装成一个gRPC或HTTP服务背后连接Doris。该服务接收一批用户ID拼装成高效的SQL查询例如SELECT * FROM user_features WHERE user_id IN (?, ?, ...)从Doris获取数据后再转换成模型需要的Tensor或Protobuf格式。踩坑记录初期我们尝试用OR条件拼接大量user_idWHERE user_id ‘A’ OR user_id ‘B’ ...性能很差。后来改为IN语句并确保一次查询的ID数量在一个合理批次如1000个性能提升了一个数量级。另外要密切关注Doris FE前端的并发连接数和查询队列高并发下需要合理调整相关参数并考虑使用连接池。3. 范式二AI反馈数据流的实时分析与模型迭代AI应用不是“部署即结束”而是一个持续迭代的过程。模型上线后我们需要实时收集它的“表现”数据用户对推荐结果的点击率、转化率风控模型的误杀率和漏杀率AI客服的对话满意度等。这些数据被称为“反馈数据”或“样本数据”。第二个范式就是关于如何实时处理这些反馈数据并将其用于模型的快速评估与迭代。3.1 构建实时样本流水线一个完整的闭环是模型产生预测 - 记录预测结果和特征 - 用户产生真实行为 - 行为反馈与预测结果关联形成标注样本 - 样本流入数据平台用于分析和新一轮训练。Doris在这里的核心作用是实时关联与聚合。假设我们有一个推荐系统推荐模型在线从范式一获取实时特征生成推荐列表并将request_id,user_id,推荐item列表、所用模型版本等日志实时写入Kafka。用户行为点击、购买日志也实时写入另一个Kafka Topic。通过Flink进行流式Join将行为日志和推荐日志通过request_id关联起来形成一条条完整的样本记录包含特征、预测结果、真实结果然后实时写入Doris的一张样本表。这张样本表的数据立刻就可以用于分析。3.2 实时模型评估与监控看板有了实时流入的样本表数据科学家和算法工程师的工作方式就被改变了。他们不再需要等待每天的离线样本产出而是可以实时计算模型指标在Doris中通过SQL几乎可以实时计算AUC、准确率、召回率、CTR等核心指标。例如每5分钟统计一次最新模型版本的AUC。-- 简化示例计算过去10分钟模型版本v2.1的CTR SELECT model_version, COUNT(*) as total_impressions, SUM(if(clicked0, 1, 0)) as total_clicks, SUM(if(clicked0, 1, 0)) / COUNT(*) as ctr FROM realtime_sample_table WHERE event_time NOW() - INTERVAL 10 MINUTE AND model_version v2.1 GROUP BY model_version;构建实时监控告警将上述SQL做成定时任务一旦发现CTR异常下跌或AUC低于阈值立即触发告警集成如Prometheus、AlertManager团队可以第一时间介入排查是模型问题、特征问题还是线上数据分布发生了漂移。多维下钻分析当发现整体指标下跌时可以快速在Doris中下钻分析“是哪个地域的用户CTR跌了”、“是新用户群体还是老用户群体”、“是某个商品类目的问题”。这种实时交互式分析能力能极大缩短问题定位时间。经验之谈实时样本表通常数据量巨大且增长极快必须做好数据生命周期管理。我们的策略是Doris中只保留最近7-14天的高频查询热数据用于实时监控和交互分析。更早的历史样本通过Doris的备份恢复功能或直接通过管道定期归档到成本更低的对象存储如S3中用于长期的离线训练和回溯分析。Doris 2.0版本对冷热数据分离的支持使得这种架构更加优雅。4. 范式三面向AI Agent的复杂、动态数据查询支持前两个范式还是“人驱动”的——人设计特征人评估模型。而AI Agent智能体的兴起带来了第三个范式由AI自主发起对数据平台的复杂、动态查询以完成其决策或任务。想象一个电商客服Agent用户问“帮我对比一下最近三个月iPhone 15和华为Mate 60在这家店的销量和客户评价趋势。” 这个Agent需要理解自然语言将其转化为一个结构化的查询意图。根据意图动态生成SQL查询语句。执行查询获取结果。将结果可能是多个表、图表的数据理解并整合用自然语言回复给用户。这里的关键挑战在于第2和第3步生成的SQL可能是复杂的、多表关联的、带有高级分析函数如窗口函数计算趋势的并且查询模式是事先无法完全预知的。4.1 Doris如何支撑Agent的即席查询强大的SQL兼容性与性能Apache Doris对标准SQL特别是OLAP场景常用的SQL有很好的支持包括复杂的JOIN、子查询、窗口函数、CTE公共表表达式等。这意味着Agent生成的绝大多数合理SQLDoris都能直接执行。其MPP架构和向量化引擎保证了即使面对即席复杂查询也能在可接受的时间内返回结果。联邦查询Federation能力Agent可能需要的数据并不全在Doris里。用户评价文本可能存在于Elasticsearch中库存细节可能在MySQL里。SelectDB增强了Doris的多源数据目录功能使其能够通过外表的方式几乎零成本地查询这些外部数据源。对于Agent来说它像是在查询一个统一的数据库无需关心数据物理存储在哪里。-- Agent生成的SQL可能类似于经过简化 WITH iphone_sales AS ( SELECT date, SUM(sales) as daily_sales FROM doris_local.sales_table WHERE product_name LIKE %iPhone 15% AND date 2023-10-01 GROUP BY date ), huawei_sales AS ( SELECT date, SUM(sales) as daily_sales FROM mysql_inventory.sales_external_table -- 这是一个映射到MySQL的外部表 WHERE product_name LIKE %Mate 60% AND date 2023-10-01 GROUP BY date ), iphone_reviews AS ( SELECT AVG(rating) as avg_rating, DATE(create_time) as date FROM es_catalog.product_reviews -- 这是一个映射到Elasticsearch的外部表 WHERE product_id IN (SELECT id FROM doris_local.products WHERE name LIKE %iPhone 15%) AND create_time 2023-10-01 GROUP BY DATE(create_time) ) SELECT s.date, i.daily_sales as iphone_sales, h.daily_sales as huawei_sales, r.avg_rating as iphone_avg_rating FROM ... -- 将上述CTE进行关联和趋势计算资源隔离与稳定性保障Agent的查询是不可预测的一个编写不当的SQL可能会消耗大量资源影响核心的特征供给和样本分析业务。SelectDB提供了完善的资源隔离能力通过资源组Resource Group可以为Agent查询分配固定的计算资源CPU、内存并设置查询超时和内存限制确保其“疯狂”的查询不会拖垮整个集群。4.2 工程实践构建Agent的数据查询层在实际系统中我们不会让Agent直接连接Doris。一个典型的架构是Agent将自然语言转化为结构化的“查询请求”包含意图、实体、过滤条件等发送给一个数据查询中间件。中间件负责a) 根据元数据信息将查询请求转换为优化过的SQLb) 管理数据库连接池执行SQLc) 处理查询结果可能进行二次聚合或格式转换d) 实施限流、熔断、降级策略保护后端数据库。这个中间件最终向Doris或通过Doris向其他数据源发起查询。Doris的稳定性和性能是这个中间件能够可靠工作的基石。我们曾遇到一个案例一个分析Agent在生成月度报告时并发发起了数十个复杂查询。得益于Doris的资源组隔离这些查询被限制在指定的资源范围内虽然自身变慢但完全没有影响线上实时特征服务的P99延迟。核心建议当为AI Agent设计数据查询层时一定要把Doris视为一个“能力强大但需要合理使用”的伙伴。建立清晰的查询模版和审核机制即使是自动生成的SQL避免笛卡尔积等灾难性查询。充分利用Doris的查询计划分析EXPLAIN功能在中间件层面对复杂查询进行初步的成本预估和拦截。5. 融合实践三大范式如何在一个系统中协同工作纸上谈兵终觉浅。我们来看一个简化的电商推荐系统案例如何将三大范式融为一体。系统目标实现一个实时个性化推荐系统并能根据效果实时调整策略同时允许运营Agent查询推荐效果数据。架构与数据流数据源用户行为日志点击、浏览、搜索、商品信息、订单数据、用户画像离线更新实时更新实时流入Kafka。范式一特征供给Flink消费Kafka中的实时行为流进行轻量聚合如用户最近1小时点击序列与批处理生成的用户静态画像、商品特征等一起通过Flink-Doris-Connector实时写入Doris的user_realtime_features和item_features表主键模型。推荐模型服务在线接收推荐请求时向特征服务发起调用特征服务从Doris中毫秒级查询出相关用户和商品的实时特征返回给模型。范式二反馈分析推荐服务每次推荐后将request_id,user_id,推荐列表,模型版本,所用特征版本等日志写入Kafkapredict_log。用户后续行为click_log,buy_log也写入Kafka。另一个Flink作业实时消费predict_log和click_log进行流式Join生成包含正负样本的realtime_sample流并实时写入Doris的sample_table。实时监控看板基于sample_table在Doris上建立一系列物化视图或定时查询实时计算各模型版本的CTR、转化率等核心指标并展示在Grafana等看板上。设置报警规则。模型迭代每小时/每天训练管道从Doris中导出最新的样本数据或结合历史归档数据进行模型重训练。新模型通过A/B测试框架上线。范式三Agent查询运营人员向运营分析Agent提问“上周上线的推荐模型V3.2在新用户和老用户上的点击率差异大吗”Agent通过中间件生成类似如下的SQL查询DorisSELECT user_type, model_version, COUNT(*) as impressions, SUM(is_click) as clicks, SUM(is_click)/COUNT(*) as ctr FROM sample_table WHERE event_date 2024-01-15 AND model_version v3.2 GROUP BY user_type, model_version;Doris快速返回结果Agent组织成自然语言和图表回复给运营人员。在这个闭环里Doris是贯穿始终的核心数据枢纽它既是实时特征的供给者又是反馈数据的分析者还是智能查询的响应者。数据在其中高速流动价值被实时挖掘和应用。6. 选型与部署SelectDB Cloud 带来的关键增强纯粹使用开源Apache Doris也能部分实现上述范式但SelectDB Cloud其云服务或SelectDB Enterprise企业版提供了许多关键增强使得构建这样的AI时代实时分析系统更加高效、稳定和易管理。极简的实时数据接入提供可视化的数据导入工作流轻松配置从Kafka、Flink、MySQL Binlog、S3等数十种数据源的实时/批量导入大大降低了链路搭建的复杂度。企业级稳定性与弹性云服务天然具备高可用、自动故障转移、弹性扩缩容的能力。对于需要7x24小时提供实时特征和样本分析的服务来说运维复杂度大幅降低。存储计算分离架构使得可以独立扩展存储或计算资源成本更优。增强的查询性能针对上述范式中的特定场景如高并发点查、复杂即席查询SelectDB进行了深度优化例如更智能的查询优化器、更高效的向量化算子、以及针对云存储的I/O优化等。完善的可观测性与管理提供丰富的监控指标集群负载、查询延迟、资源使用率等、审计日志和性能分析工具让开发者能清晰洞察系统状态快速定位性能瓶颈。安全与合规提供网络隔离、数据加密、访问控制、审计等企业级安全特性这对于处理用户行为、画像等敏感数据的AI应用至关重要。在项目初期如果团队资源有限从SelectDB Cloud开始是一个风险更低、启动更快的选择。它让你能更专注于业务逻辑和AI算法本身而不是集群部署、调优和运维的泥潭。从我实际推动项目落地的经验来看从传统数仓思维转向“AI时代实时分析”思维最大的挑战往往不是技术而是组织协作和数据流程的重构。数据工程师需要更贴近算法团队的需求算法工程师需要更理解数据系统的能力和限制。而像“Apache Doris SelectDB”这样的技术栈以其一体化的能力和清晰的演进范式正在为这种融合提供一条可行的路径。它不再只是一个数据库而是一个支撑AI数据智能闭环的基础设施。
返回列表