axelson

axelson

Scenic Core Team

Using GenStage to batch events

I am working with an external API (external from the Elixir server at least) and I want to minimize the amount of API calls to a particular endpoint. It feels like GenStage is a good fit for this use case but I’ve had a little trouble getting everything to line up.

I want to:

  • Collect all events every 100ms
  • Batch the events into batches of up to 150
  • Have three workers running in parallel to send these batched requests to the remote API

These seem to fit into three GenStage stages

  • Collector - GenStage :producer that collects events from other parts of the elixir system in a FIFO queue (based on :queue) and sends the events only in handle_demand/2
  • Batcher - GenStage :producer_consumer that asks the Collector for 1000 events every 100ms and batches them into batches of 150 (runs in :manual mode)
  • RequestSupervisor - GenStage ConsumerSupervisor that requests the batched events from the Batcher and starts workers that call the external API

The code I have seems to work but I feel like I may be going against the ethos of GenStage since I’m not really propagating demand all the way up the chain. Specifically the Batcher and the Collector both mostly ignore demand. The Batcher is set to :manual mode and asks for a static 1000 events every 100ms.

Another issue is that this setup will always incur a penalty of 100ms on each event even if there are more than 150 events that are added to the Collector at once.

But also keep in mind that I only expect about 5-10 events every 100ms (and maybe even less). But I want to have a good base for future scaling if necessary.

Any thoughts on this architecture? Is there anything that I haven’t considered that I should consider?

First Post!

axelson

axelson

Scenic Core Team

Or perhaps this setup would be better:

  • Collector - GenStage :producer that collects events from other parts of the elixir system in a FIFO queue (based on :queue).
    • When it receives events and has more than 150 events, then immediately emit the events.
    • Emit any full and partial batches of events in handle_demand/2
  • Delayer - GenStage :producer_consumer that asks the Collector for 1000 events every 100ms
    • Does this still need to run in :manual mode?
  • RequestSupervisor - GenStage ConsumerSupervisor that requests the batched events from the Delayer and starts workers that call the external API

I could probably find a better name than Delayer. Does this setup seem better? I think it might be because previously Batcher and Collector both needed to work together for Batcher to be able to actually batch events. Now batching events is the responsibility of the Collector. Also the demand can handled a little better in the Collector rather than just ignored, but now I need to prototype to make sure that I can get the batching to still work the way I want.

Where Next?

Popular in Questions Top

Brian
What is the proper way to load a module from a file in to IEX? In the python world, doing something like this pretty standard: from ....
New
vonH
In asking this question I am more interested about the expressiveness of the language itself and less concerned about the availability of...
New
LegitStack
I’m hoping you guys can give me some general advice and perhaps code examples if you’re feeling up to it. I’m very interested in Elixir,...
New
pgiesin
This should be a simple problem but I just can’t seem to figure it out. I have a standalone Elixir app that won’t find the database. Dep...
New
Werner
Hi, I’m using Ubuntu 18.04 and after updating to OTP-24.0 yesterday i have this warning when I run “mix local.hex”: 14:57:30.512 [warn] ...
New
Jim
As a follow up to my earlier question: I have the code compiling and running but not getting a successful login from the rest server. ...
New
chensan
I have a User schema with a :from_id field set to type :string: defmodule TweetBot.Repo.Migrations.CreateUsers do use Ecto.Migration ...
New
gonzofish
I’m currently trying to understand how to join three tables using Ecto. All the examples I’ve seen use 2, so maybe I’m just missing somet...
New
chewm
Hi guys, nice to meet you to the whole forum, I’m new here, I’m trying to configure visual studio code for elixir, right now the intellis...
New
romenigld
I am trying to run a deploy with docker and I successfully runned with this command: docker build -t romenigld/blog-prod . but when I t...
New

Other popular topics Top

sorentwo
Hello! tl;dr Announcing Oban, an Ecto based job processing library with a focus on reliability and historical observability. After spen...
977 41022 311
New
srinivasu
How to handle excepions in elixir? Suppose i have A, B, C ,D, E modules. and each module has get() function. A.get() method will call th...
New
gshaw
What is the idiomatic way of matching for not nil in Elixir? E.g., First way: defp halt_if_not_signed_in(conn, signed_in_account) when...
New
sergio_101
I am VERY much an elixir newbie. I have taken one elixir course and one phoenix course on Udemy. During that course, I saw the instructor...
New
lastday4you
I wanted to check elixir version in phoenix because i found that my elixir is 1.5 but when i use Enum.chunk_by it said the function is un...
New
danschultzer
None of the current solutions worked well for me, so I went ahead and built a user management system from scratch. This project took far...
548 27727 240
New
chensan
I have a User schema with a :from_id field set to type :string: defmodule TweetBot.Repo.Migrations.CreateUsers do use Ecto.Migration ...
New
josevalim
Hi everyone, One of the features added to Elixir early on to help integration with Erlang code was the idea of overridable function defi...
New
lucidguppy
I have a super simple question about elixir - how would I take a file like this foo bar baz and output a new file that enumerates th...
New
vrod
I am using the Starship cross-shell prompt – it seems pretty nice, but I get some errors: [WARN] - (starship::utils): Executing command ...
New

We're in Beta

About us Mission Statement