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

Java异步串口通信问题:whenComplete未触发+响应数据异常

串口通信类的问题及优化咨询

我正在开发一个通过串口收发命令的类,遇到两个问题:

  • whenComplete回调始终未被调用;
  • onDataReceived回调中能正确打印响应内容,但execute方法里获取的response变量是乱码,疑似错误内存访问导致。

目前用CountDownLatch获取回调数据,想请教有没有更优的实现方式?


设备类代码

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.io.IOException;

public class Device {
    private static final String port = "/dev/ttyS0";
    private static final int baud = 9600;
    private SerialHelper2 serialPort;
    private final int timeout = 500;
    private CountDownLatch responseLatch;
    private String response;

    public Device() {
    }

    public boolean open() {
        this.serialPort = new SerialHelper2(port, baud) {
            @Override
            protected void onDataReceived(final byte[] data) {
                Device.this.response = new String(data).trim();
                Serial.out.println(Device.this.response); // 这里输出正常
                Device.this.responseLatch.countDown();
            }
        };

        try {
            this.serialPort.open();
            return this.serialPort.isOpen();
        } catch (IOException e) {
            e.printStackTrace();
        }
        return false;
    }

    public CompletableFuture<String> execute(String cmd) {
        // 重置latch
        Device.this.responseLatch = new CountDownLatch(1);
        return CompletableFuture.supplyAsync(() -> {
            Device.this.serialPort.sendHex(cmd);
            try {
                if (responseLatch.await(Device.this.timeout, TimeUnit.MILLISECONDS)) {
                    Serial.out.println(Device.this.response); // 输出乱码,疑似内存访问错误
                    return Device.this.response;
                } else { // 超时
                    throw new TimeoutException(String.format("命令'%s'在%d毫秒后超时", cmd, Device.this.timeout));
                }
            } catch (Exception e) {
                throw new IllegalStateException(e);
            }
        });
    }

    public CompletableFuture<Boolean> init() {
        return this.execute("REQ")
                .thenApply(res -> res.equalsIgnoreCase("RES"));
    }
}

调用代码

device.init()
    .whenComplete((result, ex) -> {
        if (null != ex) {
            Serial.out.println("初始化失败");
        } else {
            Serial.out.println("设备初始化" + (result ? "成功" : "失败"));
        }
    });

问题分析与解决方案

1. 问题根源拆解

whenComplete未触发

  • 类成员responseLatch被多请求共享,当连续调用execute时,新的CountDownLatch会覆盖旧的,导致之前的请求永远无法完成,CompletableFuture一直处于pending状态;
  • 若主线程提前退出,supplyAsync使用的默认线程池(ForkJoinPool)可能被终止,任务未完成就中断。

乱码问题

  • 类成员response无线程同步保护,串口接收线程赋值、supplyAsync线程读取时存在可见性问题,读取线程可能拿到未完全初始化的字符串,导致乱码;
  • 额外风险:new String(data)未指定字符集,若串口数据编码与默认字符集不匹配,也可能引发乱码(但你回调内打印正常,所以核心是线程可见性问题)。

2. 更优实现:用CompletableFuture替代CountDownLatch

去掉类成员的共享变量,每个请求对应独立的CompletableFuture,彻底解决线程安全问题:

import java.nio.charset.StandardCharsets;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicReference;
import java.io.IOException;

public class Device {
    private static final String port = "/dev/ttyS0";
    private static final int baud = 9600;
    private SerialHelper2 serialPort;
    private final int timeout = 500;
    // 原子引用管理当前待处理的请求,避免多请求冲突
    private final AtomicReference<CompletableFuture<String>> pendingFuture = new AtomicReference<>();

    public Device() {
    }

    public boolean open() {
        this.serialPort = new SerialHelper2(port, baud) {
            @Override
            protected void onDataReceived(final byte[] data) {
                // 显式指定串口数据编码,根据实际情况调整
                String response = new String(data, StandardCharsets.UTF_8).trim();
                Serial.out.println(response);
                // 取出并完成当前待处理的Future
                CompletableFuture<String> future = pendingFuture.getAndSet(null);
                if (future != null && !future.isDone()) {
                    future.complete(response);
                }
            }
        };

        try {
            this.serialPort.open();
            return this.serialPort.isOpen();
        } catch (IOException e) {
            e.printStackTrace();
        }
        return false;
    }

    public CompletableFuture<String> execute(String cmd) {
        CompletableFuture<String> future = new CompletableFuture<>();
        // 检查是否已有未完成的请求,避免冲突
        if (!pendingFuture.compareAndSet(null, future)) {
            future.completeExceptionally(new IllegalStateException("已有未完成的串口请求"));
            return future;
        }

        // 发送命令,异常时清理并完成Future
        try {
            serialPort.sendHex(cmd);
        } catch (Exception e) {
            pendingFuture.set(null);
            future.completeExceptionally(e);
            return future;
        }

        // 内置超时逻辑
        CompletableFuture<Void> timeoutTask = CompletableFuture.runAsync(() -> {
            try {
                TimeUnit.MILLISECONDS.sleep(timeout);
                if (!future.isDone()) {
                    pendingFuture.set(null);
                    future.completeExceptionally(new TimeoutException(String.format("命令'%s'在%d毫秒后超时", cmd, timeout)));
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        });

        // 请求完成时取消超时任务
        return future.whenComplete((res, ex) -> timeoutTask.cancel(true));
    }

    public CompletableFuture<Boolean> init() {
        return this.execute("REQ")
                .thenApply(res -> res.equalsIgnoreCase("RES"));
    }
}

3. 关键改进点

  • 用AtomicReference管理待处理请求,避免多线程调用execute时的冲突;
  • 每个请求对应独立的CompletableFuture,消除共享变量的线程安全问题;
  • 显式指定字符集,彻底避免编码不匹配导致的乱码;
  • 内置超时逻辑,无需依赖CountDownLatch,代码更简洁;
  • CompletableFuture状态会被明确更新,whenComplete能确保触发。

4. 额外注意事项

  • 若串口存在多包响应,需在onDataReceived中实现数据拼接逻辑,直到收到完整响应再完成CompletableFuture;
  • 确保SerialHelper2的sendHex方法线程安全,或在execute中加锁保护发送逻辑;
  • 调用异步方法后,若主线程需等待回调执行,可通过CountDownLatch或CompletableFuture.join()实现,避免程序提前退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:14:51