dli

dli

Hey everyone,

I have a GenStage producer that receives real-time events from a 3rd party websocket API and produces GenStage events. I want to process, then broadcast these events via a Phoenix channel and persist them to the DB as well.

For the latter, I wonder whether there is a way to both:

  1. Create chunks with a max size, e.g. ~20,000 events
  2. Create chunks of smaller sizes when a timeout occurs before the chunk is “full”

Both types of chunks should be persisted immediately.

I already have a working GenStage consumer that uses Process.send_after/4 to schedule a flush when the first event of a chunk is processed. However, I hope that a stream-based implementation could be terser.

Is this at all possible with Streams? AFAIK, Stream.interval/1 blocks the caller and is therefore not an option + I’d need to receive these events inside Stream.chunk_while/4.

Showing Posts 1 to 4

LostKobrakai

LostKobrakai

Broadway has batching capabilities by amount and timeout.

dli

dli OP

IIUC, you’re referring to the batch_size and batch_timeout options.

Right now, Broadway would add too much complexity to the project. I’d like to avoid adding services like Redis or Kafka just to get some convenience for batching events.

WDYT, is this feasible with Streams? In the meanwhile, I also found Flow.Window but am not sure how to combine periodic windows and count windows.

LostKobrakai

LostKobrakai

Broadway can consume genstage producers. You seem to already be using genstage, so it shouldn‘t be a large change.

dli

dli OP

Gotcha. Just curious, how would you solve it without Broadway?

— 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
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
matt-savvy
Anyone here using Honeybadger? My Honeybadger account is being overwhelmed with noise from some bots. Seeing a lot of Bandit.HTTPError...
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