jnnks

jnnks

I am building a data processing pipeline which may sporadically fail. I would like to know when Producers/Consumers fail, so I can keep track of running pipelines. When I crash my producer, other processes crash with the same error message, which is confusing.

Consider the following setup:

defmodule CrashingProducer do
  use GenStage
  def init([]), do: {:producer, 0}
  def handle_demand(_, _), do: raise RuntimeError, "producer crashed"
end

# single error when commented, as expected, 
#   two error messages when included?
# Process.flag(:trap_exit, true)

{:ok, prod} =
  GenStage.start_link(CrashingProducer, [])
  |> IO.inspect(label: "prod pid")

spawn_link(fn ->
  Flow.from_stages([prod], max_demand: 1, stages: 1)
  |> Flow.run(link: false)
end)
|> IO.inspect(label: "flow pid")

:timer.sleep(1000)

When running it as is, the producer crashes, as expected with a single error.
When adding the exit trap, the GenStages started by Flow also crash, repeating the same error message. It appears as if the error originated in another process. What is the design decision behind this? It does not make sense to me.
Would it not be easier if the consumers :shutdown upon an error in the consumer? Can I do this manually?

Can I ignore these errors somehow? I already know that my producer failed, so I can shut everything else down manually. It is inconvenient to trap and handle these exits as well.

This is especially annoying when there are multiple Flow Stages and the terminal is spammed with the same error message over and over.

Showing Posts 1 to 4

jnnks

jnnks OP

I am now using Task.Supervisor.async_stream and a pipeline definition using Keyword lists. Something like this:

stages = [
  add_one: fn i -> i + 1 end,
  sub_one: fn i -> i - 1 end
]
input_stream = 1..10 |> Enum.to_list()

{:ok, sup} = Task.Supervisor.start_link(name: :tasksup)

reduce_fn = &Enum.reduce(stages, &1, fn {_name, stage}, item -> stage.(item) end)
Task.Supervisor.async_stream(sup, input_stream, reduce_fn)
|> Stream.take_while(fn
  {:ok, _} -> true
  _ -> false
end)
|> Stream.map(fn {:ok, item} -> item end)
|> Enum.to_list()
|> IO.inspect()
dimitarvp

dimitarvp

But this too will stop on the first return value that’s not {:ok, stuff}, is that your intention?

jnnks

jnnks OP

Sometimes yes, sometimes no :slight_smile:

When failed items can be discarded the In others I would replace

|> Stream.take_while(fn
  {:ok, _} -> true
  _ -> false
end)
|> Stream.map(fn {:ok, item} -> item end)

with a filter dropping everything but {:ok, item}s.

When I need to react to failures I can return the :ok/:exit tuples to the caller as well.

dimitarvp

dimitarvp

My bad, I misread the function name. Apologies for the noise.

— All posts loaded —

Where Next? Top

Trending in Questions Top

Blokh
Hey guys, I’ve got a huge CSV ( around 10 GB ) that needs to be processed hourly Do you guys have any suggestions what is the best prac...
New
kszambelanczyk
Hello! Could someone please give me a help/sample code, how to delete a file from s3 using waffle/waffle_ecto from Phoenix app. I creat...
New
Onor.io
I have what I’ve heard referred to as a “lookup table” in my database. This is a way of assigning codes to common values. One common lo...
New
Trolleger
What approach to take when sending live updates to “random” users Hi! I have a question, I have a little chat app, and when I create a DM...
New
matt-savvy
Anyone here using Honeybadger? My Honeybadger account is being overwhelmed with noise from some bots. Seeing a lot of Bandit.HTTPError...
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
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

Other Trending Topics Top

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
mcass19
ExRatatui lets you cook up rich terminal UIs in Elixir, powered by Rust’s ratatui via Rustler NIFs. Build interactive terminal applicatio...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge & Solve. They are GUI (Emerge) and State management (S...
New
netoum
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New
wintermeyer
There are three potential reasons for members of this forum to have a look at https://vutuv.de You are tired or annoyed of LinkedIn. Yo...
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews