You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Python GridDB客户端为智能电表与传感器数据创建宽表

基于GridDB Python客户端实现智能电表数据合并与窄表转宽表方案

需求梳理

  • 多台启用的智能电表每日定时上报数据,需在特定时间点所有电表完成上报后自动合并数据
  • 每台电表关联传感器数据:读数数值、电池状态、温度、网络信号强度
  • 将单电表窄表格式转换为适合报表分析的宽表格式

现有代码问题分析

  1. combine_meter_data采用客户端循环查询,每次遍历电表都重新查询全量读数,效率低下,未利用GridDB的端侧计算能力
  2. 宽表容器未定义结构,直接调用put_container(wide_table_container)会抛出参数错误
  3. 缺少“特定时间点所有电表完成上报”的判断逻辑

优化实现方案

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 06:05:55