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

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,提问作者Юлия Доружинская

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 12:57:03