ValtteriL

ValtteriL

How to make Task.async_stream return results as they are ready?

I have the following code:

defmodule Demo do
  def run do
    child =
      spawn(fn ->
        Stream.resource(
          fn -> :ok end,
          fn _ ->
            receive do
              {:msg, msg} ->
                IO.puts("# received #{inspect(msg)}")
                {[msg], nil}

              _ ->
                {[], nil}
            end
          end,
          fn _ -> :ok end
        )
        |> Stream.each(fn x ->
          IO.puts("## before #{inspect(x)}")
          x
        end)
        |> Task.async_stream(
          fn x ->
            IO.puts("### during #{inspect(x)}")
            x
          end,
          ordered: false,
          max_concurrency: 2
        )
        |> Stream.each(fn x ->
          IO.puts("#### after #{inspect(x)}")
          x
        end)
        |> Stream.run()
      end)

    1..3
    |> Enum.each(fn x -> send(child, {:msg, x}) end)
  end
end

When run, I get the following output:

iex(6)> Demo.run
# received 1
## before 1
# received 2
### during 1
## before 2
#### after {:ok, 1}
### during 2
# received 3
## before 3
### during 3
#### after {:ok, 2}

The Stream.each after the Task.async_stream is never called with the last message. Async_stream seems to withhold output until it has max_concurrency results ready.

Is there any way to make to avoid the buffering?

Background: I’m building a data processing pipeline that processes endless stream of data and cleanup after each task has to happen immediately after its done. I could do it in the task itself, but that would break encapsulation and force me into complicating things with process messaging. I previously used GenStage, which worked well at small scale, but could use performance improvement as no backpressure is needed.

Showing Posts 1 to 2

al2o3cr

al2o3cr

Try passing ordered: false to async_stream; the docs specifically mention “This is also useful when you’re using the tasks only for the side effects”.

dimitarvp

dimitarvp

Replace the Stream.each + Task.async_stream with Enum.each + Task.async, collecting the handles to all tasks and then just Task.await_many them (note that this succeeds and does not block even if the tasks are already completed).

That, or @al2o3cr’s approach.

I don’t get that. Sounds like you are complicating things by not doing it in the task itself from where I am standing.

— All posts loaded —

Where Next? Top

Trending in Questions Top

RSP87
I’m working on a project that simulates the bumbl example in the programming phoenix book. It acts almost like an email client. We have a...
New
nseaSeb
Hello, I know there is an approach for handling lists that allows for optimized traversal, but I can’t recall the specific method (somet...
New
brecabral
Documentation While reading the Scoped Routes section, I noticed that the documentation currently refers to a problem without explainin...
New
RemyXRenard
I’m seeing that a list inside a Kino.DataTable will be interpreted as a charlist, even if the Kino.configure() is set to charlists: :as_l...
New
velrest
So my question is quite simple and i have found no conclusive answer on forum, google or AI. Should we use :erlang.float for Integer to ...
New
samoloth
Hi, I’ve just set up an application with ash_authentication. There is only magic link strategy for now, so there is no confirmation add o...
New
asweet-confluent
I recently noticed that Elixir’s Logger defaults its primary log level to :debug when no :logger, :level application configuration is pre...
New

Other Trending Topics Top

JesseHerrick
Hey, I’m Jesse and I’m the main contributor behind Dexter, a full-featured, lightning-fast Elixir LSP optimized for large codebases. It s...
New
mudasobwa
I am happy to introduce the very α version of the new programming language compiled to BEAM. Welcome Cure. It has literally three kille...
New
marciok
Hi there! We created Gust: A task orchestrator inspired by Airflow. For those who have never heard about Aiflow, it’s a Python-based wor...
New
mhanberg
Hi everyone! The first release candidate for the Expert language server project is now available! We’ve published a press release detai...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Dmk
Xamal is a deployment tool for Elixir apps that deploys native releases to bare metal servers over SSH. It’s a port of GitHub - basecamp/...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews