leighshepperson

leighshepperson

This is probably a basic question about flows, but:

At the partition stage, where you can specify the number of partitions, does this number only relate to physical processors and nodes?

Or, if the number of partitions exceeds this value, does it also partition against otp processes?

At the reducer stage, you should be able to run the computations in parallel as much as possible right? So it would be an advantage to create as many partitions as possible (supposing you are not copying too much data around)?

Showing Posts 1 to 4

imetallica

imetallica

You can set as an argument the number of the partitions you want. By default, Elixir will use the same number of processors on the machine and it’s not designed to work in a distributed fashion (Flow — Flow v1.2.4).

Flow.from_enumerable([1, 2, 3], stages: 3) # To use 3 processes

If I’m not mistaken, the reducer stage is where you join everything together. So, the computations are run in parallel on the mapping stage, not on the reducer stage. So, every time you call Flow.partition/2, new processes will be used by that computations. For example:

[1, 2, 3]
|> Flow.from_enumerable()
|> Flow.map(fn x -> x + 1 end) # [2, 3, 4]
|> Flow.partition() # A new group of processes will be spawned here.
|> Flow.map(fn x -> x * 2 end) # [4, 6, 8]
|> Flow.partition() # A new group of processes will be spawned here.
|> Flow.reduce(fn -> 0 end, fn x, acc -> x + acc end)

$> 18

leighshepperson

leighshepperson OP

Hi thanks for your reply!

Perhaps I’m thinking about Flow the wrong way:

I imagined it was meant to work like map reduce. So the partition stage would be like the shuffle stage, i.e. there has already been some kind of mapping performed - preferably in parallel , giving us pairs:

{key, values}

Then, these are sent to different nodes (bucketed by key) and the values are reduced in parallel on each of the nodes. Once this has been done, the result can be obtained.

Looking at the way this is done in flow, we have: mapping stage, i.e. operations on a collection, then a partition stage, that partitions by some hash key that determines what node the values should go to, and then they are finally reduced. If the operation you want to perform on the values is associative, for example, then there is no reason why it can’t be done in parallel. This is why I was thinking it might leverage it between OTP processes in addition to the number of processor/node partitions.

Would it necessarily be a bad idea to have a version of Flow that could do this? I.e., if you know the reducer operations are associative and if the map stage splits up the input data into tmp files indexed by key that can only read by the process associated to each key , for example, then you could also partition by OTP processes?

imetallica

imetallica

Be in mind that Flow does not work in a distributed way - it only works on a single node.

I’m not sure if I understand your question. Can you elaborate more what you want to achieve?

quda

quda

Sorry, I have to revive this old topic because I have new… needs about Flow.
I know that at the level of 2017 Flow didn’t work in a distributed way.
Now is it possible to make it run on multiple nodes ? In such way to distribute partitions among nodes ?
Is there an alternative for distributed MapReduce in Elixir ?

— All posts loaded —

Where Next? Top

Trending in Questions Top

stjefim
Hello! Suppose you are building workflow (order / task / payment) processing system with the following requirements: Each workflow con...
New
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
roeland
Kia ora, We have been using elixir-google-api to connect to Google Drive. However, with the updates to Tesla due to CVEs this is now bro...
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
jaybe78
Hello, I’m developing a online persistent chat system (what’s app) like using elixir/dynamodb/aws for a mobile app(flutter). The diffic...
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

Other Trending Topics Top

garrison
Hobbes is a low-level distributed database for the Elixir programming language. Hobbes provides a simple, safe, and scalable storage lay...
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
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