larry77

Untitled

Aug 30th, 2026
97
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
text 3.13 KB | None | 0 0
  1. library(arrow)
  2. library(dplyr)
  3. library(purrr)
  4. library(future)
  5. library(future.mirai)
  6. library(futurize)
  7.  
  8.  
  9. # ============================================================
  10. # 1. PARTITION THE ORIGINAL PARQUET DATASET
  11. # ============================================================
  12.  
  13. # Number of partitions ("buckets") to create.
  14. # Every group will belong entirely to exactly one bucket.
  15. n_bucket <- 64L
  16.  
  17. # Open the Parquet dataset lazily.
  18. # At this point the full dataset is NOT loaded into R memory.
  19. ds <- open_dataset("file.parq")
  20.  
  21. # Assign each group to one of the 64 buckets.
  22. #
  23. # This assumes that `grp` is an integer.
  24. # Groups with the same value of `grp` always get the same bucket,
  25. # so a group is never split across partitions.
  26. #
  27. # write_dataset() operates through Arrow and does not require
  28. # collecting the whole 62-million-row dataset into R first.
  29. ds |>
  30. select(grp, v1, v2) |>
  31. mutate(
  32. bucket = as.integer(grp %% n_bucket)
  33. ) |>
  34. write_dataset(
  35. "bucketed",
  36. format = "parquet",
  37. partitioning = "bucket"
  38. )
  39.  
  40.  
  41. # ============================================================
  42. # 2. DEFINE HOW TO PROCESS ONE BUCKET
  43. # ============================================================
  44.  
  45. process_bucket <- function(b) {
  46.  
  47. # Open the partitioned dataset lazily.
  48. #
  49. # filter() is executed by Arrow, so only the requested
  50. # bucket is selected before anything enters R memory.
  51. x <- open_dataset("bucketed") |>
  52. filter(bucket == b) |>
  53. select(grp, v1, v2) |>
  54. collect()
  55.  
  56. # At this point ONLY this bucket is in R memory.
  57. #
  58. # Because x is now an ordinary R data frame/tibble,
  59. # my_boot_fn() can be any ordinary R function.
  60. #
  61. # All rows belonging to a group are present in this bucket.
  62. x |>
  63. summarise(
  64. boot = my_boot_fn(v1, v2),
  65. .by = grp
  66. )
  67. }
  68.  
  69.  
  70. # ============================================================
  71. # 3. PROCESS THE BUCKETS IN PARALLEL
  72. # using purrr::map() + futurize + mirai
  73. # ============================================================
  74.  
  75. # Tell the Future framework to use Mirai as its backend.
  76. #
  77. # Here at most 8 buckets will be processed simultaneously.
  78. # Each worker therefore collects only its own bucket,
  79. # rather than every worker receiving the complete dataset.
  80. plan(
  81. future.mirai::mirai_multisession,
  82. workers = 8
  83. )
  84.  
  85. # This LOOKS like an ordinary purrr::map().
  86. #
  87. # futurize() transforms the map into a parallel Future operation,
  88. # and the Futures are executed by the Mirai backend defined above.
  89. #
  90. # Each invocation of process_bucket():
  91. #
  92. # 1. reads one bucket from Parquet
  93. # 2. collects only that bucket into R
  94. # 3. groups it by `grp`
  95. # 4. runs my_boot_fn() on every complete group
  96. # 5. returns only the small group-level result
  97. #
  98. result <- 0:(n_bucket - 1L) |>
  99. map(process_bucket) |>
  100. futurize(seed = TRUE) |>
  101. list_rbind()
  102.  
  103.  
  104. # Return to ordinary sequential execution when finished.
  105. plan(sequential)
  106.  
  107.  
  108. # ============================================================
  109. # OPTIONAL: SAVE THE RESULT
  110. # ============================================================
  111.  
  112. write_parquet(
  113. result,
  114. "result.parquet"
  115. )
Advertisement
Add Comment
Please, Sign In to add comment