微信第三方平台授权事件:利用STOMP+Kafka广播给多实例Java后端
·
微信第三方平台授权事件:利用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 化微信第三方平台管理系统。
更多推荐
所有评论(0)