lud

lud

Hello,

I have this example where I would like a stream transformation to return the whole input list up to the third atom, included.

input = ["a", "b", "c", "d", :x, :y, "e", "f", :z, "g", "h", "i", :aa, :bb, "j"]

Expected output:

["a", "b", "c", "d", :x, :y, "e", "f", :z]

All I can do for now is to return:

["a", "b", "c", "d", :x, :y, "e", "f"]

Of course it would be very easy without the constraint: Further elements from the stream should not be processed once the final atom is found in the stream.

Do you know how I should do that?

Thank you!

Some test code:

defmodule Demo do
  def run do
    input = ["a", "b", "c", "d", :x, :y, "e", "f", :z, "g", "h", "i", :aa, :bb, "j"]

    wanted_atoms = 3

    input

    #
    # Classifier function. In real application it costs money to run so we never
    # want to call it if we are not going to process its result.
    |> Stream.map(fn
      "g" -> raise ~s(should not process "g")
      value when is_atom(value) -> {:atom, value}
      value -> {:other, value}
    end)

    # Take 3 atoms if we find them in the stream, take the whole stream
    # otherwise.

    |> Stream.transform(
      _start_fun = fn -> 0 end,
      _redux_fun = fn
        {:atom, value}, atom_count when atom_count == wanted_atoms - 1 -> {:halt, {:last, value}}
        {:atom, value}, atom_count -> {[value], atom_count + 1}
        {:other, value}, atom_count -> {[value], atom_count}
      end,
      __last_fun = fn
        # Not called when the reducer returns :halt
        {:last, value} -> {[value], nil}
        n when is_integer(n) -> {:halt, nil}
      end,
      _after_fun = fn acc -> :ok end
    )
    |> Enum.to_list()
    |> IO.inspect(limit: :infinity, label: "final")
  end

  defp get_score(2), do: {:ok, 1234}
  defp get_score(_), do: {:error, :nope}
end

Demo.run()

Showing Posts 1 to 5

D4no0

D4no0

Have you tried using something like Stream.take_while/2?

LostKobrakai

LostKobrakai

I guess you’d need to write your own code integrating with the Enumerable protocol to do that. Afaik none of the APIs in Stream would allow you to emit a value to the resulting enumerable and at the same time stop consuming more items from the source enumerable.

al2o3cr

al2o3cr

You can do it with a sentinel value, since each step of Stream.transform can emit multiple results:

    |> Stream.transform(
      fn -> 0 end,
      fn
        {:atom, a}, atom_count when atom_count+1 < wanted_atoms ->
          {[a], atom_count+1}
        {:atom, a}, _atom_count ->
          {[a, :__no_more_atoms], :ok}
        {_, v}, atom_count ->
          {[v], atom_count}
      end,
      fn _ -> :ok end
    )
    |> Stream.take_while(& &1 != :__no_more_atoms)

The transform produces a stream with an extra value on the end when the last wanted atom is seen, then take_while snips it off and terminates the stream.

lud

lud OP

Thank you guys!

Alright, no time to dig into the inner workings of Stream.transform today. Your solution is good enough, which is the perfect flavor of good.

input = ["a", "b", "c", "d", :x, :y, "e", "f", :z, "g", "h", "i", :aa, :bb, "j"]

wanted_atoms = 3

stop_ref = make_ref()

input
|> Stream.map(fn
  "g" -> raise ~s(should not process "g")
  value when is_atom(value) -> {:atom, value}
  value -> {:other, value}
end)
|> Stream.transform(0, fn
  {:atom, value}, atom_count when atom_count == wanted_atoms - 1 -> {[value, stop_ref], nil}
  {:atom, value}, atom_count -> {[value], atom_count + 1}
  {:other, value}, atom_count -> {[value], atom_count}
end)
|> Stream.take_while(&(&1 != stop_ref))
|> Enum.to_list()
|> IO.inspect(limit: :infinity, label: "final")
benwilson512

benwilson512

Author of Craft GraphQL APIs in Elixir with Absinthe

I have also used Stream.concat for this eg. Stream.concat(other_stream, [:ending_value])

— All posts loaded —

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
bradley
I really like the adapter patterns that ecto, nebulex, waffle, etc. use and would love find something similar for a key management servic...
New
unaware8150
Hello folks! So at work, we are seeing some situations where we have to define some “fixed” strings that are used across the codebase in...
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

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
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
mudasobwa
I fully migrated to my own harness from Anthropic/Gemini and I think it’s time to share it. Welcome DSH, the DeepSeek Harness, fully writ...
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