Scala Requests模块并发请求疑问:现有代码是串行还是并发?
问题解答
1. 代码执行方式判断
你给出的代码是串行请求。
原因:urls.map会按顺序遍历每个URL,逐个调用getBytesFromUrl方法;而requests.get(url).bytes是同步阻塞操作,必须等待当前请求的响应完全获取并转为字节数组后,才会处理下一个URL,全程没有并发执行。
2. 实现并发请求的正确方式
下面提供几种Scala中常用的并发实现方案:
方案一:使用scala.concurrent.Future(推荐)
利用Scala标准库的Future异步发起每个请求,最后统一等待所有请求完成:
import scala.concurrent.{Await, Future} import scala.concurrent.ExecutionContext.Implicits.global import scala.concurrent.duration._ object Runner extends App { def getBytesFromUrl(url: String): Array[Byte] = { requests.get(url).bytes } val urls = Seq( "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg", "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg", "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg" ) // 为每个URL创建异步执行的Future val futures: Seq[Future[Array[Byte]]] = urls.map(url => Future { getBytesFromUrl(url) }) // 等待所有Future完成,获取结果数组 val result: Seq[Array[Byte]] = Await.result(Future.sequence(futures), 10.seconds) }
说明:
Future { ... }会把同步的getBytesFromUrl放到全局执行上下文的线程池中异步执行,多个请求可同时发起。Future.sequence将Seq[Future[T]]转换为Future[Seq[T]],方便一次性等待所有请求完成。- 生产环境中如果是异步场景,建议用
onComplete或flatMap等非阻塞方式处理结果,避免Await.result的阻塞操作。
方案二:使用并行集合(ParSeq)
Scala的并行集合会自动将遍历操作分配到多个线程执行:
object Runner extends App { def getBytesFromUrl(url: String): Array[Byte] = { requests.get(url).bytes } val urls = Seq( "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg", "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg", "https://cdn.pixabay.com/photo/2014/02/27/16/10/flowers-276014__340.jpg" ) // 转为并行集合,并行执行map操作 val result: Seq[Array[Byte]] = urls.par.map(url => getBytesFromUrl(url)).seq }
说明:
.par将普通Seq转为ParSeq,遍历操作会并行执行。- 最后用
.seq转回普通Seq,方便后续统一处理。 - 这种方式简单易用,但并行度由Scala自动管理,适合简单场景。
注意事项
- 无论哪种方案,都要注意控制并发数,避免一次性发起过多请求导致目标服务器拒绝或本地资源耗尽。可以通过自定义线程池或限流工具来控制。
- 如果使用
Future,建议自定义ExecutionContext而非全局默认上下文,以便更好地管控线程资源:
import java.util.concurrent.Executors import scala.concurrent.ExecutionContext // 自定义固定大小的线程池,比如限制5个线程 val customEC = ExecutionContext.fromExecutor(Executors.newFixedThreadPool(5)) // 创建Future时指定自定义上下文 val futures = urls.map(url => Future(getBytesFromUrl(url))(customEC))
内容的提问来源于stack exchange,提问作者Shivam Sahil
相关产品推荐
相关产品推荐

