如何将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
相关产品推荐
相关产品推荐

