bgoosman
I am implementing a file processing pipeline in Elixir. Generically speaking there are a few steps file → A → B → C, and also A → D.
I’d like to emit data from each step, as it is generated, because each step can vary in throughput, and I want to consume the output immediately.. For example, step A can emit a struct with file metadata and file chunks. For example, D can emit some vectors created by sending file chunks through Bumblebee.
What’s the Elixir way for each step to have a kind of plugin architecture, where I “the pipeline orchestrator” can install a different Consumer for each step? For example, I could install a “StdioWriter” to step A, so I could see what was being emitted.
It’s kind of like a “tap” into the pipe.
Trending in Questions
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
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
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
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
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
If a change or preparation module uses Ash.Changeset.get_argument/2 or Ash.Query.get_argument/2 (or any of the other get_argument functio...
New
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
Other Trending Topics
I am happy to introduce the very α version of the new programming language compiled to BEAM.
Welcome Cure.
It has literally three kille...
New
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
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
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
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
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #library
- #deployment
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #elixirconf
- #channels
- #exunit
- #discussion
- #code-sync
- #podcasts
- #javascript
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #full-time-contract
- #podcasts-by-brainlid
- #ecto-query
- #elixirconf-us
- #blog-post
- #elixir-ls
- #ai
- #phoenix_html
- #iex
- #graphql
- #genstage
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #hex
- #security
- #metaprogramming











Showing Posts 1 to 7- Show Best Posts
- Show All (oldest first)
- Show All (newest first)
cmo
Pass a module or function as an argument and execute it in the step, publish pubsub messages or telemetry events?
gregvaughn
It sounds like you might be looking for Broadway
bgoosman
Passing a function (or list of functions) sounds like the simplest option. I guess the functions should all have the same interface.
How would you implement pubsub / telemetry?
bgoosman
cmo, did you mean this? telemetry — telemetry v1.4.2
It looks like I can use :telemetry to emit events to attached event handlers.
bgoosman
Thanks, I do need something with backpressure, so I thought GenStage or Broadway might help.
bgoosman
Hey all, I was able to accomplish what I wanted to do with the help of two github examples and a forum thread, but now I have a question about trapping exits that I could use some insight on. First here are the references, and after is my question.
The
process_filesfunction calls Supervisor.start_link() to call some GenStage processes. The first step,A, is marked:transientand terminates with reason :normal when there are no more files. All of the other steps are also:transientso they terminate whenAterminates.Ais also marked:significant, so the Supervisor will auto shutdown (because ofauto_shutdown: all_significant).process_filesis called from an ExUnit test. If I comment outProcess.flag(:trap_exit, true), why does my ExUnit test fail with exit reasonshutdown? If I uncomment it, my ExUnit test passes. I’ve renamed some things for simplicity’s sake.Am I using Supervisor correctly? Is there an entirely different way to do this?
process_files
ExUnit test
bgoosman
I found Oban and rewrote my solution quite quickly. Oban Web is great too.