lud

lud

Hello,

Is there an idiomatic way to write greedy streams in elixir?

Given this example:

1..1000
|> Stream.map(&computation_1/1)
|> Stream.map(&computation_2/1)
|> ...

If I understand streams correctly, this means that for each item we will have computation_1, then computation_2, and then we go to the next item.

What I would like is for computation_1 to be run as fast as possible without needing to be pulled from downstream. This is useful when computation_1 has variable times and computation_2 has steady times.

I could use Enum:

1..1000
|> Enum.map(&computation_1/1)
|> Stream.map(&computation_2/1)
|> ...

This is very greedy, but it blocks computation_2. It will not start until all computation_1 calls have been made.

So, is there a well-known pattern somewhere for that? In my specific case I’d like not to resort to aync_stream because of memory consumption.

Edit: hmmm while proofreading my post I realize that what I want is essentially concurrent. There is no solution on a single process. So I guess I can just run the top stream in a process, send all results to a second process, and receive in a stream from there.

Showing Posts 1 to 10

D4no0

D4no0

I think Task.async_stream/3 might be something you could use.

NVM, saw you mentioning it as not the optimal solution.

dimitarvp

dimitarvp

A nice compromizing solution – or a good prototype – would be to insert Stream.chunk_every(1000) in the middle and see if it helps.

D4no0

D4no0

If your real-world processing gets more advanced, I think you could also take a look at GenStage, it will allow to make a much more readable and configurable processing pipeline.

lud

lud OP

Not optimal but definitely an improvement in some cases. If I have a slow computation for every 10 computations, then a chunk size around 10 should help.

@D4no0 yes that’s the kind of setup that would involve two processes I was talking about in my Edit. Not sure if pulling that library would always be desirable but if yes then it would make a simple-to-write solution!

dimitarvp

dimitarvp

I agree it’s not optimal, I simply don’t know anything about your problem. It seems like something that would accelerate computation_2 a bit.

If we’re talking truly optimal I’d just have a Kafka / NATS queue and send the results of the first computation to it and computation 2 would be pulling from that queue with the maximum speed and parallelism possible i.e. you would code a solution that’s hard-bottlenecked on the speed with which items from computation_1 can be emitted.

D4no0

D4no0

I think it depends on the use-case, but I would highly recommend to read the documentation as it should contain the list of things it solves compared to self-rolled processing pipelines.

lud

lud OP

I simply don’t know anything about your problem

Oh I know, I meant it was actually a simple and good enough solution. Simplicity is a virtue :slight_smile:

My current use case is generating an ex_unit test from external resources (computation_1) and syncing that file to the disk or an external API.

This is not something where you want external services such as Kafka.

Honestly there is not much to optimize in here, I created that topic for the general case. We could reduce the scope a little bit. I think the real question can be reduced: Given two steps in a stream chain, how can we run both steps in parallel. While computation_2 runs for item i, how to make computation_1 run for item i+1.

dimitarvp

dimitarvp

Then a very normal task pool, a library or homegrown, will do the job just fine?

lud

lud OP

Yeah I’m running homegrown :smiley: One process running the top stream and sending every reesults to the parent process.

The parent process starts it’s stream by looping over receive.

dimitarvp

dimitarvp

You can do even better than that. If you use f.ex. opq or poolex you can distribute the results of computation_1 to a job pool and then it will run jobs on each complete computation – in parallel.

The way you seem to be describing your usage seems like the distribution of tasks for downstream processing is serial (by a single process running receive in a loop). And it can be made parallel.

EDIT: you can probably even use Phoenix.PubSub. :thinking:

Where Next? Top

Trending in Questions Top

katta
I having some trouble figuring out if I have set myself too strict of standards for my production server. Currently I can handle 75% of r...
New
achenet
Hello, I’m trying to build a basic Phoenix web-app, and I’d like to use Tailwind. However, when I launch mix phx.server, I get an error...
New
kpanic
Hi everyone, I am toying with the idea of building a “match maker” for giving personal help to people that wants to start coding. I sta...
New
Cxx-mlr
I’m working on a small exercise involving update_in/3, and I came up with this solution: data = %{ name: "Periodic Table", category:...
New
ChrisAmelia
I’ve got trouble wrapping my head around the order in which functions are called in this snippet (from Phoenix’s authentication): toke...
New
dillonoconnor
Is there any way to avoid the Hologram compiler running when using iex? It seems like the front-end code could potentially be disregarded...
New
thiagogsr
** (ArgumentError) expected :max_attempts to be a positive integer, got: {:@, [line: 10, column: 19], [{:max_attempts, [line: 10, column:...
New

Other Trending Topics Top

GenericJam
Edit: 2026 May 15 - This post is archived. Mob is alive!! Main docs: mob v0.7.11 — Documentation A bit of explanation for the slightly c...
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
garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
New
budgie
A little off-topic, but I feel like people here have a good head on their shoulders. I used to be quite good at making software. Was luc...
New
KristerV
Hey. Is there anyone here who creates agents in their apps? Not talking about using agents, but creating them. I’m finding it pretty diff...
New
mcass19
ExRatatui lets you cook up rich terminal UIs in Elixir, powered by Rust’s ratatui via Rustler NIFs. Build interactive terminal applicatio...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews