apache / apache/arrow

[C++] join with numeric columns has poor performance

Open
#35,334 1 comment 0 reactions 0 assignees View on GitHub
Component: C++ Component: R Priority: Critical Type: bug
Dominant language
C++
Stars
17.1k
Forks
4.3k
Avg merge
3d 18h
Merged PRs (30d)
91

Description

### Describe the bug, including details regarding any error messages, version, and platform.

When joining and the columns are `numeric` the performance scales really badly with size. I was noticing a huge speed difference when joining columns on `numeric` columns vs `int32/int64` columns when by mistake my columns were `numeric`. This is not the case for `dplyr`.

I made the following figure by creating arrow tables of increasing sizes and joined them with another table a third of the size (see bottom for code):

![image](https://user-images.githubusercontent.com/24678081/234317155-76ddf6c5-793e-4bd7-a704-1336a3aec57c.png)

As can be seen when join columns are `numeric` it scales much worse than the other ones. I'm not even sure it makes sense to join on numerics... but I think this could be a bug. Or at least something worth investigating.

Code to reproduce

```R
sizes <- c(1e3, 1e4, 1e5, 1e6)
results <- data.frame()
for (size in sizes) {
size1 <- size
size2 <- floor(size/3)

# int 32 id column
dfInt1 <- data.frame(id=sample(1:size1, replace=F),
value=runif(size1))
dfInt2 <- data.frame(id=sample(1:size2, replace=F),
value2=runif(size2))
arrowInt1 <- arrow::as_arrow_table(dfInt1)
arrowInt2 <- arrow::as_arrow_table(dfInt2)

# int 64 id column
type <- bit64::as.integer64
dfInt64_1 <- data.frame(id=type(sample(1:size1, replace=F)),
value=runif(size1))
dfInt64_2 <- data.frame(id=type(sample(1:size2, replace=F)),
value2=runif(size2))
arrowInt64_1 <- arrow::as_arrow_table(dfInt64_1)
arrowInt64_2 <- arrow::as_arrow_table(dfInt64_2)

# numeric id column
type <- as.numeric
dfNum1 <- data.frame(id=type(sample(1:size1, replace=F)),
value=runif(size1))
dfNum2 <- data.frame(id=type(sample(1:size2, replace=F)),
value2=runif(size2))
arrowNum1 <- arrow::as_arrow_table(dfNum1)
arrowNum2 <- arrow::as_arrow_table(dfNum2)

arrowResultsInt <- arrowInt1 |> dplyr::inner_join(arrowInt2, by='id') |> dplyr::compute()
arrowResultsInt64 <- arrowInt64_1 |> dplyr::inner_join(arrowInt64_2, by='id') |> dplyr::compute()
arrowResultsNum <- arrowNum1 |> dplyr::inner_join(arrowNum2, by='id') |> dplyr::compute()

res <- microbenchmark::microbenchmark(arrowInt32=arrowInt1 |> dplyr::inner_join(arrowInt2, by='id') |> dplyr::compute(),
arrowInt64=arrowInt64_1 |> dplyr::inner_join(arrowInt64_2, by='id') |> dplyr::compute(),
arrowNum=arrowNum1 |> dplyr::inner_join(arrowNum2, by='id') |> dplyr::compute(),
dplyrNum=dfNum1 |> dplyr::inner_join(dfNum2, by='id'),
dplyrInt=dfInt1 |> dplyr::inner_join(dfInt2, by='id'),
check = NULL, times=10
)
res <- summary(res)
res$size <- size

results <- rbind(res, results)

}

ggplot2::ggplot(data=results, ggplot2::aes(x=size, y=mean, group=expr, color=expr)) +
ggplot2::geom_line() + ylab('mean time (ms)')
```

### Component(s)

R

Contributor guide

Open the contributing guide

Research direction

Start with the supplied R benchmark using arrow::as_arrow_table(), dplyr::inner_join(), and compute() to reproduce the scaling difference for numeric and integer columns. Trace the join path from this reproduction and compare its behavior across the three key types. Done means identifying and addressing the numeric-column performance regression, with benchmark results showing improved scaling.

Written by the indexing model from the issue text.

Assessment

Tech stack
cpp, r
Domain
data, performance
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.