DolphinDB实时统计分析:滑动窗口计算

发布时间:2026/9/12 16:51:53

DolphinDB实时统计分析:滑动窗口计算 摘要本文深入讲解DolphinDB实时统计分析技术。从滑动窗口原理到窗口函数应用从实时统计到趋势计算从多维度分析到性能优化全面介绍滑动窗口计算的核心方法。通过丰富的代码示例帮助读者掌握实时统计分析的核心技能。一、滑动窗口概述1.1 滑动窗口原理滑动窗口数据流窗口1窗口2窗口3统计结果1.2 窗口类型类型说明滚动窗口固定大小不重叠滑动窗口固定大小可重叠会话窗口基于活动间隔1.3 应用场景场景窗口类型分钟级统计滚动窗口移动平均滑动窗口会话分析会话窗口二、滚动窗口2.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(temp_agg,60000,//1分钟窗口[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)2.2 多窗口聚合//多窗口聚合defmultiWindowAgg(){//1分钟窗口 agg1mcreateTimeSeriesEngine(agg_1m,60000,[avg(temperature)asavg_temp],agg_1m_result,timestamp,device_id)//5分钟窗口 agg5mcreateTimeSeriesEngine(agg_5m,300000,[avg(temperature)asavg_temp],agg_5m_result,timestamp,device_id)//1小时窗口 agg1hcreateTimeSeriesEngine(agg_1h,3600000,[avg(temperature)asavg_temp],agg_1h_result,timestamp,device_id)}2.3 自定义聚合//自定义聚合函数defmyAgg(data){returndict(STRING,ANY,[[mean,avg(data)],[std,std(data)],[skewness,skewness(data)],[kurtosis,kurtosis(data)]])}三、滑动窗口3.1 移动平均//移动平均defmovingAverage(data,window){returnselect timestamp,mavg(temperature,window)asmafromdata}//使用 ttable(1..100asid,now()1..100*1000astimestamp,rand(20.0..30.0,100)astemperature)resultmovingAverage(t,10)3.2 移动统计//移动统计defmovingStats(data,window){returnselect timestamp,temperature,mavg(temperature,window)asmoving_avg,mstd(temperature,window)asmoving_std,mmax(temperature,window)asmoving_max,mmin(temperature,window)asmoving_min,msum(temperature,window)asmoving_sumfromdata}3.3 指数移动平均//指数移动平均defema(data,alpha0.1){resultarray(DOUBLE,data.rows())result[0]data[0]for(iin1..data.rows()){result[i]alpha*data[i](1-alpha)*result[i-1]}returnresult}四、实时统计4.1 实时计数//实时计数 share table(1:0,time_windowcount,[TIMESTAMP,LONG])ascount_result countEnginecreateTimeSeriesEngine(count_engine,60000,[count(*)ascount],count_result,timestamp)subscribeTable(,sensor_stream,count,-1,countEngine,true)4.2 实时求和//实时求和 share table(1:0,time_windowtotal_temp,[TIMESTAMP,DOUBLE])assum_result sumEnginecreateTimeSeriesEngine(sum_engine,60000,[sum(temperature)astotal_temp],sum_result,timestamp)subscribeTable(,sensor_stream,sum,-1,sumEngine,true)4.3 实时百分位//实时百分位defpercentileAgg(data,p){returnpercentile(data,p)}//使用 share table(1:0,time_windowp50p90p99,[TIMESTAMP,DOUBLE,DOUBLE,DOUBLE])aspercentile_result percentileEnginecreateTimeSeriesEngine(percentile_engine,60000,[percentile(temperature,50)asp50,percentile(temperature,90)asp90,percentile(temperature,99)asp99],percentile_result,timestamp)subscribeTable(,sensor_stream,percentile,-1,percentileEngine,true)五、趋势计算5.1 变化率//变化率defcalculateChangeRate(data){returnselect timestamp,temperature,deltas(temperature)aschange,deltas(temperature)/prev(temperature)*100aschange_ratefromdata}5.2 累计统计//累计统计defcumulativeStats(data){returnselect timestamp,temperature,cumsum(temperature)ascum_sum,cumavg(temperature)ascum_avg,cummax(temperature)ascum_max,cummin(temperature)ascum_minfromdata}5.3 趋势检测//趋势检测defdetectTrend(data,window10){resultselect timestamp,temperature,mavg(temperature,window)asma,iif(temperaturemavg(temperature,window),up,iif(temperaturemavg(temperature,window),down,stable))astrendfromdatareturnresult}六、多维度分析6.1 分组统计//分组统计defgroupStats(data,groupCol){returnselect avg(temperature)asavg_temp,max(temperature)asmax_temp,min(temperature)asmin_temp,std(temperature)asstd_temp,count(*)ascountfromdata group byeval(groupCol)}6.2 多维度聚合//多维度聚合defmultiDimAgg(data){returnselect device_id,bar(timestamp,1h)ashour,avg(temperature)asavg_temp,max(temperature)asmax_temp,min(temperature)asmin_tempfromdata group by device_id,bar(timestamp,1h)}6.3 交叉分析//交叉分析defcrossAnalysis(data){returnselect device_id,iif(temperature25,high,low)astemp_level,count(*)ascount,avg(humidity)asavg_humidityfromdata group by device_id,iif(temperature25,high,low)}七、性能优化7.1 增量计算//增量计算defincrementalAgg(newData,existingStats){//更新统计量 newCountexistingStats.countnewData.rows()newSumexistingStats.sumsum(newData.temperature)newAvgnewSum/newCountreturndict(STRING,ANY,[[count,newCount],[sum,newSum],[avg,newAvg]])}7.2 并行计算//并行计算defparallelAgg(data,numPartitions4){resultsarray(ANY,0)for(iin0..numPartitions){partitionselect*fromdata where device_id%numPartitionsi results.append!(aggPartition(partition))}//合并结果returnmergeResults(results)}八、实战案例8.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_tempmax_tempmin_tempstd_tempcount,[TIMESTAMP,SYMBOL,DOUBLE,DOUBLE,DOUBLE,DOUBLE,LONG])asagg_result//3.创建时间序列引擎 aggEnginecreateTimeSeriesEngine(sensor_agg,60000,[avg(temperature)asavg_temp,max(temperature)asmax_temp,min(temperature)asmin_temp,std(temperature)asstd_temp,count(*)ascount],agg_result,timestamp,device_id)subscribeTable(,sensor_stream,agg,-1,aggEngine,true)//4.移动统计 share table(1:0,timestampdevice_idtemperaturema_10ma_30,[TIMESTAMP,SYMBOL,DOUBLE,DOUBLE,DOUBLE])asma_resultdefcalculateMA(data){insert into ma_result select timestamp,device_id,temperature,mavg(temperature,10)asma_10,mavg(temperature,30)asma_30fromdata context by device_id}subscribeTable(,sensor_stream,ma,-1,def(msg){calculateMA(msg)},true)//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)//6.查询接口defgetLatestStats(){returnselect*fromagg_result order by time_window desc limit10}addFunctionView(getLatestStats)print(实时统计分析系统启动完成)九、总结本文详细介绍了DolphinDB实时统计分析滚动窗口时间序列引擎、多窗口聚合滑动窗口移动平均、移动统计、指数移动平均实时统计实时计数、实时求和、实时百分位趋势计算变化率、累计统计、趋势检测多维度分析分组统计、多维度聚合、交叉分析性能优化增量计算、并行计算思考题如何选择合适的窗口大小如何优化滑动窗口计算性能如何处理乱序数据参考资料DolphinDB时间序列引擎DolphinDB流计算系列预告下一篇将介绍实时关联分析敬请期待
延伸阅读

更多相关文章

2026/9/11 19:51:22

第七届心理健康与教育、人文发展国际学术会议(MHEHD 2026)

第七届心理健康与教育、人文发展国际学术会议(MHEHD2026)将于2026年8月17-19日在中国长沙隆重召开。会议主要围绕“心理健康”“人文教育”等研究领域展开讨论。旨在为心理健康与人文教育的专家学者及企业发展人提供一个分享研究成果、讨论存在的问题与挑…

2026/9/12 16:50:53

2026年学术论文降AI工具评测与最佳实践指南

1. 项目背景与核心需求 2026年的学术环境正在经历一场前所未有的变革。随着AI生成内容的普及,学术诚信面临全新挑战。各大高校和期刊编辑部纷纷引入AI检测工具,导致大量论文因"AI率过高"被退回或质疑。在这个背景下,论文降AI工具应…

2026/9/12 16:50:53

用一句自然语言让浏览器替你干活:MidScene.js 上手记

用一句自然语言让浏览器替你干活:MidScene.js 上手记 【免费下载链接】midscene GUI Agent for E2E Testing 项目地址: https://gitcode.com/GitHub_Trending/mid/midscene MidScene.js 是一款 AI 浏览器自动化工具,你用中文或英文写一句话&#…

2026/9/12 16:45:53

环境变量与密钥管理实战:彻底告别硬编码密码

前阵子接手一个外包项目的交接,代码拉到本地,随手翻到config.js里面静静躺着一段password: "Pssw0rd2022"。我当场截图发到工作群,问这是谁的,群里安静了十分钟,最后有人在私聊里回了一句"先跑起来再说…

2026/9/12 2:05:33

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/12 3:55:12

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/12 10:09:03

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/12 0:04:17

MATLAB仿生优化框架:长鼻浣熊算法多策略融合实现

简介:本资源是一份面向智能优化算法研究者与MATLAB初学者的仿生智能算法实践代码包,聚焦于长鼻浣熊优化算法(COA)的多策略改进与性能验证。针对传统COA易陷局部最优、收敛精度不足等问题,作者融合Circle映射初始化提升…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 JavaWeb 的校园一卡通管理系统的设计与实现 基于 JavaWeb 的校园卡业务管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 Java 的图书馆借阅管理平台的搭建与实现 基于 Java 的图书馆综合管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 6:29:36

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/12 14:32:17

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/12 6:37:43

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码