Repository navigation
[FR] gforce mean and sum in parallel #3042
Description
Activity
Related: #2533
@st-pasha thanks, added to first issue.
I posted here last week, but there were some mistakes in my post and it was long and confusing, so I deleted it again. Sorry for that. This is my 2nd attempt:
I have a function which I would like to apply in grouped-by fashion. Applying this function takes much longer as the
data.tableinternals involved in grouping etc. Such a scenario lends itself to be run in parallel but currently this does not work out as one could (naïvely) hope.Let's say, I have data like
tblnew_dt <- function(n = 3e6) { data.table::data.table( a = rep(LETTERS, each = n), b = rep(letters, times = n), c = runif(length(letters) * n, 0, 1), d = runif(length(letters) * n, 1, 2), e = runif(length(letters) * n, 1, 3), f = rnorm(length(letters) * n, 0, 1), g = rnorm(length(letters) * n, 1, 2), h = rnorm(length(letters) * n, 2, 1) ) } tbl <- new_dt()
and two functions
col_var()andplus_one()likecol_var <- function(x) lapply(x, var) plus_one <- function(x) lapply(x, function(y) y + 1)
Now if I split the data into n_core partitions (by adding a group index), run in parallel with
parallel::mcparallel()and within each partition do something likedt[GroupIndex == i, c(use_cols) := fun(.SD), by = group_by, .SDcols = use_cols]
this works fine with
col_var()but not withplus_one(). The reason for this is serialisation of the results for communication from child to master process. Forcol_var(),nrow(res)is26 ^ 2 / n_cores, whereas forplus_one(),nrow(res)isn / n_cores(andn = 3e6).For something like this to work efficiently, we need a way to place a
data.tablein shared memory (see #3104). This leads me to two questions:- Do the maintainers feel that this is in-scope for
data.table? - I'd be interested in having a go at this. Are there people willing to offer some guidance/discussion on implementation strategy (there are several options for shared memory), package API decisions, etc.?
Some of my experiment code is available from here.
- Do the maintainers feel that this is in-scope for
@nbenn
GroupIndex == iis probably going to be inefficient, since in order to compute this filtering condition, you'd have to run through the entire columnGroupIndex, in each thread. That is unless the table is known to be sorted byGroupIndex, in which case a simple binary search would be sufficient.Also, please note that many groupby tasks do not actually require putting data.table into shared memory and forking R process via
parallel::mcparallel(), they can be implemented using the built-in OpenMP library.GroupIndex == iis probably going to be inefficient, since in order to compute this filtering condition, you'd have to run through the entire column GroupIndex, in each thread.Sure, this isn't great. You could make the grouping col an index, or create a vector like
grp_vec <- parallel::splitIndices(nrow(dt), n_cores)
and use this to subset the
data.table. But this is beside the point here:- if the time spent evaluating
funis large, all of this becomes irrelevant - this group indexing could easily be done in C to speed things up
- the subsetting in forked processes is what currently is working. The problem ist getting the results back to the master process (that is why I'm proposing a shared memory implementation).
Also, please note that many groupby tasks do not actually require putting data.table into shared memory and forking R process via parallel::mcparallel(), they can be implemented using the built-in OpenMP library.
What do you mean by many groupby tasks? A handful of functions like
sum(),mean(),var()etc.? This is not what I am talking about here. I am talking about enabling arbitrary functions to be applied in parallel grouped-by fashion.- if the time spent evaluating
@nbenn Great, it helps to clear things up -- that's why we have this discussion in the first place!
When you say "what do you mean by many groupby tasks" -- I really do mean the handful of functions like
sum(),mean(), etc. They are not "many" by the simple count, but they account for the majority of use cases thatdata.tableusers face. And they are probably what @jangorecki had in mind when creating this issue, given that the benchmark he posted uses sums and means.Which is NOT to say that your use case is somehow less important -- clearly it is important to you and to many people like you. So it should be implemented too.
So, I think this discussion helped us understand the following: there are 2 (at least) use cases for "parallel grouping": (1) simple group-by reductions like
sum(),mean(), etc; (2) arbitrary group-by operation with.SD, like computing a regression. These two use cases require radically different approaches: one uses OpenMP, the other process forks. I would say it also means that this issue ought to be split into two: one dedicated to use-case 1, and the other to use-case 2.Yes, @st-pasha I agree that splitting this issue makes sense. The reason why I posted here is that there are already quite a few issues surrounding this topic of "parallel group-by operation" that I did not want to create yet another new issue. But as the two scenarios you are outlining above require different approaches, it does make sense. Thank you for engaging in this discussion.
- changed the title
[-][FR] grouping in parallel[/-][+][FR] gforce mean and sum in parallel[/+]on Jan 11, 2019 The focus of this FR at the top was the 5 grouping tests on https://h2oai.github.io/db-benchmark/.
There is much left to do on extending to other gforce functions and grouping arbitrary functions (existing or new FRs should be refined) but I'd say this FR has been achieved for 1.12.0 which was the target Jan suggested. Have updated the title and closed.
For future reference, results were tweeted here : https://twitter.com/MattDowle/status/1074746218645938176
Reacted by Jan Gorecki
As highlighted in recently updated grouping benchmark https://h2oai.github.io/db-benchmark/ data.table is already lagged behind some other tools, precisely speaking those that can perform aggregation using multiple cores. To keep up with the competition we need to parallelize grouping.
Related issues:
+,sumand many others #2919 aggregate but not group by - parallelism applied on different (lower level) loopWe should try to make it for 1.12.0.