kwando

kwando

I have built a Broadway pipeline with a custom producer that reads from a single queue but I want to start broadway with 3 producers where each one reads from a separate queue. Is there any way to accomplish this?

Showing Posts 1 to 5

josevalim

josevalim

Creator of Elixir

You can start three different topologies but use the same callback module for each of them. Something like

defmodule MyBroadway do
  use Broadway

  def start_link(queue) do
    Broadway.start_link(__MODULE__, ...., queue_name: queue)
  end
end

and then in your supervision tree:

children = [
  {MyBroadway, "queue-1"},
  {MyBroadway, "queue-2"},
  {MyBroadway, "queue-3"}
]

If you want them to be consumed precisely from the same topology (let’s say because you need to share batchers), then that’s a feature that yo need to add to your producer.

kwando

kwando OP

Thanks, that is a good idea! In this case I want to limit the number of processors though.

My producer is a bit involved, it fetches “batched items” which will each be processed as a Broadway message and I can’t fetch a new batch of items until all items in a batch are processed. I got this working by letting the producer keep track of which items it has emitted.
Fetching a batch of items is a costly operation in this case (~30s) and I don’t want the pipeline to just idle while this is happening so while I’m processing a batch from queue 1 I’m now fetching next batch from queue 2 concurrently.

I might have just forced a square peg into round hole with my current solution, but I got something that works now at least. :slight_smile:

tmpduarte

tmpduarte

I have something similar but a bit more contrived :grinning:.

I have an iot service (saas) where I would like to create a queue for every iot device that registers, so that each one of the millions of devices has its own queue and have broadway consuming / processing the messages sent. Is this even possible? As in, have broadway pipelines that consume messages from queues created at runtime?

This is just one possible MQ architecture I’m thinking about, the other one would be to have a single queue and every iot device just sends their messages to this common queue. Is this better than the previous solution? Won’t I possibly run into a queue length limit problem?

ijunaidfarooq

ijunaidfarooq

has something changed?

I am trying to start the same way, 3 times but its throwing error as

If using maps as child specifications, make sure the :id keys are unique.
If using a module or {module, arg} as child, use Supervisor.child_spec/2 to change the :id, for example:

    children = [
      Supervisor.child_spec({MyWorker, arg}, id: :my_worker_1),
      Supervisor.child_spec({MyWorker, arg}, id: :my_worker_2)
    ]
rogerweb

rogerweb

Jose’s example is a simplification. You have to give a different id for each supervised process (like described in the error message) and a different name when calling Broadway’s start_link/2. You can do both at the same time by defining a child_spec function in the module you use Broadway.

— All posts loaded —

Where Next? Top

Trending in Questions Top

RSP87
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
nseaSeb
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
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
brecabral
Documentation While reading the Scoped Routes section, I noticed that the documentation currently refers to a problem without explainin...
New
velrest
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
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
FlyingNoodle
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

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
marciok
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
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Dmk
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
netoum
Corex is an accessible, unstyled UI component library for Phoenix that integrates Zag.js state machines using Vanilla JavaScript and Live...
New
webofbits
With AI doing more of the implementation work, I’ve been wondering how much coding I should deliberately keep doing myself. My main conc...
#ai
New

We're in Beta

About us Mission Statement

Options

Thread Display Mode




Thread Preview

Skip Thread Previews