如何使用Python GridDB客户端为智能电表与传感器数据创建宽表
基于GridDB Python客户端实现智能电表数据合并与窄表转宽表方案
需求梳理
- 多台启用的智能电表每日定时上报数据,需在特定时间点所有电表完成上报后自动合并数据
- 每台电表关联传感器数据:读数数值、电池状态、温度、网络信号强度
- 将单电表窄表格式转换为适合报表分析的宽表格式
现有代码问题分析
combine_meter_data采用客户端循环查询,每次遍历电表都重新查询全量读数,效率低下,未利用GridDB的端侧计算能力- 宽表容器未定义结构,直接调用
put_container(wide_table_container)会抛出参数错误 - 缺少“特定时间点所有电表完成上报”的判断逻辑
优化实现方案
1. 利用GridDB SQL JOIN完成数据关联
在GridDB端执行JOIN查询,直接关联电表信息和对应的传感器读数,减少客户端与数据库的交互次数,提升效率。
2. 实现全量上报校验逻辑
通过统计指定时间点的上报电表数量,与启用的电表总数对比,确认是否满足合并条件。
3. 定义宽表容器结构并完成转换
根据宽表需求设计容器字段,将关联后的窄表数据转换为宽表格式存储。
完整优化代码
import griddb_python as griddb import datetime def check_all_meters_reported(gridstore, target_timestamp): # 查询启用的电表总数 meter_count_query = gridstore.query("meters_container", "SELECT COUNT(*) FROM meters_container WHERE enabled = TRUE") meter_count = meter_count_query.fetch_one()[0] # 查询该时间点已上报的电表数量 reported_count_query = gridstore.query( "readings_container", f"SELECT COUNT(DISTINCT meter_id) FROM readings_container WHERE timestamp = TIMESTAMP('{target_timestamp}')" ) reported_count = reported_count_query.fetch_one()[0] return reported_count == meter_count def combine_meter_data(gridstore, target_timestamp): # 使用GridDB SQL JOIN关联电表与读数数据 join_query = gridstore.query( "", f""" SELECT m.meter_id, m.meter_name, r.timestamp, r.reading_value, r.battery_status, r.temperature, r.network_signal_strength FROM meters_container m JOIN readings_container r ON m.meter_id = r.meter_id WHERE r.timestamp = TIMESTAMP('{target_timestamp}') AND m.enabled = TRUE """ ) # 转换为宽表格式:以时间戳为主键,每个电表的指标作为独立字段 wide_data = {} for row in join_query.fetch(False): meter_id, meter_name, ts, reading_val, battery, temp, signal = row ts_str = ts.strftime("%Y-%m-%d %H:%M:%S") if ts_str not in wide_data: wide_data[ts_str] = { 'timestamp': ts, } # 为每个电表的指标添加独立字段 wide_data[ts_str][f"{meter_name}_reading"] = reading_val wide_data[ts_str][f"{meter_name}_battery"] = battery wide_data[ts_str][f"{meter_name}_temperature"] = temp wide_data[ts_str][f"{meter_name}_signal"] = signal return list(wide_data.values()) def main(): factory = griddb.StoreFactory.get_instance() # 替换为实际集群信息 cluster_info = { "host": "your_host", "port": "your_port", "cluster_name": "your_cluster", "username": "admin", "password": "admin" } try: gridstore = factory.get_store(**cluster_info) # 示例目标时间点(可替换为定时任务触发的时间) target_ts = "2023-09-14 12:00:00" # 校验所有电表是否完成上报 if not check_all_meters_reported(gridstore, target_ts): print(f"Not all meters reported at {target_ts}, skipping merge") return # 合并并转换为宽表数据 wide_table_data = combine_meter_data(gridstore, target_ts) # 定义宽表容器结构(根据实际电表数量和字段调整) # 示例:假设有两个电表Meter_A和Meter_B,动态生成字段或提前定义 wide_container_fields = [ ["timestamp", griddb.Type.TIMESTAMP], ["Meter_A_reading", griddb.Type.FLOAT], ["Meter_A_battery", griddb.Type.INT], ["Meter_A_temperature", griddb.Type.FLOAT], ["Meter_A_signal", griddb.Type.INT], ["Meter_B_reading", griddb.Type.FLOAT], ["Meter_B_battery", griddb.Type.INT], ["Meter_B_temperature", griddb.Type.FLOAT], ["Meter_B_signal", griddb.Type.INT] ] wide_container_info = griddb.ContainerInfo( "MeterWideTable", wide_container_fields, griddb.ContainerType.COLLECTION, griddb.IndexInfo([griddb.IndexType.DEFAULT, "timestamp"]) ) wide_container = gridstore.put_container(wide_container_info) # 插入宽表数据 for data in wide_table_data: # 按容器字段顺序整理数据 row = [ data['timestamp'], data['Meter_A_reading'], data['Meter_A_battery'], data['Meter_A_temperature'], data['Meter_A_signal'], data['Meter_B_reading'], data['Meter_B_battery'], data['Meter_B_temperature'], data['Meter_B_signal'] ] wide_container.put(row) print(f"Wide table data for {target_ts} stored successfully") except griddb.GSException as e: print(f"GridDB Error: {e}") for i in range(e.get_error_stack_size()): print(f"Error Stack {i}: {e.get_message(i)}") if __name__ == "__main__": main()
关键优化点说明
- 端侧计算优化:使用GridDB的SQL JOIN替代客户端循环查询,减少数据传输量,提升处理效率
- 上报校验:通过COUNT(DISTINCT)统计已上报电表数量,确保只有全量上报时才执行合并
- 宽表转换:将每个电表的指标转换为独立字段,符合报表分析的需求
- 容器结构规范:明确定义宽表容器的字段类型和索引,避免数据插入错误
内容的提问来源于stack exchange,提问作者eugene nyamari
相关产品推荐
相关产品推荐

