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时间序列引擎

相关新闻

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

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

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

2026/7/24 8:01:17 阅读更多 →
揭秘!24小时AI客服领域,究竟哪家才是优秀服务商?

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

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

2026/7/24 3:56:58 阅读更多 →
某智驾大牛创业

某智驾大牛创业

作者:钟声编辑:Mark出品:红色星际头图:智能驾驶图片据悉,国内某头部智驾公司端到端模型技术大牛Z投身创业,并且已经拿到融资。Z不仅是该头部公司内部最年轻的对标阿里P10级别技术负责⼈,更是业内…

2026/7/24 3:52:37 阅读更多 →

最新新闻

医疗票据OCR技术解析与API对接实战

医疗票据OCR技术解析与API对接实战

1. 医疗票据OCR的技术痛点与行业需求 医疗票据的数字化处理一直是医院和医保系统的老大难问题。每天门诊大厅里堆积如山的发票、住院部源源不断的结算单,传统的人工录入方式不仅效率低下,还容易出错。我曾亲眼见过某三甲医院的财务科,20多名工…

2026/7/24 8:21:45 阅读更多 →
医疗票据OCR技术:医院数字化转型的核心解决方案

医疗票据OCR技术:医院数字化转型的核心解决方案

1. 医疗票据OCR技术为何成为医院数字化刚需 每天清晨,当三甲医院结算窗口排起长队时,财务人员面前堆积如山的医保单据正暴露着传统医疗票据处理的三大痛点:人工录入速度慢(熟练员工处理单张票据需3-5分钟)、差错率高&a…

2026/7/24 8:21:45 阅读更多 →
FCA-RL框架:强化学习在动态出行调度中的应用

FCA-RL框架:强化学习在动态出行调度中的应用

1. 项目概述:FCA-RL框架的核心价值 在网约车、共享出行等实时服务领域,市场供需关系往往呈现剧烈波动。传统静态调度算法在面对突发天气、节假日高峰或区域性活动时,常出现响应迟滞、资源错配等问题。我们团队提出的FCA-RL(Feedba…

2026/7/24 8:21:45 阅读更多 →
交互式网络图:从NHL斯坦利杯历史挖掘球员与球队的冠军关联

交互式网络图:从NHL斯坦利杯历史挖掘球员与球队的冠军关联

1. 先搞清楚这个网络图到底能看什么 如果你对 NHL(国家冰球联盟)的斯坦利杯冠军历史感兴趣,这个交互式网络图项目值得花十分钟试试。它不是简单地把历年冠军列个表,而是把所有冠军球队、球员、教练这些实体用连线关系可视化出来。…

2026/7/24 8:21:45 阅读更多 →
FCA-RL框架:强化学习在网约车动态资源分配中的应用

FCA-RL框架:强化学习在网约车动态资源分配中的应用

1. 项目概述在出行服务领域,动态市场环境下的资源分配效率一直是行业痛点。我们团队提出的FCA-RL框架,正是针对网约车平台中出行服务商面临的这一核心挑战。这个方案通过强化学习技术,实现了在预算约束条件下的动态投资策略优化,帮…

2026/7/24 8:21:45 阅读更多 →
AI Agent如何用《论语》智慧解决现代问题

AI Agent如何用《论语》智慧解决现代问题

1. 项目背景与核心思路"半部论语治天下"这个说法最早出自宋代赵普的典故,形容儒家经典《论语》蕴含的治理智慧。而今天,我尝试用AI Agent技术来重构这个千年命题——不是简单地文本处理,而是构建一个能理解、应用《论语》智慧的智能…

2026/7/24 8:20:45 阅读更多 →

日新闻

用Highcharts 创建可拖拽三维散点立方体3D图表

用Highcharts 创建可拖拽三维散点立方体3D图表

该案例基于Highcharts scatter3d 三维散点图实现空间立方体散点可视化,核心特色:三维 X/Y/Z 三轴空间,所有散点分布在 0~10 立方体空间内;散点使用径向渐变实现立体 3D 圆球质感;支持鼠标 / 触屏拖拽画布,…

2026/7/24 0:00:29 阅读更多 →
AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls:进程创建路径上的 DLL 入口

AppCertDlls:进程创建路径上的 DLL 入口 AppCertDlls 位于 HKLM\System\CurrentControlSet\Control\Session Manager\AppCertDlls。本文的程序功能是只读列出这个键在 64 位和 32 位注册表视图中的全部值,并显示每条值的来源、名称、类型和可安全显示的数…

2026/7/24 0:00:29 阅读更多 →
我的编程之路:第一篇博客

我的编程之路:第一篇博客

大家好,我是一名编程初学者,同时这也是我编程学习之路上的第一篇博客。在这里,我想要向大家介绍我的一些想法和规划。a.自我介绍我是一个刚刚接触编程的新手,目前在学习c语言,我对编程世界充满了强烈的好奇。当然&…

2026/7/24 0:00:29 阅读更多 →

周新闻

Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/24 3:59:20 阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/24 1:23:39 阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/23 17:49:47 阅读更多 →

月新闻