微信第三方平台授权事件:利用STOMP+Kafka广播给多实例Java后端

微信第三方平台在接收到公众号/小程序授权、取消授权等事件时,仅向预设的 URL 推送一次 HTTP 通知。在多实例部署的 Java 后端架构中,若仅一个实例处理该事件,其他实例将无法同步状态。本文通过 接收微信事件 → 写入 Kafka → 广播至所有实例 → 通过 STOMP WebSocket 推送至前端,实现全集群状态一致。

微信事件接收与Kafka投递

使用 Spring Web MVC 接收微信事件并异步写入 Kafka:

package wlkankan.cn.wechat.event;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class WeChatEventHandler {

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @PostMapping("/wechat/notify")
    public String handleWeChatEvent(@RequestBody String xmlPayload) {
        // 验签、解析 XML 略
        // 广播事件到 Kafka topic
        kafkaTemplate.send("wechat.auth.events", xmlPayload);
        return "success";
    }
}

在这里插入图片描述

Kafka消费者广播至本地STOMP

每个后端实例独立消费 Kafka 消息,并通过本地 STOMP Broker 转发:

package wlkankan.cn.wechat.kafka;

import wlkankan.cn.wechat.stomp.StompMessageService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class AuthEventKafkaConsumer {

    @Autowired
    private StompMessageService stompService;

    @KafkaListener(topics = "wechat.auth.events", groupId = "wechat-auth-group")
    public void consumeAuthEvent(String xmlPayload) {
        // 解析事件类型(如 authorized / unauthorized)
        String eventType = parseEventType(xmlPayload);
        String authorizerAppid = extractAuthorizerAppid(xmlPayload);

        // 广播至所有 WebSocket 客户端(含本实例)
        stompService.broadcastAuthEvent(eventType, authorizerAppid, xmlPayload);
    }

    private String parseEventType(String xml) {
        // 简化:实际应解析 <InfoType>
        if (xml.contains("<InfoType><![CDATA[authorized]]></InfoType>")) {
            return "authorized";
        } else if (xml.contains("<InfoType><![CDATA[unauthorized]]></InfoType>")) {
            return "unauthorized";
        }
        return "unknown";
    }

    private String extractAuthorizerAppid(String xml) {
        // 提取 <AuthorizerAppid>
        int start = xml.indexOf("<AuthorizerAppid><![CDATA[") + 24;
        int end = xml.indexOf("]]></AuthorizerAppid>");
        return xml.substring(start, end);
    }
}

STOMP消息服务封装

使用 SimpMessagingTemplate 向订阅路径推送消息:

package wlkankan.cn.wechat.stomp;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.simp.SimpMessagingTemplate;
import org.springframework.stereotype.Service;

@Service
public class StompMessageService {

    @Autowired
    private SimpMessagingTemplate messagingTemplate;

    public void broadcastAuthEvent(String eventType, String appid, String rawXml) {
        AuthEventMessage message = new AuthEventMessage();
        message.setEventType(eventType);
        message.setAuthorizerAppid(appid);
        message.setTimestamp(System.currentTimeMillis());

        // 广播到 /topic/wechat/auth 路径,所有订阅者接收
        messagingTemplate.convertAndSend("/topic/wechat/auth", message);
    }

    public static class AuthEventMessage {
        private String eventType;
        private String authorizerAppid;
        private long timestamp;

        // getters & setters
        public String getEventType() { return eventType; }
        public void setEventType(String eventType) { this.eventType = eventType; }
        public String getAuthorizerAppid() { return authorizerAppid; }
        public void setAuthorizerAppid(String authorizerAppid) { this.authorizerAppid = authorizerAppid; }
        public long getTimestamp() { return timestamp; }
        public void setTimestamp(long timestamp) { this.timestamp = timestamp; }
    }
}

STOMP Broker配置

启用简单内存 Broker 支持广播(生产环境建议用 RabbitMQ 或 Redis):

package wlkankan.cn.wechat.config;

import org.springframework.context.annotation.Configuration;
import org.springframework.messaging.simp.config.MessageBrokerRegistry;
import org.springframework.web.socket.config.annotation.EnableWebSocketMessageBroker;
import org.springframework.web.socket.config.annotation.StompEndpointRegistry;
import org.springframework.web.socket.config.annotation.WebSocketMessageBrokerConfigurer;

@Configuration
@EnableWebSocketMessageBroker
public class StompConfig implements WebSocketMessageBrokerConfigurer {

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        registry.addEndpoint("/ws").setAllowedOriginPatterns("*").withSockJS();
    }

    @Override
    public void configureMessageBroker(MessageBrokerRegistry registry) {
        registry.enableSimpleBroker("/topic");
        registry.setApplicationDestinationPrefixes("/app");
    }
}

前端订阅示例(简略)

前端通过 SockJS 订阅 /topic/wechat/auth

var socket = new SockJS('/ws');
var stompClient = Stomp.over(socket);
stompClient.connect({}, function () {
    stompClient.subscribe('/topic/wechat/auth', function (msg) {
        var event = JSON.parse(msg.body);
        console.log('Received auth event:', event.eventType, event.authorizerAppid);
        // 更新 UI 或缓存
    });
});

该方案确保任意实例收到微信授权事件后,全集群所有节点(及前端)实时同步状态,避免因实例隔离导致的数据不一致,适用于 SaaS 化微信第三方平台管理系统。

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐