Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- library(arrow)
- library(dplyr)
- library(purrr)
- library(future)
- library(future.mirai)
- library(futurize)
- # ============================================================
- # 1. PARTITION THE ORIGINAL PARQUET DATASET
- # ============================================================
- # Number of partitions ("buckets") to create.
- # Every group will belong entirely to exactly one bucket.
- n_bucket <- 64L
- # Open the Parquet dataset lazily.
- # At this point the full dataset is NOT loaded into R memory.
- ds <- open_dataset("file.parq")
- # Assign each group to one of the 64 buckets.
- #
- # This assumes that `grp` is an integer.
- # Groups with the same value of `grp` always get the same bucket,
- # so a group is never split across partitions.
- #
- # write_dataset() operates through Arrow and does not require
- # collecting the whole 62-million-row dataset into R first.
- ds |>
- select(grp, v1, v2) |>
- mutate(
- bucket = as.integer(grp %% n_bucket)
- ) |>
- write_dataset(
- "bucketed",
- format = "parquet",
- partitioning = "bucket"
- )
- # ============================================================
- # 2. DEFINE HOW TO PROCESS ONE BUCKET
- # ============================================================
- process_bucket <- function(b) {
- # Open the partitioned dataset lazily.
- #
- # filter() is executed by Arrow, so only the requested
- # bucket is selected before anything enters R memory.
- x <- open_dataset("bucketed") |>
- filter(bucket == b) |>
- select(grp, v1, v2) |>
- collect()
- # At this point ONLY this bucket is in R memory.
- #
- # Because x is now an ordinary R data frame/tibble,
- # my_boot_fn() can be any ordinary R function.
- #
- # All rows belonging to a group are present in this bucket.
- x |>
- summarise(
- boot = my_boot_fn(v1, v2),
- .by = grp
- )
- }
- # ============================================================
- # 3. PROCESS THE BUCKETS IN PARALLEL
- # using purrr::map() + futurize + mirai
- # ============================================================
- # Tell the Future framework to use Mirai as its backend.
- #
- # Here at most 8 buckets will be processed simultaneously.
- # Each worker therefore collects only its own bucket,
- # rather than every worker receiving the complete dataset.
- plan(
- future.mirai::mirai_multisession,
- workers = 8
- )
- # This LOOKS like an ordinary purrr::map().
- #
- # futurize() transforms the map into a parallel Future operation,
- # and the Futures are executed by the Mirai backend defined above.
- #
- # Each invocation of process_bucket():
- #
- # 1. reads one bucket from Parquet
- # 2. collects only that bucket into R
- # 3. groups it by `grp`
- # 4. runs my_boot_fn() on every complete group
- # 5. returns only the small group-level result
- #
- result <- 0:(n_bucket - 1L) |>
- map(process_bucket) |>
- futurize(seed = TRUE) |>
- list_rbind()
- # Return to ordinary sequential execution when finished.
- plan(sequential)
- # ============================================================
- # OPTIONAL: SAVE THE RESULT
- # ============================================================
- write_parquet(
- result,
- "result.parquet"
- )
Advertisement
Add Comment
Please, Sign In to add comment