首页/新闻资讯/正文详情

DolphinDB实时聚合计算:多维度聚合

发布时间:2026/9/24 23:07:36 来源:云帆数科 栏目:资讯中心
DolphinDB实时聚合计算:多维度聚合
目录摘要一、聚合计算概述1.1 聚合类型1.2 聚合函数1.3 聚合维度二、基础聚合2.1 单表聚合2.2 分组聚合2.3 条件聚合三、多维度聚合3.1 多列分组3.2 Cube聚合3.3 Rollup聚合四、层级聚合4.1 组织层级4.2 时间层级4.3 上卷下钻五、实时聚合引擎5.1 时间序列聚合5.2 多度量聚合5.3 自定义聚合六、聚合优化6.1 增量聚合6.2 并行聚合6.3 预聚合七、实战案例7.1 完整实时聚合系统八、总结参考资料摘要本文深入讲解DolphinDB实时聚合计算技术。从聚合函数到多维度聚合从层级聚合到实时汇总从分组统计到聚合优化全面介绍实时聚合计算的核心方法。通过丰富的代码示例帮助读者掌握多维度聚合的核心技能。一、聚合计算概述1.1 聚合类型聚合计算单维度聚合聚合结果多维度聚合层级聚合1.2 聚合函数函数说明sum求和avg平均值max最大值min最小值count计数std标准差1.3 聚合维度维度说明时间维度按时间聚合设备维度按设备聚合产品维度按产品聚合区域维度按区域聚合二、基础聚合2.1 单表聚合//单表聚合defbasicAggregation(data){returnselectsum(temperature)astotal,avg(temperature)asmean,max(temperature)asmax_val,min(temperature)asmin_val,count(*)ascount,std(temperature)asstd_valfromdata}2.2 分组聚合//分组聚合defgroupAggregation(data,groupCol){returnselecteval(groupCol)asgroup_key,sum(temperature)astotal,avg(temperature)asmean,count(*)ascountfromdata group byeval(groupCol)}2.3 条件聚合//条件聚合defconditionalAggregation(data){returnselectsum(iif(temperature25,temperature,0))ashigh_temp_sum,sum(iif(temperature25,temperature,0))aslow_temp_sum,count(iif(temperature25,1,0))ashigh_count,count(iif(temperature25,1,0))aslow_countfromdata}三、多维度聚合3.1 多列分组//多列分组聚合defmultiDimAggregation(data){returnselect device_id,bar(timestamp,1h)ashour,sum(temperature)astotal,avg(temperature)asmean,max(temperature)asmax_val,min(temperature)asmin_val,count(*)ascountfromdata group by device_id,bar(timestamp,1h)}3.2 Cube聚合//Cube聚合多维度组合defcubeAggregation(data){//按设备聚合 byDeviceselect device_id,allashour,sum(temperature)astotal,avg(temperature)asmeanfromdata group by device_id//按时间聚合 byHourselectallasdevice_id,bar(timestamp,1h)ashour,sum(temperature)astotal,avg(temperature)asmeanfromdata group by bar(timestamp,1h)//按设备和时间聚合 byBothselect device_id,bar(timestamp,1h)ashour,sum(temperature)astotal,avg(temperature)asmeanfromdata group by device_id,bar(timestamp,1h)//合并returnbyDevice.union(byHour).union(byBoth)}3.3 Rollup聚合//Rollup聚合层级聚合defrollupAggregation(data){//层级设备-车间-工厂//设备级别 deviceLevelselect device_id,workshop,factory,sum(temperature)astotalfromdata group by device_id,workshop,factory//车间级别 workshopLevelselectallasdevice_id,workshop,factory,sum(temperature)astotalfromdata group by workshop,factory//工厂级别 factoryLevelselectallasdevice_id,allasworkshop,factory,sum(temperature)astotalfromdata group by factoryreturndeviceLevel.union(workshopLevel).union(factoryLevel)}四、层级聚合4.1 组织层级//组织层级聚合defhierarchyAggregation(data,hierarchy){resultsarray(ANY,0)for(levelinhierarchy){aggselecteval(level)aslevel_key,sum(temperature)astotal,avg(temperature)asmeanfromdata group byeval(level)results.append!(agg)}returnresults}4.2 时间层级//时间层级聚合deftimeHierarchyAggregation(data){//分钟级 minuteselect bar(timestamp,1m)astime,avg(temperature)asmeanfromdata group by bar(timestamp,1m)//小时级 hourselect bar(timestamp,1h)astime,avg(temperature)asmeanfromdata group by bar(timestamp,1h)//天级 dayselect date(timestamp)astime,avg(temperature)asmeanfromdata group by date(timestamp)returndict(STRING,ANY,[[minute,minute],[hour,hour],[day,day]])}4.3 上卷下钻//上卷聚合到更高层级defrollup(data,fromLevel,toLevel){returnselecteval(toLevel)aslevel,sum(temperature)astotal,avg(temperature)asmeanfromdata group byeval(toLevel)}//下钻展开到更低层级defdrilldown(data,fromLevel,toLevel,filter){filteredselect*fromdata whereeval(filter)returnselecteval(toLevel)aslevel,sum(temperature)astotal,avg(temperature)asmeanfromfiltered group byeval(toLevel)}五、实时聚合引擎5.1 时间序列聚合//创建流表 share streamTable(100000:0,device_idtimestamptemperaturehumidity,[SYMBOL,TIMESTAMP,DOUBLE,DOUBLE])assensor_stream//创建聚合结果表 share table(1:0,time_windowdevice_idavg_tempmax_tempmin_tempcount,[TIMESTAMP,SYMBOL,DOUBLE,DOUBLE,DOUBLE,LONG])asagg_result//创建聚合引擎 aggEnginecreateTimeSeriesEngine(sensor_agg,60000,[avg(temperature)asavg_temp,max(temperature)asmax_temp,min(temperature)asmin_temp,count(*)ascount],agg_result,timestamp,device_id)//订阅 subscribeTable(,sensor_stream,agg,-1,aggEngine,true)5.2 多度量聚合//多度量聚合 share table(1:0,time_windowdevice_idavg_tempavg_humidmax_tempmin_temp,[TIMESTAMP,SYMBOL,DOUBLE,DOUBLE,DOUBLE,DOUBLE])asmulti_agg multiAggEnginecreateTimeSeriesEngine(multi_agg,60000,[avg(temperature)asavg_temp,avg(humidity)asavg_humid,max(temperature)asmax_temp,min(temperature)asmin_temp],multi_agg,timestamp,device_id)subscribeTable(,sensor_stream,multi_agg,-1,multiAggEngine,true)5.3 自定义聚合//自定义聚合函数defcustomAgg(data){returndict(STRING,ANY,[[mean,avg(data)],[median,med(data)],[mode,mode(data)],[range,max(data)-min(data)],[iqr,percentile(data,75)-percentile(data,25)]])}六、聚合优化6.1 增量聚合//增量聚合 sharedict(STRING,ANY)asaggStatedefincrementalAgg(newData){for(rowinnewData){keyrow.device_idif(notaggState.has(key)){aggState[key]dict(STRING,ANY,[[sum,0.0],[count,0],[max,-infinity],[min,infinity]])}stateaggState[key]state[sum]row.temperature state[count]1state[max]max(state[max],row.temperature)state[min]min(state[min],row.temperature)}}6.2 并行聚合//并行聚合defparallelAgg(data,numWorkers4){resultsarray(ANY,0)//分区处理for(iin0..numWorkers){partitionselect*fromdata where device_id%numWorkersi results.append!(aggPartition(partition))}//合并结果returnmergeAggResults(results)}defmergeAggResults(results){totalSumsum(each(def(r){r.sum},results))totalCountsum(each(def(r){r.count},results))returndict(STRING,ANY,[[sum,totalSum],[count,totalCount],[avg,totalSum/totalCount]])}6.3 预聚合//预聚合表 share table(1:0,device_idhourpre_sumpre_countpre_maxpre_min,[SYMBOL,TIMESTAMP,DOUBLE,LONG,DOUBLE,DOUBLE])aspre_agg//定时预聚合defpreAggregationTask(){while(true){nownow()hourStartbar(now,1h)//聚合最近一小时数据 aggselect device_id,sum(temperature)aspre_sum,count(*)aspre_count,max(temperature)aspre_max,min(temperature)aspre_minfromsensor_stream where timestamphourStart group by device_id pre_agg.append!(agg)sleep(3600000)}}七、实战案例7.1 完整实时聚合系统//实时聚合计算系统//1.创建数据流 share streamTable(100000:0,device_idtimestamptemperaturehumiditypressure,[SYMBOL,TIMESTAMP,DOUBLE,DOUBLE,DOUBLE])assensor_stream enableTablePersistence(sensor_stream,true,true,1000000)//2.创建聚合结果表 share table(1:0,time_windowdevice_idavg_tempavg_humidmax_tempmin_tempcount,[TIMESTAMP,SYMBOL,DOUBLE,DOUBLE,DOUBLE,DOUBLE,LONG])asagg_result//3.创建聚合引擎 aggEnginecreateTimeSeriesEngine(sensor_agg,60000,[avg(temperature)asavg_temp,avg(humidity)asavg_humid,max(temperature)asmax_temp,min(temperature)asmin_temp,count(*)ascount],agg_result,timestamp,device_id)subscribeTable(,sensor_stream,agg,-1,aggEngine,true)//4.多维度聚合接口defgetMultiDimAgg(startTime,endTime){tloadTable(dfs://sensor_db,sensor_data)returnselect device_id,date(timestamp)asdate,bar(timestamp,1h)ashour,avg(temperature)asavg_temp,max(temperature)asmax_temp,min(temperature)asmin_temp,count(*)ascountfromt where timestamp between startTimeandendTime group by device_id,date(timestamp),bar(timestamp,1h)}addFunctionView(getMultiDimAgg)//5.模拟数据defgenerateMockData(){while(true){datatable(take(1..10,10)asdevice_id,take(now(),10)astimestamp,rand(20.0..30.0,10)astemperature,rand(40.0..60.0,10)ashumidity,rand(1000.0..1020.0,10)aspressure)sensor_stream.append!(data)sleep(5000)}}submitJob(mock_data,模拟数据,generateMockData)print(实时聚合计算系统启动完成)八、总结本文详细介绍了DolphinDB实时聚合计算基础聚合单表聚合、分组聚合、条件聚合多维度聚合多列分组、Cube聚合、Rollup聚合层级聚合组织层级、时间层级、上卷下钻实时聚合引擎时间序列聚合、多度量聚合、自定义聚合聚合优化增量聚合、并行聚合、预聚合思考题如何设计高效的多维度聚合如何优化实时聚合性能如何处理聚合中的数据倾斜参考资料DolphinDB聚合函数DolphinDB时间序列引擎

相关推荐

微信小程序二维码生成深度解析:weapp-qrcode架构设计与最佳实践
微信小程序二维码生成深度解析:weapp-qrcode架构设计与最佳实践

微信小程序二维码生成深度解析:weapp-qrcode架构设计与最佳实践 【免费下载链接】weapp-qrcode weapp.qrcode.js 在 微信小程序 中,快速生成二维码 项目地址: https://gitcode.com/gh_mirrors/we/weapp-qrcode 在微信小程序开发中,二维… · 2026/9/24 23:07:19

电子证件照片制作全教程:手机免费操作、微信支付宝流程、标准尺寸底色大全
电子证件照片制作全教程:手机免费操作、微信支付宝流程、标准尺寸底色大全

2026年各类线上报名、证件办理、入职存档、签证申请等场景,均需要合规的电子证件照。很多人常因尺寸不符、底色错误、画质模糊、文件大小超标导致上传审核失败。本文整理全套手机免费制作方法,涵盖微信、支付宝主流操作流程,明确通用标准尺寸… · 2026/9/11 21:10:40

揭秘!24小时AI客服领域,究竟哪家才是优秀服务商?
揭秘!24小时AI客服领域,究竟哪家才是优秀服务商?

在当今数字化时代,24小时AI客服成为众多企业提升服务水平与效率的重要工具,它能为企业实现全天候客户接待,带来更好的客户体验。那么,该领域有哪些优秀服务商呢?行业背景与痛点随着市场对客服服务要求的提升&#xff0… · 2026/9/21 12:29:43

Vite 掉坑自救:process is not defined、打包慢与微前端方案详解
Vite 掉坑自救:process is not defined、打包慢与微前端方案详解

看到这个标题我愣了一下。作为一个从 webpack 时代一路折腾过来的前端,我第一反应是:Vite 不是尤雨溪自己做的吗?自己出个工具干掉自己?这剧本不对。但等我顺着热搜词翻了一圈,发现大家真正在纠结的根本不是“Vize 能不… · 2026/9/24 23:07:31

PySide6表格多格式录入与主从表联动实现
PySide6表格多格式录入与主从表联动实现

做桌面端业务工具的老哥应该都有同感:表格控件是绕不开的核心组件,但你很少能只靠一个QTableWidget走天下。至少在我最近做的进销存桌面工具里,光表格录入就折腾了两个多星期——用户要求在一个表格里同时敲文本、数字、日期,还得… · 2026/9/24 23:07:31

C++多态详解:虚函数、对象切片与动态绑定原理
C++多态详解:虚函数、对象切片与动态绑定原理

先说一个我在C初学群里见过无数次的场景&#xff1a;有人写了一个Shape基类&#xff0c;派生出Circle和Rectangle&#xff0c;然后想用一个vector把不同类型的图形装在一起&#xff0c;遍历调用area()。大多数人的第一反应是&#xff1a;std::vector<Shape> shapes; shap… · 2026/9/24 23:07:31

异构固定翼无人机集群协同搜索:Matlab复现中的运动学约束与避障决策解析
异构固定翼无人机集群协同搜索:Matlab复现中的运动学约束与避障决策解析

把固定翼无人机集群撒到一片完全未知的区域里去找目标&#xff0c;这件事听起来像是"多画几条搜索路径就行"&#xff0c;但真正在Matlab里把一整套方法跑通之后&#xff0c;我会告诉你&#xff1a;这套系统涉及的运动学约束、决策分配、避障协同&#xff0c;每一个模… · 2026/9/24 23:07:31

推荐系统召回阶段UserCF实战指南:原理、Spark实现与踩坑笔记
推荐系统召回阶段UserCF实战指南:原理、Spark实现与踩坑笔记

在推荐系统里&#xff0c;“召回”决定了整个推荐效果的天花板。很多人一上来就盯着精排模型、多目标排序&#xff0c;却忽略了一个基础又耐打的方法&#xff1a;UserCF&#xff0c;也就是基于用户的协同过滤。这篇文章不打算讲那些花哨的深度学习框架&#xff0c;只围绕一个主… · 2026/9/24 23:07:31

AI微服务底座向导式安装实战:Ollama+Qdrant+Dify+网关一键部署
AI微服务底座向导式安装实战:Ollama+Qdrant+Dify+网关一键部署

我一直觉得&#xff0c;AI 应用开发里最劝退人的环节不是写代码&#xff0c;而是搭环境。你想做一个带知识库问答的智能体&#xff0c;背后要跑模型推理、向量检索、应用编排&#xff0c;再来个 API 网关做统一入口&#xff0c;这一整套微服务底座手动配下来&#xff0c;光依赖… · 2026/9/24 23:07:18

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介&#xff1a;这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源&#xff0c;围绕YOLOv8实现渔船作业监控系统&#xff0c;可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件&#xff0c;约24.21MB&#xff0c;以70个Python源码文件为… · 2026/9/24 0:00:13

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介&#xff1a;面向时间序列数据建模的一维卷积神经网络完整实现&#xff0c;适合深度学习入门者及需要快速验证时序模型的研究者&#xff0c;能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小&#xff0c;只有3KB&#xff0c;内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L&#xff0c;而是舌尖上的L最近在几个方言群和语音教学社群里&#xff0c;反复看到有人发一句&#xff1a;“也说字母L&#xff1a;柔软的长舌”。初看以为是英语发音课笔记&#xff0c;点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码