mfclarke
We have ConsumerSupervisor which starts 1 child per event. Is it possible to use this as a :producer_consumer ?
[A] → [B(1, 2, 3, 4, 5, …)] → [C]
Effectively what I’m after is being able to get B to subscribe to A and C to subscribe to B. Then manage concurrency of B just with :max_demand. As opposed to starting an arbitrary number of regular :producer_consumer stages for B and manually handling subscriptions of C to each B.
The example of using a ConsumerSupervisor in the gen_stage repo returns a Supervisor style term in it’s init/1 but the docs show a regular GenStage style return value where you specify :producer, :producer_consumer or :consumer: ConsumerSupervisor — gen_stage v1.3.2
So I’m thinking this might be possible but I’m not understanding how to set it up.
Trending in Questions
Other Trending Topics
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
- #ai
- #elixirconf-us
- #blog-post
- #elixir-ls
- #phoenix_html
- #iex
- #graphql
- #genstage
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #hex
- #security
- #metaprogramming










Showing Posts 1 to 10- Show Best Posts
- Show All (oldest first)
- Show All (newest first)
peerreynders
Based on my GenStage experiments and reading the ConsumerSupervisor documentation here are my thoughts:
ConsumerSupervisor is a Supervisor process which creates a child process for each event that it receives. In general Supervisors are designed to be as simple as possible so that they can focus on one thing: “supervising child processes” (while not taking on any additional responsibility which may cause them to crash) - this one just happens to dynamically create children for any events received. By this nature the ConsumerSupervisor can only really be a Consumer in a GenStage scenario as a supervisor typically doesn’t collect and manage results from its children.
Your particular
[A] -> [B(1, 2, 3, 4, 5, ...)] -> [C]scenario could be realized with C as a Producer (if there is a (producer-) consumer D). Essentially the children of the ConsumerSupervisor would simply deliver their result to C viaGenStage.cast/2orGenStage.call/3- i.e. flow to C would be entirely governed by the:min_demand,:max_demandsettings on B, the ConsumerProducer.The documentation of
initis strange to say the least. There is aConsumerSupervisor.init/1callback that is analogous toSupervisor.init/1- so I would expectConsumerSupervisor.init/1function to be analogous toSupervisor.init/2. The text of theConsumerSupervisor.init/2function seems very similar to the text of the Genstage.init/2 callback - so I suspect a copy/paste error. TheConsumerSupervisor.init/1function code simply prepares a tuple containing the supervisor flags and child specifications.mfclarke
Many thanks for the reply @peerreynders. Your first point makes perfect sense, and matches how the example works and the description of ConsumerSupervisor in the docs. I think there’s a mistake in
init/1in the docs then.The second point I’m not sure I understand though. C can’t possibly be a
:producerbecause it’s supposed to receive events from B (B being a single:producer_consumerstage or an array of:producer_consumerstages). Unless I’m totally misunderstanding you!peerreynders
GenServerat C would suffice.GenStage.handle_cast/2,GenStage.handle_call/3,GenStage.handle_info/2, etc. The difference is that the arrival of these events isn’t regulated.:min_demand/:max_demandconfiguration.mfclarke
Ok so I think what you’re suggesting is to make B a
ConsumerSupervisorthat sends events to C viacast/2/call/3. Where C is just a regularGenServeroutside of the GenStage structure:[A] -> [B] ..cast/call.. [C]In this way C wouldn’t be able to provide back pressure.
That would almost work for me, but (you guessed it) C needs to provide some back pressure, which is part of the reason why I’m using gen_stage in the first place
I’m starting to look into
Flowwhich seems to be able to this kind of thing out of the box, but sacrificing the flexibility of having full fledged GenStages. It seems like if Flow can do it, then it should be possible to implement the same thing with what Flow is built on top of - GenStage. Maybe I need to dig into the Flow source a bit…peerreynders
In a sense ConsumerSupervisor + Producer “is a Consumer Producer” where
:min_demand/max_demandprovides the back pressure regulation for the Producer in front of it and the (B) Producer is regulated by the Consumer behind it. So you could fashion something this way:A(Producer) -> (Bc(ConsumerSupervisor) ~cast/call~> Bp(Producer)) -> C(Consumer)BcConsumer portion ofBBpProducer portion ofBmfclarke
Hmm, nice! Thank you @peerreynders
Trevoke
Hi,
Are you saying that the consumer supervisor Bc would trigger a task that does some work then sends a message to a producer (Bp) to queue up additional work that consumer C would then do more work on?
peerreynders
I’m not quite sure what exactly you are asking. (Bp) was only introduced because of the statement
So I assume that C was going to do some further work on the results that where generated by the tasks spawned by (Bc). Now whether results have to be queued at (Bp) is entirely dependent on how demand is handled within (Bc/Bp). That part of the solution was never discussed - for a regular producer-consumer:
i.e. there is code in the GenStage behaviour that automatically forwards demand from the final consumer of the pipeline to the producer at the beginning of the pipeline - that forwarding mechanism would need to be manually added to (Bc/Bp). Bp could receive demand via
handle_demand/2(responding with an empty list of events), while “informing” Bc of the demand.I suppose Bc could handle subscription manually and use
GenStage.ask/3to get the demand to A once Bc actually “knows” what the demand is. At this point it’s A’s responsibility to provide no more than is demanded (and store demand if it can’t provide it all).If back pressure is propagated in this manner there shouldn’t be a need to queue any results at Bp - it can simply release any result to C immediately via a
{:noreply, [result], state}tuple (see for examplehandle_cast/2) as there is no actual obligation to “batch” the events. Bp would then primarily exist to receive demand from C and to forward it to Bc.Trevoke
Oh, I’m sorry, I think I misunderstood your earlier diagram.
Here’s what I’m looking for, which might be misguided (I’m still very much a beginner to the GenStage world):
ProducerA has eventsConsumerSupervisorB requests the events from A and spawnsTaskworkers up tomax_demand.The work from B’s workers can be picked up by another
ConsumerSupervisor, CI’m not quite sure how to connect B and C, and I thought that your diagram would allow me to do so. Is this what you were talking about? Now I’m getting the impression that you aren’t, and that I’m just not understanding what you’re saying (the ocean of GenStage is vast and I am just getting my sea legs).
peerreynders
I’m by no means a GenStage expert - I have just happened to play around with it for a few days 4 months ago - and looking at the documentation right now I still find it a bit hard to digest.
For example it’s only when I started to write the code that I noticed that any of the callbacks that can return a tuple that includes
[events]can release events to the next stage - skimming through the examples in the documentation you could easily be left with the impression thathandle_demand/2is the place where events are released - when in fact it is only one place where events are released.If you have events during the
handle_demand/2callback then by all means release what you have that doesn’t exceed the demand. But as a producer any unfulfilled demand has to be “stored” - because the consumer is only going to issue another “demand” if it wants more than it already asked for.So whenever a producer “acquires events” it can release them immediately provided it has “stored demand”.
So there are situations where
handle_demand/2will simply return an empty event list (because there are simply none available at the time) - and the events are only released later when the producer somehow gets ahold of them - and then they are released as a result of the callback that “delivers” the events to the producer.Now a consumer is a
GenStagethat is at very end of the pipeline and is ultimately considered to be the “constraining operation” - that is why it gets to set the demand that propagates up all the way to the producer at the beginning of the pipeline (who has to honour that constraint). SoConsumerSupervisoris inherently designed to be at the end of that pipeline.My guess is that it doesn’t need to be a producer because conceptually as it is the “constraining operation” back pressure isn’t an issue with any processing stages that follow - they can simply be implemented as
GenServers because they can deal with any volume that theConsumerSupervisoris capable of throwing at them.Which brings me to the pertinent point - why would you think that you need
ConsumerSupervisorfeeding anotherConsumerSupervisor? It’s a necessary question to ensure that we aren’t running into the XY problem - i.e. there may be a solution to your actual problem that has nothing to do withConsumerSupervisors.