CAIOHOBORGHI
Hello guys,
I`m starting to learn Elixir with the purpose to create a distributed application to do some validations in a list that is updated every min.
Here`s my ideia:
Create a DynamicSupervisor with N workers to isolate the validations(one process one validation of some item of the list)
The list contains about 1.000 items, where each item contains 10 properties.
A process need to validate some of the properties of one item. (eg: If property A from item N is bigger than 10, print it)
The problem is that this list is updated every min and I dont know which is the best way to emit the updates
to every worker. My first ideia is to use ETS but I`m a little concerned about every process doing a select every min in the database, doesn´t look like a good idea.
The complete flux is(I think it should be):
→ Supervisor gets the update from an external socket application
→ Filter the list according to each process item
→ Send to process only the properties of its item
But it needs to be fast, and scalable ![]()
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
- #elixirconf-us
- #ai
- #blog-post
- #elixir-ls
- #phoenix_html
- #iex
- #graphql
- #genstage
- #websockets
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #metaprogramming
- #hex
- #security










Showing Posts 1 to 10- Show Best Posts
- Show All (oldest first)
- Show All (newest first)
jkmrto
hey!
Normally it is not good to do any task on the supervisor.
Instead a valid approach would be using GenStage.
So, with this approach you could think about having a producer, the process that will receive the messages, and many consumers to process messages.
In the Producer you just will need to receive the messages and broadcast them. The filtering part can be done almost automatically using the
selectoroption for the subscription of the consumers.Hope it helps!
CAIOHOBORGHI
Hello jkrmrto, thanks for answering so fast!
I`ve dived into some GenStage tutorials and ended up with something close for what I need, but I still have some issues.
This is my current code
And this is my application.exs
The problem is that the selector is not working, the consumer A stills receives messages B and C, and the same thing happens in B and C Consumers.
I know I can place an if-statement in handle_events function of Consumer, but I think it should be better using your suggestion(filtering with selector)
Could you please analyze my code and help me see where is the mistake?
This is my console output after running “mix run --no-halt”
jkmrto
hey,
You should use a BroadcastDispatcher. So, when initializing the producer:
{:producer, list, dispatcher: GenStage.BroadcastDispatcher}lud
What do you want to broadcast ? The fact that the list has been updated (a single message every minute), or each list items (1000 messages / minute) ?
Edit:
Also
selectoris only supported by BroadcastDispatcher, and you function would not match:Selector:
fn %{key: key} -> state[:code] == key endEvents:
["A", "B", "C"]%{key: key}can never match"A","B"or"C"Also, what do you want to do with the validations ? If you need to update the list then you cannot do it in parallel without handling concurrent writes to the list by two validators. If you just want to print the items then it is fine.
CAIOHOBORGHI
I want to broadcast the list(actualy, a property “lastInfos” of the list, where each worker will receive only lastInfos where state[:code] == event[:code])
The case is:
Every minute I receive a new list, and each worker will be looking at one item, something like this:
List: [
%{code: “A”, lastInfos: {}},
%{code: “B”, lastInfos: {}},
%{code: “C”, lastInfos: {}}
]
Where each worker will print the property lastInfos of the item being observed.
My objective is that, every time I update the list in Producer, the Consumers need to receive the update.
A Consumer will never change the list, he needs only to receive the most recent data of the item observed.
lud
I see,
In that case GenStage can do the job with a consumer per code.
But my code is far from good:
That is why I would not use GenStage if I required process isolation for handling each item. Although, as it is a functional immutable language, I guess you are fine with handling lists of items.
Edit: You could also use the PartitionDispatcher, with the same code as above except for those differences:
It makes more sense to me, but all the possible codes must be known beforehand. For example if you write
partitions: ["A", "B"], then a list item with"C"will crash the producer. On the other hand, with the broadcast dispatcher, it is the opposite, as if you declare A,B,C, then a list item that would have D would be ignored (an maybe stay in the producer memory for ever ? I don’t know). So I would rather have the process to crash and use partitions that have unhandled events.mpope
If you don’t want to bring in external deps, you can leverage pg, which is a way to register processes under a specific group. A gen_server can partition the list between the processes in a group. You mentioned it your goal was a distributed application, which works will with pg. It is eventually consistent across all nodes.
CAIOHOBORGHI
So,
Thanks for the response, I think your code will work.
The scenario is this:
The list of [“A”, “B”, “C”] will actually be a Map list where the primary key is the string, something like:
And the list (with codes and infos) is updated every minute, meaning that new codes (non existing at the first time) could be added.
Just to knowledge, I`ll be getting this list from an external source(through a socket connection)
With that in mind, do you still recomend using a BroadcastDispatcher ?
Is there a way to put every process into the supervision of a DynamicSupervisor?(To be respawned when failed, or killed when the code doesnt exists on the newest list)
Thanks for your time!
lud
Well, in that case, what do you actually want to be sent to an isolated process ?
%{name: "John", age: 18, nickname: "Doe"}or the whole%{infos: [...], code: "A"}?Same for what you want to be killed if something is not known, is it a process that would receive the whole list or an item?
You have many options here.
First I would use a
Registryto register your validators with the key (e.g."A") they know. So you can send them the whole list if the key is registered, or reject the list if you have no validator for it (or use a default validator, or no validation).This registered process can be a simple
GenServerthat would validate each item in a loop. Or if you really need isolation, would useTask.async_streamto validate the items.This registered process could be a
GenStageproducerlike we did above, receiving the whole list and returning it as events with the default dispatcher, and you then add aConsumerSupervisorthat would listen to this producer and would spawn a new process for each item in the list. This is basically how theopqpackage works.Or, as you receive only one list at a time, could send the list to a unique producer, that would transform
%{code: "A", items: [%{name: "mary"}, %{name: "john"}]}into events like this:[{"A", %{name: "mary"}, {"A", %{name: "john"}], and then have an uniqueConsumerSupervisorthat would spawn a task (isoalted process) for each event, and the task would call a dispatch function:Just do not over-engineer your solution because a full minute is very large, even for a million of items.
CAIOHOBORGHI
Well, I though of sending it separated for performance improvement…
What I want to send separated is the entire item
The code is used only as a “primary key” of the consumer, cause I need one consumer for each code.
The list will grow with time, and I guess that it can become a problem in the future if I build a Producer that sends the entire list to each consumer.
I dont know if you understand it clearly, but each consumer will do a quick validation in the item designed for it.
The case that I want to kill a consumer, is if the new list(that is updated every min) doesnt contain the code “he’s” watching.
Like, in the first minute, I`ll have a list like this
And, in the next minute, I`ll get a new list like
In this case, I have to kill consumer that is looking for “B” and start a new consumer, that will look for “Not B”.
So, the new consumer “Not B” has to remain alive untils the item “Not B” is present on the subsequent updates of the list.
The producer will be connected to a socket, that sends new updates of the list each minute, but this time can change in a not-so-close future(for a smaller time interval)