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

如何用Python GridDB客户端关联支付与智能电表并实现触发器

基于GridDB的支付与智能电表管理项目问题

我正在开发一个基于Python GridDB客户端的支付与智能电表管理项目,已完成GridDB环境搭建,并将CSV中的支付数据导入到名为Payments的Time Series容器中,导入代码如下:

import datetime
import griddb_python
import csv
import random
import time

griddb = griddb_python
factory = griddb.StoreFactory.get_instance()

gridstore = factory.get_store(
    host="",
    port=,
    cluster_name="",
    username="",
    password=""
)

gridstore.drop_container("Payments")
expInfo = griddb.ExpirationInfo(1, griddb.TimeUnit.MINUTE, 5)

conInfo = griddb.ContainerInfo("Payments",
    [["timestamp", griddb.Type.TIMESTAMP],
    ["amount", griddb.Type.FLOAT],
    ["transaction_id", griddb.Type.STRING],
    ["payment_phone_number", griddb.Type.STRING]],
    griddb.ContainerType.TIME_SERIES, expiration=expInfo)

ts = gridstore.put_container(conInfo)

with open("payments.csv") as csvfile:
    reader = csv.DictReader(csvfile)
    for row in reader:
        timestamp = datetime.datetime.strptime(row["timestamp"], "%Y-%m-%d %H:%M:%S")
        amount = float(row["amount"])
        transaction_id = row["transaction_id"]
        payment_phone_number = row["payment_phone_number"]
        ts.put([timestamp, amount, transaction_id, payment_phone_number])

目前我遇到两个核心问题,寻求解决方案:


1. 支付与智能电表关联的高效实现

现有一个独立的智能电表容器,需要通过匹配transaction_id将支付记录与对应电表关联:若某笔支付的transaction_id在智能电表容器中存在匹配项,就将二者关联。我已编写如下函数,但希望得到更高效的实现方案:

def link_payments_to_smartmeters(gridstore, payments_container, smartmeters_container): 

    # 获取所有支付数据
    payments_data = gridstore.query(payments_container, "select *")

    # 遍历每笔支付
    for payment in payments_data:
        # 获取transaction_id(注:原代码索引有误,此处修正为Payments容器的第三列)
        transaction_id = payment[2]

        # 查询匹配的智能电表数据
        smartmeters_data = gridstore.query(smartmeters_container, "select * where smartmeter_id = ?", transaction_id)

        # 若找到匹配项,更新支付记录关联电表
        if smartmeters_data:
            smartmeter = smartmeters_data[0]
            gridstore.put(payments_container, [payment[0], payment[1], smartmeter[0], smartmeter[1]])

优化方案:

  • 添加索引加速查询:在智能电表容器的smartmeter_id字段上创建索引,大幅降低单条查询的耗时:
    # 创建智能电表容器时指定索引
    conInfo_meter = griddb.ContainerInfo("SmartMeters",
        [["smartmeter_id", griddb.Type.STRING],
        ["meter_reading", griddb.Type.FLOAT],
        # 其他字段...
        ],
        griddb.ContainerType.COLLECTION,
        index_info={"smartmeter_id": griddb.IndexType.DEFAULT}
    )
    
  • 预加载电表数据到内存字典:一次性将所有电表的smartmeter_id与对应数据加载到Python字典,避免循环中重复发起DB查询:
    def link_payments_to_smartmeters(gridstore, payments_container, smartmeters_container):
        # 预加载电表数据到字典,key为smartmeter_id
        meter_query = gridstore.query(smartmeters_container, "select *")
        meter_dict = {row[0]: row for row in meter_query}
    
        # 批量获取支付数据
        payments_data = gridstore.query(payments_container, "select *")
        update_batch = []
    
        for payment in payments_data:
            transaction_id = payment[2]
            if transaction_id in meter_dict:
                smartmeter = meter_dict[transaction_id]
                updated_payment = [payment[0], payment[1], smartmeter[0], smartmeter[1]]
                update_batch.append(updated_payment)
        
        # 批量更新支付容器,减少DB交互次数
        if update_batch:
            gridstore.put_all(payments_container, update_batch)
    
  • 使用GridDB JOIN查询:若容器类型支持JOIN,直接通过SQL JOIN一次性获取关联数据,再批量更新:
    # 执行JOIN查询获取关联记录
    join_query = gridstore.query(
        f"SELECT p.*, m.* FROM {payments_container} p JOIN {smartmeters_container} m ON p.transaction_id = m.smartmeter_id"
    )
    # 处理JOIN结果并批量更新Payments容器
    

2. GridDB中实现电表传感器触发器

完成支付与电表关联后,需要实现触发器,在检测到支付事件时自动执行电表读数更新、令牌计算等操作。以下是几种可行的实现方式:

方式1:Python客户端容器事件监听

利用GridDB的容器事件通知功能,监听Payments容器的PUT事件,当有新支付记录(或关联后的记录)插入时触发业务逻辑:

def payment_trigger(event):
    # 获取触发事件的支付记录
    payment_record = event.get_row()
    transaction_id = payment_record[2]
    
    # 查询对应电表数据
    meter_container = gridstore.get_container("SmartMeters")
    meter_data = gridstore.query(meter_container, "select * where smartmeter_id = ?", transaction_id)
    
    if meter_data:
        # 执行电表读数更新(自定义计算逻辑)
        updated_reading = meter_data[0][1] + calculate_consumption()
        # 生成用电令牌(根据支付金额)
        token = generate_token(payment_record[1])
        
        # 更新电表容器
        gridstore.put(meter_container, [transaction_id, updated_reading, token])

# 注册事件监听:监听Payments容器的PUT事件
payments_container = gridstore.get_container("Payments")
payments_container.add_event_listener(griddb.ContainerEvent.PUT, payment_trigger)

方式2:定时轮询+事件驱动

若容器事件监听不满足需求,可采用定时轮询方式,定期检查Payments容器的新支付记录并执行触发操作:

import time
import threading

def process_payment_trigger(payment, gridstore):
    # 封装电表更新、令牌计算逻辑
    transaction_id = payment[2]
    meter_container = gridstore.get_container("SmartMeters")
    meter_data = gridstore.query(meter_container, "select * where smartmeter_id = ?", transaction_id)
    if meter_data:
        updated_reading = meter_data[0][1] + calculate_consumption()
        token = generate_token(payment[1])
        gridstore.put(meter_container, [transaction_id, updated_reading, token])

def poll_payments_for_trigger(gridstore):
    last_check_time = datetime.datetime.now()
    while True:
        # 查询上次检查后新增的支付记录
        query = gridstore.query(
            "Payments",
            "select * where timestamp > ?",
            last_check_time
        )
        new_payments = list(query)
        
        for payment in new_payments:
            process_payment_trigger(payment, gridstore)
        
        last_check_time = datetime.datetime.now()
        time.sleep(60)  # 每分钟轮询一次

# 启动轮询线程
trigger_thread = threading.Thread(target=poll_payments_for_trigger, args=(gridstore,))
trigger_thread.daemon = True
trigger_thread.start()

方式3:GridDB内置SQL触发器

GridDB支持通过SQL定义触发器,在容器数据变化时自动执行预定义操作。例如,创建触发器在Payments容器插入记录时更新SmartMeters容器:

-- 通过GridDB CLI或JDBC执行
CREATE TRIGGER update_meter_on_payment
AFTER INSERT ON Payments
FOR EACH ROW
BEGIN
    UPDATE SmartMeters
    SET meter_reading = meter_reading + 10, token = generate_token(NEW.amount)
    WHERE smartmeter_id = NEW.transaction_id;
END;

注:需确认GridDB版本支持SQL触发器,且generate_token需为已注册的存储过程。


内容的提问来源于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 05:52:15