Python WebSocket下如何给OCPP协议保持连接的充电桩客户端主动发消息
问题原因分析
- 异步语法错误:
on_getlocallistversion函数没有加async关键字,函数内部直接使用await会直接抛出语法错误,这是你报错的最直接原因。 - 事件绑定冲突:你给
Heartbeat事件同时绑定了两个处理函数on_getlocallistversion和on_hearbeat,事件触发时逻辑会混乱,执行不符合预期。 - 方法调用错误:
route_message是OCPP框架内部用来处理充电桩客户端发送到服务端的消息的方法,不能用来向客户端发消息,调用该方法会触发框架内部逻辑错误。 - 代码格式错误:你贴的示例代码缩进混乱,
await self.route_message行没有正确缩进在on_getlocallistversion函数内,也会导致执行错误。
正确实现方案
1. 全局维护已连接的充电桩实例
首先需要用全局容器存储所有在线的充电桩连接实例,后续主动发消息需要从这个容器里取对应实例:
# 全局字典,key是充电桩id,value是对应的client_main实例 connected_charge_points = {}
修改连接建立的原有逻辑,增加实例存储和断开清理逻辑:
charge_point_id = path.strip('/') cp = client_main(charge_point_id, websocket) # 存入全局字典 connected_charge_points[charge_point_id] = cp logging.info(charge_point_id) print(charge_point_id) print(path) await websocket.send(json.dumps([2,"222", "GetLocalListVersion", {}])) # 连接断开时从字典删除实例,避免内存泄漏 try: await cp.start() finally: if charge_point_id in connected_charge_points: del connected_charge_points[charge_point_id]
2. 修复心跳事件处理逻辑
删除重复的事件绑定,给异步方法加async关键字,直接用框架内置的连接对象发消息:
class client_main(cp): errors = False if not errors: @on('BootNotification') def on_boot_notitication(self, charge_point_vendor, charge_point_model, charge_point_serial_number,firmware_version, meter_type, **kwargs): return call_result.BootNotificationPayload( status="Accepted", current_time=date_str.replace('+00:00','Z'), interval=60 ) @on('Heartbeat') async def on_heartbeat(self, **kwargs): # 心跳触发时发送业务消息,直接用_connection属性调用websocket的send方法 await self._connection.send(json.dumps([2,"222","GetLocalListVersion",{}])) return call_result.HeartbeatPayload( current_time=datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S')+"Z" )
3. 实现定期主动发消息保活
新增后台异步协程,定期遍历所有在线充电桩发送保活/业务消息:
import asyncio async def periodic_push_task(): while True: # 间隔时间可自行调整,单位秒 await asyncio.sleep(30) # 遍历所有在线充电桩 for cp_id, cp_instance in connected_charge_points.items(): try: # 这里替换为你需要发送的保活或业务消息 await cp_instance._connection.send(json.dumps([2,"custom_msg_id","你的消息类型",{}])) print(f"已向充电桩{cp_id}发送定期消息") except Exception as e: print(f"向充电桩{cp_id}发消息失败:{str(e)}") # 服务启动时同时启动这个后台任务 async def service_start(): # 原有服务启动逻辑,比如启动websocket监听等 asyncio.create_task(periodic_push_task()) # 剩余原有服务启动代码
4. 任意场景主动推送消息
如果需要在其他业务逻辑中给指定充电桩发消息,直接从全局字典取对应实例即可:
async def send_msg_to_charge_point(charge_point_id, msg_content): if charge_point_id not in connected_charge_points: return False, "充电桩不在线" try: cp_instance = connected_charge_points[charge_point_id] await cp_instance._connection.send(json.dumps(msg_content)) return True, "发送成功" except Exception as e: return False, f"发送失败:{str(e)}"
内容的提问来源于stack exchange,提问作者Юлия Доружинская
相关产品推荐
相关产品推荐

