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

如何在指定时长内读取InputStream并避免返回后追加额外订单?

问题分析与修复方案

你的代码存在三个核心问题,导致返回结果后仍可能追加数据:

  1. 中断检查时机滞后:原代码在解析完订单后才检查中断状态,此时即使线程被中断,已解析的订单仍会被添加到列表
  2. 阻塞IO不响应中断:bufferedReader.readLine()是阻塞操作,线程被中断时不会主动退出阻塞,会继续等待输入
  3. 集合线程不安全:LinkedList不是线程安全集合,主线程返回列表时,读取线程可能仍在修改它,导致数据不一致

修复后代码

方案1:关闭流唤醒阻塞(适合允许关闭外部传入流的场景)

InputStream source;
ObjectMapper objectMapper;

public OrderStreamReader(InputStream source, ObjectMapper objectMapper) {
    this.source = source;
    this.objectMapper = objectMapper;
}

public List<OrderStream> get(Duration maxTime) {
    ExecutorService executor = Executors.newSingleThreadExecutor();
    ReaderOrdersStream readerOrder = new ReaderOrdersStream();
    Future<?> future = executor.submit(readerOrder);

    try {
        future.get(maxTime.toMillis(), TimeUnit.MILLISECONDS);
    } catch (TimeoutException e) {
        future.cancel(true);
        // 关闭流唤醒阻塞的readLine
        try {
            source.close();
        } catch (IOException ex) {
            // 忽略流关闭异常,或根据业务需求处理
        }
    } catch (InterruptedException | ExecutionException e) {
        throw new RuntimeException(e);
    } finally {
        executor.shutdownNow();
        // 等待线程完全终止,避免返回后仍有数据写入
        try {
            if (!executor.awaitTermination(500, TimeUnit.MILLISECONDS)) {
                source.close();
            }
        } catch (InterruptedException | IOException ex) {
            Thread.currentThread().interrupt();
        }
    }

    return readerOrder.getOrders();
}

private class ReaderOrdersStream implements Runnable {
    // 使用线程安全集合避免并发修改
    private final List<OrderStream> orders = Collections.synchronizedList(new LinkedList<>());

    @Override
    public void run() {
        try (BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(source))) {
            String rawOrder;
            while (!Thread.currentThread().isInterrupted()) {
                // 循环开头先检查中断状态
                if (Thread.currentThread().isInterrupted()) {
                    break;
                }

                rawOrder = bufferedReader.readLine();
                if (rawOrder == null) {
                    break;
                }

                // 解析前再次检查中断
                if (Thread.currentThread().isInterrupted()) {
                    break;
                }

                var currentOrder = objectMapper.readValue(rawOrder, OrderStream.class);

                // 添加前最后一次检查中断
                if (!Thread.currentThread().isInterrupted()) {
                    orders.add(currentOrder);
                }
            }
        } catch (IOException e) {
            // 忽略流关闭导致的异常,视为正常终止
            String msg = e.getMessage();
            if (!(msg != null && (msg.contains("Stream closed") || msg.contains("Socket closed")))) {
                throw new RuntimeException(e);
            }
        }
    }

    // 返回不可修改的副本,避免后续意外修改
    public List<OrderStream> getOrders() {
        return Collections.unmodifiableList(new ArrayList<>(orders));
    }
}

方案2:使用可中断输入流(不关闭原流的场景)

如果不能关闭外部传入的InputStream,可以自定义可中断输入流包装类,让阻塞IO响应中断:

// 自定义可中断输入流
class InterruptibleInputStream extends InputStream {
    private final InputStream delegate;

    public InterruptibleInputStream(InputStream delegate) {
        this.delegate = delegate;
    }

    @Override
    public int read() throws IOException {
        checkInterrupt();
        return delegate.read();
    }

    @Override
    public int read(byte[] b) throws IOException {
        checkInterrupt();
        return delegate.read(b);
    }

    @Override
    public int read(byte[] b, int off, int len) throws IOException {
        checkInterrupt();
        return delegate.read(b, off, len);
    }

    private void checkInterrupt() throws IOException {
        if (Thread.currentThread().isInterrupted()) {
            throw new IOException("Thread interrupted");
        }
    }

    // 重写其他read相关方法,均添加中断检查
}

然后在ReaderOrdersStream中替换输入流:

try (BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(new InterruptibleInputStream(source)))) {
    // 循环逻辑同方案1
}

关键修改说明

  1. 提前中断检查:在循环开头、读取行后、解析后多次检查中断状态,确保中断后立即停止所有操作
  2. 线程安全集合:用Collections.synchronizedList包装集合,最后返回不可修改副本,彻底避免并发修改问题
  3. 唤醒阻塞IO:通过关闭流或自定义可中断流,解决readLine()不响应中断的问题
  4. 等待线程终止:在finally块中等待线程完全停止后再返回结果,确保返回的列表状态稳定

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 07:47:05