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

如何在Camunda中传递带Socket的Java对象至Service Task执行

可行!复用Socket绑定的Presse/Bohrer实例的解决方案

你的核心问题是Camunda默认会为每个服务任务创建新的JavaDelegate实例(通过无参构造),导致无法复用你在TCP连接时创建的、已经绑定Socket的Presse/Bohrer对象。下面是几种落地的解决方案,按实现复杂度排序:

方案1:将实例存入流程变量(最直接)

这个方案的核心是:在TCP服务器创建Presse/Bohrer实例后,把它存入对应的Camunda流程实例变量,然后在服务任务的Delegate中直接取出复用。

步骤1:拆分Presse的业务逻辑与Delegate

首先把Presse的业务代码和Camunda的JavaDelegate解耦,让Presse专注于Socket通信,用一个专门的Delegate来调用它:

// 修改后的Presse类,去掉JavaDelegate实现
public class Presse {
    String param = "5";
    private HexToByteConverter hexToByteConverter = new HexToByteConverter();
    Socket socket;
    InputStream in;
    OutputStream out;

    // 带Socket参数的构造方法(保留)
    public Presse(Socket clientSocket) throws IOException {
        this.socket = clientSocket;
        this.in = socket.getInputStream();
        this.out = socket.getOutputStream();
    }

    // 业务方法:发送消息
    public void sendMessage(String durationParam) throws IOException {
        // 根据参数生成对应的字节数组
        byte[] pressen1Hex = hexToByteConverter.hexStringToByteArray(
            "33333333003d0064000600000004004001c9c78900010000006e00000000000000000000000000010000000000140000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000"+durationParam
        );
        byte[] pressen2Hex = hexToByteConverter.hexStringToByteArray(
            "3333333300400065000a00000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000"
        );
        byte[] pressen3Hex = hexToByteConverter.hexStringToByteArray(
            "3333333300400065001400000000004001c9c6e900010000006e000000000000000000000000000100000000001e0000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000"
        );

        // 发送消息
        out.write(pressen1Hex);
        out.write(pressen2Hex);
        out.write(pressen3Hex);
        out.flush();
    }
}

// 专门的Camunda服务任务Delegate
public class PresseTaskDelegate implements JavaDelegate {
    @Override
    public void execute(DelegateExecution execution) throws Exception {
        // 从流程变量中取出已绑定Socket的Presse实例
        Presse presse = (Presse) execution.getVariable("presseInstance");
        if (presse == null) {
            throw new RuntimeException("No Presse instance found for this process!");
        }

        // 获取流程参数
        String durationParam = (String) execution.getVariable("durationglobal");
        // 执行Socket发送逻辑
        presse.sendMessage(durationParam);
    }
}

步骤2:在TCP服务器中存储实例到流程变量

在ServerAzyklisch创建Presse实例后,通过Camunda的RuntimeService把实例存入对应的流程变量(需要确保你能拿到目标流程实例ID,比如从客户端消息、配置或提前启动流程获取):

public class ServerAzyklisch implements Runnable, JavaDelegate {
    // ... 保留原有代码

    public void run() {
        // ... 保留原有循环逻辑
        threadPool.submit(() -> {
            try {
                Socket finalSocket = socket;
                Presse p = new Presse(finalSocket);
                Bohrer b = new Bohrer(finalSocket);

                // 关键:将实例存入流程变量
                // 这里需要替换为你的实际流程实例ID,比如从客户端接收的消息中获取
                String processInstanceId = "your-target-process-instance-id";
                RuntimeService runtimeService = ProcessEngines.getDefaultProcessEngine().getRuntimeService();
                runtimeService.setVariable(processInstanceId, "presseInstance", p);
                runtimeService.setVariable(processInstanceId, "bohrerInstance", b);

            } catch (IOException e) {
                e.printStackTrace();
            }
        });
        // ... 保留原有代码
    }

    // ... 保留原有代码
}

步骤3:更新BPMN模型

把「Pressen」服务任务的Java Class配置为com.example.workflow.PresseTaskDelegate(而不是原来的Presse),同理「Bohren」任务配置对应的BohrerTaskDelegate。


方案2:全局缓存存储实例(适合多流程/多客户端场景)

如果你的系统有多个流程实例或多个客户端连接,可以用线程安全的全局缓存存储实例,通过标识(比如Socket ID、客户端ID)关联到流程变量:

步骤1:创建全局缓存类

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

public class DeviceInstanceCache {
    // 线程安全的缓存,key可以是Socket的唯一标识或客户端ID
    private static final Map<String, Presse> PRESSE_CACHE = new ConcurrentHashMap<>();
    private static final Map<String, Bohrer> BOHRER_CACHE = new ConcurrentHashMap<>();

    public static void putPresse(String key, Presse presse) {
        PRESSE_CACHE.put(key, presse);
    }

    public static Presse getPresse(String key) {
        return PRESSE_CACHE.get(key);
    }

    public static void putBohrer(String key, Bohrer bohrer) {
        BOHRER_CACHE.put(key, bohrer);
    }

    public static Bohrer getBohrer(String key) {
        return BOHRER_CACHE.get(key);
    }

    // 客户端断开时清理实例,避免内存泄漏
    public static void removePresse(String key) {
        PRESSE_CACHE.remove(key);
    }
}

步骤2:TCP服务器中存入缓存+流程变量

threadPool.submit(() -> {
    try {
        Socket finalSocket = socket;
        // 用Socket的唯一标识作为缓存key
        String socketKey = finalSocket.getInetAddress() + ":" + finalSocket.getPort();
        
        Presse p = new Presse(finalSocket);
        Bohrer b = new Bohrer(finalSocket);
        
        // 存入缓存
        DeviceInstanceCache.putPresse(socketKey, p);
        DeviceInstanceCache.putBohrer(socketKey, b);
        
        // 将key存入流程变量
        String processInstanceId = "your-target-process-instance-id";
        RuntimeService runtimeService = ProcessEngines.getDefaultProcessEngine().getRuntimeService();
        runtimeService.setVariable(processInstanceId, "presseSocketKey", socketKey);
        runtimeService.setVariable(processInstanceId, "bohrerSocketKey", socketKey);

        // 监听客户端断开,清理缓存
        finalSocket.close();
        DeviceInstanceCache.removePresse(socketKey);
    } catch (IOException e) {
        e.printStackTrace();
    }
});

步骤3:Delegate中从缓存取实例

public class PresseTaskDelegate implements JavaDelegate {
    @Override
    public void execute(DelegateExecution execution) throws Exception {
        String socketKey = (String) execution.getVariable("presseSocketKey");
        Presse presse = DeviceInstanceCache.getPresse(socketKey);
        if (presse == null) {
            throw new RuntimeException("Presse instance not found for socket: " + socketKey);
        }

        String durationParam = (String) execution.getVariable("durationglobal");
        presse.sendMessage(durationParam);
    }
}

关键注意事项

  1. 线程安全:缓存必须用线程安全的集合(如ConcurrentHashMap),避免多线程环境下的并发问题
  2. 资源清理:客户端断开连接时,要及时从缓存/流程变量中移除实例,防止内存泄漏
  3. 流程实例关联:确保每个流程实例对应正确的设备实例,避免不同流程复用错误的Socket连接
  4. 异常处理:在Delegate中要处理实例不存在的情况,避免流程中断

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 22:27:41