如何用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
相关产品推荐
相关产品推荐

