如何按分/时/日/周聚合DynamoDB毫秒级时间戳数据并补全缺失值?
DynamoDB多时间维度聚合查询方案(含空维度补0)
一、先把数据存对:优化表结构
每30秒上报一条数据,毫秒级时间戳,要做多维度聚合,首先得在存储阶段就把维度信息提前算好,避免后续查询时实时计算拖慢速度。
给原始数据表(比如叫sensor_raw_data)这么设计:
- 分区键:
device_id(如果是全局数据,就用固定值比如global) - 排序键:
timestamp(毫秒级整数,保证时序) - 额外字段:提前计算好
minute_key(格式yyyyMMddHHmm,比如202311151433)、hour_key(yyyyMMddHH)、day_key(yyyyMMdd)、week_key(yyyyWW,WW是ISO周数),还有上报的数值value。
示例数据:
{ "device_id": "dev_001", "timestamp": 1699999999000, "value": 25.6, "minute_key": "202311151433", "hour_key": "2023111514", "day_key": "20231115", "week_key": "202346" }
二、两种聚合方案:实时查/预聚合
1. 实时查询(数据量小的时候用)
如果你的数据量不大,不用提前算,查出来在应用层处理就行,步骤很简单:
- 第一步:先把要查询的时间范围内的所有维度键列出来。比如查14:00到14:10的分钟数据,就生成10个
minute_key:202311151400、202311151401...202311151410。 - 第二步:用DynamoDB的
Query拉取这个时间范围内的所有数据,按维度键分组算平均值。 - 第三步:把预先列好的维度键和聚合结果一一匹配,没找到的就填0。
Python示例代码:
import boto3 from collections import defaultdict dynamodb = boto3.resource('dynamodb') raw_table = dynamodb.Table('sensor_raw_data') # 生成目标分钟的key集合(14:00到14:10) target_minutes = [f"2023111514{str(i).zfill(2)}" for i in range(0, 11)] # 对应的毫秒时间戳范围 start_ts = 1699999200000 # 2023-11-15 14:00:00 end_ts = 1699999800000 # 2023-11-15 14:10:00 # 拉取数据 response = raw_table.query( KeyConditionExpression="device_id = :dev_id AND timestamp BETWEEN :start AND :end", ExpressionAttributeValues={ ":dev_id": "dev_001", ":start": start_ts, ":end": end_ts }, ProjectionExpression="minute_key, value" ) # 分组算平均值 agg_data = defaultdict(list) for item in response['Items']: agg_data[item['minute_key']].append(float(item['value'])) avg_results = {k: sum(v)/len(v) for k, v in agg_data.items()} # 补0,生成最终结果 final_result = {minute: avg_results.get(minute, 0.0) for minute in target_minutes} print(final_result)
2. 预聚合(数据量大的时候必用)
数据量上来之后,实时查询会慢到离谱,这时候就得提前把聚合结果算好存在另一个表里。
预聚合表设计
表名就叫sensor_agg_data:
- 分区键:
device_id - 排序键:
agg_key(格式是维度类型#维度值,比如minute#202311151433、hour#2023111514) - 字段:
avg_value(平均值)、count(该维度下的数据条数,用来更新平均值)、updated_at(最后更新时间)
预聚合逻辑
用Lambda触发原始表的写入事件,每次有新数据进来,就更新四个维度的预聚合记录:
- 拿新数据的四个维度键,分别生成对应的
agg_key - 对每个
agg_key,查预聚合表有没有这条记录:- 有记录:用
(旧平均值*旧条数 + 新值)/(旧条数+1)计算新平均值,然后原子更新avg_value和count - 没记录:直接插入新记录,
avg_value设为新值,count设为1
- 有记录:用
查询预聚合数据并补0
和实时查询逻辑差不多,先列全目标维度键,再查预聚合表,最后补0:
agg_table = dynamodb.Table('sensor_agg_data') # 要查2023-11-15全天的小时数据 target_hours = [f"20231115{str(i).zfill(2)}" for i in range(0, 24)] # 查询该设备所有小时维度的聚合数据 response = agg_table.query( KeyConditionExpression="device_id = :dev_id AND begins_with(agg_key, :prefix)", ExpressionAttributeValues={ ":dev_id": "dev_001", ":prefix": "hour#" }, ProjectionExpression="agg_key, avg_value" ) # 整理聚合结果 agg_results = {} for item in response['Items']: hour_key = item['agg_key'].split('#')[1] agg_results[hour_key] = float(item['avg_value']) # 补0得到最终结果 final_result = {hour: agg_results.get(hour, 0.0) for hour in target_hours}
三、踩坑提醒
- 维度键的生成要统一:比如周数必须用ISO标准,不然不同系统生成的周数会乱,导致聚合出错。
- 预聚合要保证原子性:更新预聚合表时,必须用DynamoDB的
UpdateExpression做原子操作,避免并发写入导致数据错误。比如:
# 原子更新的示例表达式 update_expr = """ SET avg_value = (:old_avg * :old_count + :new_val) / (:old_count + 1), count = count + :incr, updated_at = :ts """ expr_vals = { ":old_avg": existing_avg, ":old_count": existing_count, ":new_val": new_value, ":incr": 1, ":ts": current_timestamp }
- 补0的核心是先列全所有要查的维度键:不管有没有数据,先把时间范围内的所有维度列出来,再和聚合结果匹配,缺失的就填0,不能只靠查询结果来生成最终数据。
内容的提问来源于stack exchange,提问作者Pankaj Verma
相关产品推荐
相关产品推荐

