PageSourceSearch

https://zio.dev/assets/js/61a831a0.7ad30a78.js

js zio.dev collected 2026-10-03 19:29:24 UTC 15,477 bytes, 1 lines download raw bytes

1"use strict";(globalThis.webpackChunkzio_site||=[]).push([[17796],{21282(e,n,u){u.r(n),u.d(n,{assets:()=>r,contentTitle:()=>l,default:()=>c,frontMatter:()=>a,metadata:()=>i,toc:()=>o});const i=JSON.parse('{"id":"reference/concurrency/queue","title":"Queue","description":"Queue is a lightweight in-memory queue built on ZIO with composable and transparent back-pressure. It is fully asynchronous (no locks or blocking), purely-functional and type-safe.","source":"@site/versioned_docs/version-1.0.18/reference/concurrency/queue.md","sourceDirName":"reference/concurrency","slug":"/reference/concurrency/queue","permalink":"/1.0.18/reference/concurrency/queue","draft":false,"unlisted":false,"editUrl":"https://github.com/zio/zio/edit/series/2.x/versioned_docs/version-1.0.18/reference/concurrency/queue.md","tags":[],"version":"1.0.18","frontMatter":{"id":"queue","title":"Queue"},"sidebar":"overview_sidebar","previous":{"title":"Promise","permalink":"/1.0.18/reference/concurrency/promise"},"next":{"title":"Hub","permalink":"/1.0.18/reference/concurrency/hub"}}');var t=u(74848),s=u(28453);const a={id:"queue",title:"Queue"},l=void 0,r={},o=[{value:"Creating a queue",id:"creating-a-queue",level:2},{value:"Adding items to a queue",id:"adding-items-to-a-queue",level:2},{value:"Consuming Items from a Queue",id:"consuming-items-from-a-queue",level:2},{value:"Shutting Down a Queue",id:"shutting-down-a-queue",level:2},{value:"Transforming queues",id:"transforming-queues",level:2},{value:"ZQueue#map",id:"zqueuemap",level:3},{value:"ZQueue#mapM",id:"zqueuemapm",level:3},{value:"ZQueue#contramapM",id:"zqueuecontramapm",level:3},{value:"ZQueue#bothWith",id:"zqueuebothwith",level:3},{value:"Additional Resources",id:"additional-resources",level:2}];function d(e){const n={a:"a",code:"code",h2:"h2",h3:"h3",li:"li",p:"p",pre:"pre",ul:"ul",...(0,s.R)(),...e.components};return(0,t.jsxs)(t.Fragment,{children:[(0,t.jsxs)(n.p,{children:[(0,t.jsx)(n.code,{children:"Queue"})," is a lightweight in-memory queue built on ZIO with composable and transparent back-pressure. It is fully asynchronous (no locks or blocking), purely-functional and type-safe."]}),"\n",(0,t.jsxs)(n.p,{children:["A ",(0,t.jsx)(n.code,{children:"Queue[A]"})," contains values of type ",(0,t.jsx)(n.code,{children:"A"})," and has two basic operations: ",(0,t.jsx)(n.code,{children:"offer"}),", which places an ",(0,t.jsx)(n.code,{children:"A"})," in the ",(0,t.jsx)(n.code,{children:"Queue"}),", and ",(0,t.jsx)(n.code,{children:"take"})," which removes and returns the oldest value in the ",(0,t.jsx)(n.code,{children:"Queue"}),"."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"import zio._\n\nval res: UIO[Int] = for {\n  queue <- Queue.bounded[Int](100)\n  _ <- queue.offer(1)\n  v1 <- queue.take\n} yield v1\n"})}),"\n",(0,t.jsx)(n.h2,{id:"creating-a-queue",children:"Creating a queue"}),"\n",(0,t.jsxs)(n.p,{children:["A ",(0,t.jsx)(n.code,{children:"Queue"})," can be bounded (with a limited capacity) or unbounded."]}),"\n",(0,t.jsx)(n.p,{children:"There are several strategies to process new values when the queue is full:"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:["The default ",(0,t.jsx)(n.code,{children:"bounded"})," queue is back-pressured: when full, any offering fiber will be suspended until the queue is able to add the item;"]}),"\n",(0,t.jsxs)(n.li,{children:["A ",(0,t.jsx)(n.code,{children:"dropping"})," queue will drop new items when the queue is full;"]}),"\n",(0,t.jsxs)(n.li,{children:["A ",(0,t.jsx)(n.code,{children:"sliding"})," queue will drop old items when the queue is full."]}),"\n"]}),"\n",(0,t.jsx)(n.p,{children:"To create a back-pressured bounded queue:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val boundedQueue: UIO[Queue[Int]] = Queue.bounded[Int](100)\n"})}),"\n",(0,t.jsx)(n.p,{children:"To create a dropping queue:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val droppingQueue: UIO[Queue[Int]] = Queue.dropping[Int](100)\n"})}),"\n",(0,t.jsx)(n.p,{children:"To create a sliding queue:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val slidingQueue: UIO[Queue[Int]] = Queue.sliding[Int](100)\n"})}),"\n",(0,t.jsx)(n.p,{children:"To create an unbounded queue:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val unboundedQueue: UIO[Queue[Int]] = Queue.unbounded[Int]\n"})}),"\n",(0,t.jsx)(n.h2,{id:"adding-items-to-a-queue",children:"Adding items to a queue"}),"\n",(0,t.jsxs)(n.p,{children:["The simplest way to add a value to the queue is ",(0,t.jsx)(n.code,{children:"offer"}),":"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val res1: UIO[Unit] = for {\n  queue <- Queue.bounded[Int](100)\n  _ <- queue.offer(1)\n} yield ()\n"})}),"\n",(0,t.jsxs)(n.p,{children:["When using a back-pressured queue, offer might suspend if the queue is full: you can use ",(0,t.jsx)(n.code,{children:"fork"})," to wait in a different fiber."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val res2: UIO[
1Unit] = for {\n  queue <- Queue.bounded[Int](1)\n  _ <- queue.offer(1)\n  f <- queue.offer(1).fork // will be suspended because the queue is full\n  _ <- queue.take\n  _ <- f.join\n} yield ()\n"})}),"\n",(0,t.jsxs)(n.p,{children:["It is also possible to add multiple values at once with ",(0,t.jsx)(n.code,{children:"offerAll"}),":"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val res3: UIO[Unit] = for {\n  queue <- Queue.bounded[Int](100)\n  items = Range.inclusive(1, 10).toList\n  _ <- queue.offerAll(items)\n} yield ()\n"})}),"\n",(0,t.jsx)(n.h2,{id:"consuming-items-from-a-queue",children:"Consuming Items from a Queue"}),"\n",(0,t.jsxs)(n.p,{children:["The ",(0,t.jsx)(n.code,{children:"take"})," operation removes the oldest item from the queue and returns it. If the queue is empty, this will suspend, and resume only when an item has been added to the queue. As with ",(0,t.jsx)(n.code,{children:"offer"}),", you can use ",(0,t.jsx)(n.code,{children:"fork"})," to wait for the value in a different fiber."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:'val oldestItem: UIO[String] = for {\n  queue <- Queue.bounded[String](100)\n  f <- queue.take.fork // will be suspended because the queue is empty\n  _ <- queue.offer("something")\n  v <- f.join\n} yield v\n'})}),"\n",(0,t.jsxs)(n.p,{children:["You can consume the first item with ",(0,t.jsx)(n.code,{children:"poll"}),". If the queue is empty you will get ",(0,t.jsx)(n.code,{children:"None"}),", otherwise the top item will be returned wrapped in ",(0,t.jsx)(n.code,{children:"Some"}),"."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val polled: UIO[Option[Int]] = for {\n  queue <- Queue.bounded[Int](100)\n  _ <- queue.offer(10)\n  _ <- queue.offer(20)\n  head <- queue.poll\n} yield head\n"})}),"\n",(0,t.jsxs)(n.p,{children:["You can consume multiple items at once with ",(0,t.jsx)(n.code,{children:"takeUpTo"}),". If the queue doesn't have enough items to return, it will return all the items without waiting for more offers."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val taken: UIO[List[Int]] = for {\n  queue <- Queue.bounded[Int](100)\n  _ <- queue.offer(10)\n  _ <- queue.offer(20)\n  list  <- queue.takeUpTo(5)\n} yield list\n"})}),"\n",(0,t.jsxs)(n.p,{children:["Similarly, you can get all items at once with ",(0,t.jsx)(n.code,{children:"takeAll"}),". It also returns without waiting (an empty list if the queue is empty)."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val all: UIO[List[Int]] = for {\n  queue <- Queue.bounded[Int](100)\n  _ <- queue.offer(10)\n  _ <- queue.offer(20)\n  list  <- queue.takeAll\n} yield list\n"})}),"\n",(0,t.jsx)(n.h2,{id:"shutting-down-a-queue",children:"Shutting Down a Queue"}),"\n",(0,t.jsxs)(n.p,{children:["It is possible with ",(0,t.jsx)(n.code,{children:"shutdown"})," to interrupt all the fibers that are suspended on ",(0,t.jsx)(n.code,{children:"offer*"})," or ",(0,t.jsx)(n.code,{children:"take*"}),". It will also empty the queue and make all future calls to ",(0,t.jsx)(n.code,{children:"offer*"})," and ",(0,t.jsx)(n.code,{children:"take*"})," terminate immediately."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val takeFromShutdownQueue: UIO[Unit] = for {\n  queue <- Queue.bounded[Int](3)\n  f <- queue.take.fork\n  _ <- queue.shutdown // will interrupt f\n  _ <- f.join // Will terminate\n} yield ()\n"})}),"\n",(0,t.jsxs)(n.p,{children:["You can use ",(0,t.jsx)(n.code,{children:"awaitShutdown"})," to execute an effect when the queue is shut down. This will wait until the queue is shut down. If the queue is already shutdown, it will resume right away."]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val awaitShutdown: UIO[Unit] = for {\n  queue <- Queue.bounded[Int](3)\n  p <- Promise.make[Nothing, Boolean]\n  f <- queue.awaitShutdown.fork\n  _ <- queue.shutdown\n  _ <- f.join\n} yield ()\n"})}),"\n",(0,t.jsx)(n.h2,{id:"transforming-queues",children:"Transforming queues"}),"\n",(0,t.jsxs)(n.p,{children:["A ",(0,t.jsx)(n.code,{children:"Queue[A]"}
1)," is in fact a type alias for ",(0,t.jsx)(n.code,{children:"ZQueue[Any, Any, Nothing, Nothing, A, A]"}),".\nThe signature for the expanded version is:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"trait ZQueue[RA, RB, EA, EB, A, B]\n"})}),"\n",(0,t.jsx)(n.p,{children:"Which is to say:"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsxs)(n.li,{children:["The queue may be offered values of type ",(0,t.jsx)(n.code,{children:"A"}),". The enqueueing operations require an environment of type ",(0,t.jsx)(n.code,{children:"RA"})," and may fail with errors of type ",(0,t.jsx)(n.code,{children:"EA"}),";"]}),"\n",(0,t.jsxs)(n.li,{children:["The queue will yield values of type ",(0,t.jsx)(n.code,{children:"B"}),". The dequeueing operations require an environment of type ",(0,t.jsx)(n.code,{children:"RB"})," and may fail with errors of type ",(0,t.jsx)(n.code,{children:"EB"}),"."]}),"\n"]}),"\n",(0,t.jsxs)(n.p,{children:["Note how the basic ",(0,t.jsx)(n.code,{children:"Queue[A]"})," cannot fail or require any environment for any of its operations."]}),"\n",(0,t.jsx)(n.p,{children:"With separate type parameters for input and output, there are rich composition opportunities for queues:"}),"\n",(0,t.jsx)(n.h3,{id:"zqueuemap",children:"ZQueue#map"}),"\n",(0,t.jsx)(n.p,{children:"The output of the queue may be mapped:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val mapped: UIO[String] = \n  for {\n    queue  <- Queue.bounded[Int](3)\n    mapped = queue.map(_.toString)\n    _      <- mapped.offer(1)\n    s      <- mapped.take\n  } yield s\n"})}),"\n",(0,t.jsx)(n.h3,{id:"zqueuemapm",children:"ZQueue#mapM"}),"\n",(0,t.jsx)(n.p,{children:"We may also use an effectful function to map the output. For example,\nwe could annotate each element with the timestamp at which it was dequeued:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"import java.util.concurrent.TimeUnit\nimport zio.clock._\n\nval currentTimeMillis = currentTime(TimeUnit.MILLISECONDS)\n\nval annotatedOut: UIO[ZQueue[Any, Clock, Nothing, Nothing, String, (Long, String)]] =\n  for {\n    queue <- Queue.bounded[String](3)\n    mapped = queue.mapM { el =>\n      currentTimeMillis.map((_, el))\n    }\n  } yield mapped\n"})}),"\n",(0,t.jsx)(n.h3,{id:"zqueuecontramapm",children:"ZQueue#contramapM"}),"\n",(0,t.jsxs)(n.p,{children:["Similarly to ",(0,t.jsx)(n.code,{children:"mapM"}),", we can also apply an effectful function to\nelements as they are enqueued. This queue will annotate the elements\nwith their enqueue timestamp:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val annotatedIn: UIO[ZQueue[Clock, Any, Nothing, Nothing, String, (Long, String)]] =\n  for {\n    queue <- Queue.bounded[(Long, String)](3)\n    mapped = queue.contramapM { el: String =>\n      currentTimeMillis.map((_, el))\n    }\n  } yield mapped\n"})}),"\n",(0,t.jsx)(n.p,{children:"This queue has the same type as the previous one, but the timestamp is\nattached to the elements when they are enqueued. This is reflected in\nthe type of the environment required by the queue for enqueueing."}),"\n",(0,t.jsxs)(n.p,{children:["To complete this example, we could combine this queue with ",(0,t.jsx)(n.code,{children:"mapM"})," to\ncompute the time that the elements stayed in the queue:"]}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"import zio.duration._\n\nval timeQueued: UIO[ZQueue[Clock, Clock, Nothing, Nothing, String, (Duration, String)]] =\n  for {\n    queue <- Queue.bounded[(Long, String)](3)\n    enqueueTimestamps = queue.contramapM { el: String =>\n      currentTimeMillis.map((_, el))\n    }\n    durations = enqueueTimestamps.mapM { case (enqueueTs, el) =>\n      currentTimeMillis\n        .map(dequeueTs => ((dequeueTs - enqueueTs).millis, el))\n    }\n  } yield durations\n"})}),"\n",(0,t.jsx)(n.h3,{id:"zqueuebothwith",children:"ZQueue#bothWith"}),"\n",(0,t.jsx)(n.p,{children:"We may also compose two queues together into a single queue that\nbroadcasts offers and takes from both of the queues:"}),"\n",(0,t.jsx)(n.pre,{children:(0,t.jsx)(n.code,{className:"language-scala",children:"val fromComposedQueues: UIO[(Int, String)] = \n  for {\n    q1       <- Queue.bounded[Int](3)\n    q2       <- Queue.bounded[Int](3)\n    q2Mapped =  q2.map(_.toString)\n    both     =  q1.bothWith(q2Mapped)((_, _))\n    _        <- both.offer(1)\n    iAndS    <- both.take\n    (i, s)   =  iAndS\n  } yield (i, s)\n"})}),"\n",(0,t.jsx)(n.h2,{id:"additional-resources",children:"Additional Resources"}),"\n",(0,t.jsxs)(n.ul,{children:["\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://www.slideshare.net/jdegoe
1s/zio-queue",children:"ZIO Queue Talk by John De Goes @ ScalaWave 2018"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://www.slideshare.net/wiemzin/psug-zio-queue",children:"ZIO Queue Talk by Wiem Zine El Abidine @ PSUG 2018"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://medium.com/@wiemzin/elevator-control-system-using-zio-c718ae423c58",children:"Elevator Control System using ZIO"})}),"\n",(0,t.jsx)(n.li,{children:(0,t.jsx)(n.a,{href:"https://blog.softwaremill.com/scalaz-8-io-vs-akka-typed-actors-vs-monix-part-1-5672657169e1",children:"Scalaz 8 IO vs Akka (typed) actors vs Monix"})}),"\n"]})]})}function c(e={}){const{wrapper:n}={...(0,s.R)(),...e.components};return n?(0,t.jsx)(n,{...e,children:(0,t.jsx)(d,{...e})}):d(e)}},28453(e,n,u){u.d(n,{R:()=>a,x:()=>l});var i=u(96540);const t={},s=i.createContext(t);function a(e){const n=i.useContext(s);return i.useMemo(function(){return"function"==typeof e?e(n):{...n,...e}},[n,e])}function l(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(t):e.components||t:a(e.components),i.createElement(s.Provider,{value:n},e.children)}}}]);

Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.