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

如何将DicomInputStream拆分为两个缓冲InputStream供双API消费?

问题分析

你当前代码的核心问题是**PipedInputStream默认缓冲区过小(仅1024字节)导致的死锁**:

  • 主线程通过TeeInputStream向PipedOutputStream写入数据时,一旦缓冲区填满,主线程会被阻塞,等待消费者读取释放空间。
  • 如果Future中的storageApiClient.forwardInstance启动不及时,或者读取速度跟不上主线程的写入速度,就会形成双向阻塞:主线程等缓冲区空,Future等主线程写数据,最终导致Future永远无法完成。

同时,SplittableInputStream的挂起问题也是因为内存缓冲无法承载大文件的字节量,导致内存溢出或阻塞。


可行解决方案

方案一:临时文件缓冲(推荐,大文件友好)

将原始流先写入临时文件,再给两个API分别提供独立的文件输入流。这种方式不占用大量内存,且两个API可以各自独立读取,完全避免线程阻塞问题。

import java.io.{File, FileInputStream, FileOutputStream}
import java.nio.file.Files
import scala.concurrent.{Await, Future}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.Duration

// 创建自动删除的临时文件
val tempFile = Files.createTempFile("dicom_temp_", ".dcm").toFile
tempFile.deleteOnExit()

// 将原始流写入临时文件
val fileOut = new FileOutputStream(tempFile)
try {
  dicomInputStream.transferTo(fileOut)
} finally {
  dicomInputStream.close()
  fileOut.close()
}

// 启动两个线程分别处理两个API请求
val storageFuture = Future {
  val fileIn = new FileInputStream(tempFile)
  try {
    storageApiClient.forwardInstance(fileIn, calledAet, "123", instanceUID)
  } finally {
    fileIn.close()
  }
}

val healthcareFuture = Future {
  val fileIn = new FileInputStream(tempFile)
  try {
    healthcareApiClient.stowRs(fileIn, calledAet)
  } finally {
    fileIn.close()
  }
}

// 等待两个任务全部完成
Await.result(storageFuture.zip(healthcareFuture), Duration.Inf)

方案二:阻塞队列缓冲(无磁盘IO,内存可控)

通过BlockingQueue实现字节块的中转,主线程读取原始流并将字节块放入队列,两个消费者线程从队列取数据并封装成自定义InputStream供API消费。可以通过限制队列容量避免内存溢出。

import java.io.InputStream
import java.util.concurrent.ArrayBlockingQueue
import scala.concurrent.{Await, Future}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.Duration

// 定义结束标记
private val END_MARKER = new Array[Byte](0)
// 队列容量设为100个8KB字节块,可根据内存情况调整
val byteQueue = new ArrayBlockingQueue[Array[Byte]](100)

// 自定义InputStream,从阻塞队列读取数据
class QueueInputStream(queue: ArrayBlockingQueue[Array[Byte]]) extends InputStream {
  private var currentBuffer: Array[Byte] = _
  private var currentPos: Int = 0

  override def read(): Int = {
    if (currentBuffer == null || currentPos >= currentBuffer.length) {
      currentBuffer = queue.take()
      if (currentBuffer == END_MARKER) return -1
      currentPos = 0
    }
    val byte = currentBuffer(currentPos) & 0xFF
    currentPos += 1
    byte
  }

  override def read(b: Array[Byte], off: Int, len: Int): Int = {
    if (currentBuffer == null || currentPos >= currentBuffer.length) {
      currentBuffer = queue.take()
      if (currentBuffer == END_MARKER) return -1
      currentPos = 0
    }
    val readLen = math.min(len, currentBuffer.length - currentPos)
    System.arraycopy(currentBuffer, currentPos, b, off, readLen)
    currentPos += readLen
    readLen
  }
}

// 主线程:读取原始流,将字节块放入队列
val producerFuture = Future {
  val buffer = new Array[Byte](8192) // 每次读取8KB
  var bytesRead = 0
  try {
    while ({ bytesRead = dicomInputStream.read(buffer); bytesRead != -1 }) {
      val copy = new Array[Byte](bytesRead)
      System.arraycopy(buffer, 0, copy, 0, bytesRead)
      byteQueue.put(copy)
    }
  } finally {
    byteQueue.put(END_MARKER)
    dicomInputStream.close()
  }
}

// 两个消费者线程分别处理API请求
val storageFuture = Future {
  val inputStream = new QueueInputStream(byteQueue)
  try {
    storageApiClient.forwardInstance(inputStream, calledAet, "123", instanceUID)
  } finally {
    inputStream.close()
  }
}

val healthcareFuture = Future {
  val inputStream = new QueueInputStream(byteQueue)
  try {
    healthcareApiClient.stowRs(inputStream, calledAet)
  } finally {
    inputStream.close()
  }
}

// 等待所有任务完成
Await.result(producerFuture.zip(storageFuture).zip(healthcareFuture), Duration.Inf)

方案三:调整PipedInputStream缓冲区(仅适合小文件)

如果坚持使用PipedStream,可以增大缓冲区大小降低阻塞概率,但仍存在因消费速度不匹配导致死锁的风险,不适合300MB以上的大文件。

import java.io.{PipedInputStream, PipedOutputStream, TeeInputStream}
import scala.concurrent.{Await, Future}
import scala.concurrent.ExecutionContext.Implicits.global
import scala.concurrent.duration.Duration

// 增大缓冲区到64KB
val pipedIn = new PipedInputStream(65536)
val tee = new TeeInputStream(dicomInputStream, new PipedOutputStream(pipedIn))

val storageFuture = Future {
  try {
    storageApiClient.forwardInstance(pipedIn, calledAet, "123", instanceUID)
  } finally {
    pipedIn.close()
  }
}

try {
  healthcareApiClient.stowRs(tee, calledAet)
} finally {
  tee.close()
}

Await.result(storageFuture, Duration.Inf)

内容的提问来源于stack exchange,提问作者Dominic Bou-Samra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 23:49:51