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

Milo OPC UA客户端重连后无法重建订阅问题求助

解决Milo客户端重连后订阅丢失且无法重建的问题

核心问题分析

Milo客户端在服务器重启重连后,订阅转移失败导致订阅丢失,直接在onSessionActive或onSubscriptionTransferFailed回调中同步创建订阅会阻塞事件循环线程,导致createSubscription返回null Future,代码停滞,且订阅无法重建。

解决方案步骤

  1. 异步执行订阅重建逻辑
    不要在Milo的回调线程(事件循环线程)中直接执行订阅创建操作,否则会阻塞客户端状态机(FSM),导致后续操作无法进行。将重建逻辑放到客户端的异步线程池中执行。

  2. 彻底清理旧订阅状态
    调用createSubscription前,必须先彻底关闭并清理SubscriptionManager的旧状态,避免残留资源干扰新订阅创建:

    • 先调用subscriptionManager.shutdown()并等待完成,确保旧订阅资源完全释放
    • 再调用subscriptionManager.clearSubscriptions()清理订阅记录
  3. 提前保存订阅配置
    在onSubscriptionTransferFailed回调中记录所有失败订阅的关键信息(扫描率、监控项列表、回调逻辑等),方便后续重建时复用。

代码实现示例

import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
import org.eclipse.milo.opcua.sdk.client.api.subscription.UaMonitoredItem;
import org.eclipse.milo.opcua.sdk.client.api.subscription.UaSubscription;
import org.eclipse.milo.opcua.sdk.client.subscription.SubscriptionManager;
import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode;
import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn;
import org.eclipse.milo.opcua.stack.core.types.structured.MonitoredItemCreateRequest;

// 保存订阅配置的辅助类
static class SubscriptionConfig {
    private final double scanRate;
    private final List<UaMonitoredItem> monitoredItems;

    public SubscriptionConfig(double scanRate, List<UaMonitoredItem> monitoredItems) {
        this.scanRate = scanRate;
        this.monitoredItems = monitoredItems;
    }

    public double getScanRate() {
        return scanRate;
    }

    public List<UaMonitoredItem> getMonitoredItems() {
        return monitoredItems;
    }
}

// 客户端初始化时注册回调
OpcUaClient client = ...; // 初始化你的客户端
List<SubscriptionConfig> failedSubscriptions = new ArrayList<>();

client.getSubscriptionManager().addSubscriptionTransferFailedListener((subscription, statusCode) -> {
    // 记录失败的订阅配置
    failedSubscriptions.add(new SubscriptionConfig(subscription.getPublishInterval(), subscription.getMonitoredItems()));
});

client.getSessionListener().onSessionActive(session -> {
    // 异步执行订阅重建
    CompletableFuture.runAsync(() -> {
        SubscriptionManager subscriptionManager = client.getSubscriptionManager();
        try {
            // 关闭旧订阅管理器并等待完成
            subscriptionManager.shutdown().get(5, TimeUnit.SECONDS);
            // 清理旧订阅记录
            subscriptionManager.clearSubscriptions();

            // 逐个重建订阅和监控项
            for (SubscriptionConfig config : failedSubscriptions) {
                // 创建新订阅,等待结果返回
                UaSubscription newSubscription = subscriptionManager.createSubscription(config.getScanRate()).get(10, TimeUnit.SECONDS);
                
                // 重新创建监控项
                List<MonitoredItemCreateRequest> requests = new ArrayList<>();
                for (UaMonitoredItem item : config.getMonitoredItems()) {
                    requests.add(new MonitoredItemCreateRequest(
                        item.getNodeId(),
                        item.getAttributeId(),
                        null,
                        item.getMonitoringParameters(),
                        true,
                        null
                    ));
                }

                newSubscription.createMonitoredItems(
                    TimestampsToReturn.Both,
                    requests,
                    (s, t, items) -> {
                        // 处理监控项创建结果,比如恢复回调逻辑
                        for (UaMonitoredItem item : items) {
                            if (item.getStatusCode().isGood()) {
                                // 重新绑定之前的数值变化回调
                                item.setValueConsumer((it, value) -> {
                                    // 你的业务逻辑
                                });
                            }
                        }
                    }
                );
            }

            // 清空失败订阅列表,避免重复重建
            failedSubscriptions.clear();
        } catch (Exception e) {
            e.printStackTrace();
            // 可添加重试逻辑或日志报警
        }
    }, client.getExecutorService());
});

关键注意事项

  • 避免阻塞回调线程:Milo的回调线程负责处理状态机事件,阻塞会导致客户端无法响应后续的会话或订阅事件。
  • 等待shutdown完成:shutdown()是异步操作,必须等待其完成后再执行清理和创建操作,否则会出现状态不一致。
  • 恢复监控项回调:重建监控项后,需要重新绑定之前的数值变化回调,确保业务逻辑正常运行。

内容的提问来源于stack exchange,提问作者obeliksz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 13:08:08