Skip to content

ARROW-16703: [R] Refactor map_batches() so it can stream results - #13650

Merged
paleolimbot merged 7 commits into
apache:masterfrom
paleolimbot:map-batches
Jul 23, 2022
Merged

ARROW-16703: [R] Refactor map_batches() so it can stream results#13650
paleolimbot merged 7 commits into
apache:masterfrom
paleolimbot:map-batches

Conversation

@paleolimbot

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@paleolimbot

Copy link
Copy Markdown
MemberAuthor

@wjones127 I'd love your review on this when you have a chance!

A note that I'm going to keep this as a draft until #13397/ARROW-16444 is merged because after it is we can pass this type of record batch reader directly into the query engine (i.e., the awkward (function(x) x$read_table()) additions in this PR can disappear because we'll be running most queries with the required event loop for SafeCallIntoR() to work.

@wjones127wjones127 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think as_record_batch_reader.function is so cool! 😍 As a follow-up, we should add an example to the datasets vignette. It seems like it might be useful to show how to use it to generate a larger-than-memory simulated dataset.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would collect() not work here?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It doesn't right now because the requisite RunWithCapturedR() isn't there yet (it gets added here: https://github.com/apache/arrow/pull/13397/files#diff-0d1ff6f17f571f6a348848af7de9c05ed588d3339f46dd3bcf2808489f7dca92R132-R144 )

Details
library(arrow, warn.conflicts=FALSE)
#> Some features are not enabled in this build of Arrow. Run `arrow_info()` for more information.source_reader<-RecordBatchReader$create(
batches=list(
as_record_batch(mtcars[1:10, ]),
as_record_batch(mtcars[11:20, ]),
as_record_batch(mtcars[21:nrow(mtcars), ])
)
)
reader<-source_reader|> map_batches(~rbind(as.data.frame(.), as.data.frame(.))) dplyr::collect(reader)
#> Error in `dplyr::collect()` at r/R/dplyr-collect.R:43:48:#> ! NotImplemented: Call to R from a non-R thread without calling RunWithCapturedR#> /Users/deweydunnington/Desktop/rscratch/arrow/cpp/src/arrow/record_batch.h:242 ReadNext(&batch)#> /Users/deweydunnington/Desktop/rscratch/arrow/cpp/src/arrow/util/iterator.h:428 it_.Next()#> /Users/deweydunnington/Desktop/rscratch/arrow/cpp/src/arrow/compute/exec/exec_plan.cc:559 iterator_.Next()#> /Users/deweydunnington/Desktop/rscratch/arrow/cpp/src/arrow/record_batch.cc:337 ReadNext(&batch)#> /Users/deweydunnington/Desktop/rscratch/arrow/cpp/src/arrow/record_batch.cc:351 ToRecordBatches()

Created on 2022-07-19 by the reprex package (v2.0.1)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah right I should have read your earlier comment to completion. Thanks for explaining!

Comment threadr/R/dataset-scan.R Outdated

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should probably add a test to make sure the dots are being passed through.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

That was a really good catch (it needed to be c(list(batch), dots))!

@paleolimbot
paleolimbot marked this pull request as ready for review July 22, 2022 17:15
@paleolimbot

Copy link
Copy Markdown
MemberAuthor

@wjones127 any thoughts on whether or not this is too large of a change to merge? I don't know how many users we have of map_batches()...this does restrict what one can do with the result of map_batches() (for example, an exec plan that ends with head() won't work). One option is to add a lazy = TRUE|FALSE argument to allow the streaming behaviour but default to non-streaming like we did before.

@wjones127

Copy link
Copy Markdown
Member

any thoughts on whether or not this is too large of a change to merge? I don't know how many users we have of map_batches()...this does restrict what one can do with the result of map_batches() (for example, an exec plan that ends with head() won't work). One option is to add a lazy = TRUE|FALSE argument to allow the streaming behaviour but default to non-streaming like we did before.

On one hand we do mark this as experimental, so I'm no super worried about breaking it right now. But also I find "an exec plan that ends with head() won't work" alarming. Why wouldn't that work?

@paleolimbot

Copy link
Copy Markdown
MemberAuthor

It is a bit alarming...it's because some exec plans that end in head() rely on an R-level RecordBatchReader where the ExecPlan is executing stuff in the background. The user-defined functions PR only works when we evaluate the whole plan into a table (so that we can guarantee all the R function calls have happened before we hop out of the event loop). I have ARROW-17178 for the follow-up to fix it, but it sounds like it would be best to add a lazy argument with a default that mirrors the current behaviour so that users can opt-in to the even more experimental behaviour.

@wjones127

Copy link
Copy Markdown
Member

Okay that’s less alarming than I thought. Let’s add that extra param then; it sounds like the best solution for now. Hopefully we can make a goal to stabilize and document in the next version though; we’ve been messing with it for a while.

@paleolimbot
paleolimbot merged commit 70904df into apache:masterJul 23, 2022
@paleolimbot
paleolimbot deleted the map-batches branch July 23, 2022 12:10
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = ee2e944 and contender = 70904df. 70904df is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Failed ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed ⬇️1.57% ⬆️0.07%] test-mac-arm
[Failed ⬇️0.54% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.86% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Failed] 70904dff ec2-t3-xlarge-us-east-2
[Failed] 70904dff test-mac-arm
[Failed] 70904dff ursa-i9-9960x
[Finished] 70904dff ursa-thinkcentre-m75q
[Failed] ee2e9448 ec2-t3-xlarge-us-east-2
[Finished] ee2e9448 test-mac-arm
[Finished] ee2e9448 ursa-i9-9960x
[Finished] ee2e9448 ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

kou pushed a commit that referenced this pull request Feb 20, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@paleolimbot@wjones127@ursabot