Help implementing 'parallel' writing to a RealizationSink with BiocParallel
Chưa có ai nhận issue này.
Đánh giá
- Độ khó
- 5/5
- Thời gian dự kiến
- Hơn một tuần
- Mức phù hợp với người mới
- 25/100
- Loại issue
- Tính năng
- Độ rõ ràng
- Cần làm rõ
- Mức độ hoạt động
- Đình trệ
- Công nghệ
- r
- Lĩnh vực
- backend-api-design, data, distributed-systems
Hướng nghiên cứu
Bắt đầu với DelayedArray:::RealizationSink, write_block_to_sink(), RegularArrayGrid() và các lệnh gọi BiocParallel::bplapply() trong ví dụ đồ chơi. So sánh kết quả giữa HDF5Array, arrayRealizationSink trong bộ nhớ và RleArray với SerialParam(), MulticoreParam() và SnowParam(); được xem là hoàn thành khi việc ghi các block an toàn và DelayedArray được trả về chứa tất cả các giá trị đã tạo.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Mô tả
I'm trying to implement writing to an arbitrary RealizationSink backend via BiocParallel::bplapply() with an arbitrary BiocParallelParam backend. That is, I want to be able to construct blocks of output, possibly in parallel, and then safely write each block to a HDF5Array, RleArray, etc. (controllable by the user) after each block is generated.
I've managed to get this set up for an HDF5RealizationSink (my main use case) by using the inter-process locks available in BiocParallel to ensure the writes are safe (@mtmorgan am I using these correctly?). But I've not had any luck using the same code for an arrayRealizationSink. Nor can I get it to work for the RleArray backend, although here the problem seems to be slightly different.
Below is a toy example that tries to construct a DelayedMatrix column-by-column using this strategy, in conjunction with different DelayedArray backends:
suppressPackageStartupMessages(library(DelayedArray))
suppressPackageStartupMessages(library(BiocParallel))
FUN_with_IPC_lock <- function(b, nrow, sink, grid, ipcid) {
# Ensure DelayedArray package is loaded on worker (needed for backends
# other than MulticoreParam).
suppressPackageStartupMessages(library("DelayedArray"))
if (is(sink, "HDF5RealizationSink")) {
# Ensure HDF5Array package is loaded on worker, if required.
# NOTE: Feels a little clunky that this is necessary; would be good if
# this was automatically loaded, properly.
suppressPackageStartupMessages(library("HDF5Array"))
}
message("Processing block = ", b, " on PID = ", Sys.getpid())
tmp <- matrix(seq_len(nrow) + (b - 1L) * nrow, ncol = 1)
ipclock(ipcid)
write_block_to_sink(tmp, sink, grid[[b]])
ipcunlock(ipcid)
}
g <- function(FUN, BPPARAM, nrow = 100L, ncol = 10L) {
# Construct RealizationSink
backend <- getRealizationBackend()
message("The backend is ", if (is.null(backend)) "NULL" else backend)
sink <- DelayedArray:::RealizationSink(dim = c(nrow, ncol), type = "integer")
# Consruct ArrayGrid over columns of RealizationSink
grid <- RegularArrayGrid(dim(sink), c(nrow(sink), 1L))
# Fill the RealizationSink
n_block <- length(grid)
# Construct an IPC mutex ID
ipcid <- ipcid()
# Apply function with BiocParallel
bplapply(seq_len(n_block), FUN, BPPARAM = BPPARAM, nrow = nrow, sink = sink,
grid = grid, ipcid = ipcid)
# Return the RealizationSink as a DelayedArray
as(sink, "DelayedArray")
}
# Everything works with the HDF5Array backend
setRealizationBackend("HDF5Array")
#> Loading required package: rhdf5
# Works
g(FUN_with_IPC_lock, SerialParam())
#> The backend is HDF5Array
#> Processing block = 1 on PID = 5740
#> Processing block = 2 on PID = 5740
#> Processing block = 3 on PID = 5740
#> Processing block = 4 on PID = 5740
#> Processing block = 5 on PID = 5740
#> Processing block = 6 on PID = 5740
#> Processing block = 7 on PID = 5740
#> Processing block = 8 on PID = 5740
#> Processing block = 9 on PID = 5740
#> Processing block = 10 on PID = 5740
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] 1 101 201 301 . 601 701 801 901
#> [2,] 2 102 202 302 . 602 702 802 902
#> [3,] 3 103 203 303 . 603 703 803 903
#> [4,] 4 104 204 304 . 604 704 804 904
#> [5,] 5 105 205 305 . 605 705 805 905
#> ... . . . . . . . . .
#> [96,] 96 196 296 396 . 696 796 896 996
#> [97,] 97 197 297 397 . 697 797 897 997
#> [98,] 98 198 298 398 . 698 798 898 998
#> [99,] 99 199 299 399 . 699 799 899 999
#> [100,] 100 200 300 400 . 700 800 900 1000
# Works: the IPC lock does its job!
g(FUN_with_IPC_lock, MulticoreParam(2L))
#> The backend is HDF5Array
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] 1 101 201 301 . 601 701 801 901
#> [2,] 2 102 202 302 . 602 702 802 902
#> [3,] 3 103 203 303 . 603 703 803 903
#> [4,] 4 104 204 304 . 604 704 804 904
#> [5,] 5 105 205 305 . 605 705 805 905
#> ... . . . . . . . . .
#> [96,] 96 196 296 396 . 696 796 896 996
#> [97,] 97 197 297 397 . 697 797 897 997
#> [98,] 98 198 298 398 . 698 798 898 998
#> [99,] 99 199 299 399 . 699 799 899 999
#> [100,] 100 200 300 400 . 700 800 900 1000
# Works: the IPC lock does its job!
g(FUN_with_IPC_lock, SnowParam(2L))
#> The backend is HDF5Array
#> Loading required package: HDF5Array
#> Loading required package: rhdf5
#> Processing block = 1 on PID = 5755
#> Processing block = 2 on PID = 5755
#> Processing block = 3 on PID = 5755
#> Processing block = 4 on PID = 5755
#> Processing block = 5 on PID = 5755
#> Loading required package: HDF5Array
#> Loading required package: rhdf5
#> Processing block = 6 on PID = 5761
#> Processing block = 7 on PID = 5761
#> Processing block = 8 on PID = 5761
#> Processing block = 9 on PID = 5761
#> Processing block = 10 on PID = 5761
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] 1 101 201 301 . 601 701 801 901
#> [2,] 2 102 202 302 . 602 702 802 902
#> [3,] 3 103 203 303 . 603 703 803 903
#> [4,] 4 104 204 304 . 604 704 804 904
#> [5,] 5 105 205 305 . 605 705 805 905
#> ... . . . . . . . . .
#> [96,] 96 196 296 396 . 696 796 896 996
#> [97,] 97 197 297 397 . 697 797 897 997
#> [98,] 98 198 298 398 . 698 798 898 998
#> [99,] 99 199 299 399 . 699 799 899 999
#> [100,] 100 200 300 400 . 700 800 900 1000
# Only SerialParam() works for the in-memory backend. The in-memory backend,
# an arrayRealizationSink, is implemented as an array inside an environment.
setRealizationBackend(NULL)
# Works
g(FUN_with_IPC_lock, SerialParam())
#> The backend is NULL
#> Processing block = 1 on PID = 5740
#> Processing block = 2 on PID = 5740
#> Processing block = 3 on PID = 5740
#> Processing block = 4 on PID = 5740
#> Processing block = 5 on PID = 5740
#> Processing block = 6 on PID = 5740
#> Processing block = 7 on PID = 5740
#> Processing block = 8 on PID = 5740
#> Processing block = 9 on PID = 5740
#> Processing block = 10 on PID = 5740
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] 1 101 201 301 . 601 701 801 901
#> [2,] 2 102 202 302 . 602 702 802 902
#> [3,] 3 103 203 303 . 603 703 803 903
#> [4,] 4 104 204 304 . 604 704 804 904
#> [5,] 5 105 205 305 . 605 705 805 905
#> ... . . . . . . . . .
#> [96,] 96 196 296 396 . 696 796 896 996
#> [97,] 97 197 297 397 . 697 797 897 997
#> [98,] 98 198 298 398 . 698 798 898 998
#> [99,] 99 199 299 399 . 699 799 899 999
#> [100,] 100 200 300 400 . 700 800 900 1000
# Doesn't work: sink isn't filled. The IPC mutex isn't doing its job?
g(FUN_with_IPC_lock, MulticoreParam(2L))
#> The backend is NULL
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] NA NA NA NA . NA NA NA NA
#> [2,] NA NA NA NA . NA NA NA NA
#> [3,] NA NA NA NA . NA NA NA NA
#> [4,] NA NA NA NA . NA NA NA NA
#> [5,] NA NA NA NA . NA NA NA NA
#> ... . . . . . . . . .
#> [96,] NA NA NA NA . NA NA NA NA
#> [97,] NA NA NA NA . NA NA NA NA
#> [98,] NA NA NA NA . NA NA NA NA
#> [99,] NA NA NA NA . NA NA NA NA
#> [100,] NA NA NA NA . NA NA NA NA
# Doesn't work: sink isn't filled. The IPC mutex isn't doing its job?
g(FUN_with_IPC_lock, SnowParam(2L))
#> The backend is NULL
#> Processing block = 6 on PID = 5777
#> Processing block = 7 on PID = 5777
#> Processing block = 8 on PID = 5777
#> Processing block = 9 on PID = 5777
#> Processing block = 10 on PID = 5777
#> Processing block = 1 on PID = 5771
#> Processing block = 2 on PID = 5771
#> Processing block = 3 on PID = 5771
#> Processing block = 4 on PID = 5771
#> Processing block = 5 on PID = 5771
#> <100 x 10> DelayedMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] NA NA NA NA . NA NA NA NA
#> [2,] NA NA NA NA . NA NA NA NA
#> [3,] NA NA NA NA . NA NA NA NA
#> [4,] NA NA NA NA . NA NA NA NA
#> [5,] NA NA NA NA . NA NA NA NA
#> ... . . . . . . . . .
#> [96,] NA NA NA NA . NA NA NA NA
#> [97,] NA NA NA NA . NA NA NA NA
#> [98,] NA NA NA NA . NA NA NA NA
#> [99,] NA NA NA NA . NA NA NA NA
#> [100,] NA NA NA NA . NA NA NA NA
# Only SerialParam() works for the RleArray backend.
setRealizationBackend("RleArray")
# Works
g(FUN_with_IPC_lock, SerialParam())
#> The backend is RleArray
#> Processing block = 1 on PID = 5740
#> Processing block = 2 on PID = 5740
#> Processing block = 3 on PID = 5740
#> Processing block = 4 on PID = 5740
#> Processing block = 5 on PID = 5740
#> Processing block = 6 on PID = 5740
#> Processing block = 7 on PID = 5740
#> Processing block = 8 on PID = 5740
#> Processing block = 9 on PID = 5740
#> Processing block = 10 on PID = 5740
#> <100 x 10> RleMatrix object of type "integer":
#> [,1] [,2] [,3] [,4] ... [,7] [,8] [,9] [,10]
#> [1,] 1 101 201 301 . 601 701 801 901
#> [2,] 2 102 202 302 . 602 702 802 902
#> [3,] 3 103 203 303 . 603 703 803 903
#> [4,] 4 104 204 304 . 604 704 804 904
#> [5,] 5 105 205 305 . 605 705 805 905
#> ... . . . . . . . . .
#> [96,] 96 196 296 396 . 696 796 896 996
#> [97,] 97 197 297 397 . 697 797 897 997
#> [98,] 98 198 298 398 . 698 798 898 998
#> [99,] 99 199 299 399 . 699 799 899 999
#> [100,] 100 200 300 400 . 700 800 900 1000
# Doesn't work: Don't understand why.
g(FUN_with_IPC_lock, MulticoreParam(2L))
#> The backend is RleArray
#> Error in validObject(.Object): invalid class "ChunkedRleArraySeed" object:
#> length of object data [0] does not match object dimensions
#> [product 1000]
# Doesn't work: don't understand why I get this error.
g(FUN_with_IPC_lock, SnowParam(2L))
#> The backend is RleArray
#> Processing block = 6 on PID = 5793
#> Processing block = 7 on PID = 5793
#> Processing block = 8 on PID = 5793
#> Processing block = 9 on PID = 5793
#> Processing block = 10 on PID = 5793
#> Processing block = 1 on PID = 5787
#> Processing block = 2 on PID = 5787
#> Processing block = 3 on PID = 5787
#> Processing block = 4 on PID = 5787
#> Processing block = 5 on PID = 5787
#> Error in validObject(.Object): invalid class "ChunkedRleArraySeed" object:
#> length of object data [0] does not match object dimensions
#> [product 1000]
Created on 2018-05-21 by the reprex package (v0.2.0).
Session info
devtools::session_info()
#> Session info -------------------------------------------------------------
#> setting value
#> version R version 3.5.0 (2018-04-23)
#> system x86_64, darwin15.6.0
#> ui X11
#> language (EN)
#> collate en_AU.UTF-8
#> tz America/New_York
#> date 2018-05-21
#> Packages -----------------------------------------------------------------
#> package * version date source
#> backports 1.1.2 2017-12-13 CRAN (R 3.5.0)
#> base * 3.5.0 2018-04-24 local
#> BiocGenerics * 0.27.0 2018-05-01 Bioconductor
#> BiocParallel * 1.15.3 2018-05-11 Bioconductor
#> compiler 3.5.0 2018-04-24 local
#> datasets * 3.5.0 2018-04-24 local
#> DelayedArray * 0.7.0 2018-05-01 Bioconductor
#> devtools 1.13.5 2018-02-18 CRAN (R 3.5.0)
#> digest 0.6.15 2018-01-28 CRAN (R 3.5.0)
#> evaluate 0.10.1 2017-06-24 CRAN (R 3.5.0)
#> graphics * 3.5.0 2018-04-24 local
#> grDevices * 3.5.0 2018-04-24 local
#> HDF5Array * 1.9.0 2018-05-01 Bioconductor
#> htmltools 0.3.6 2017-04-28 CRAN (R 3.5.0)
#> IRanges * 2.15.13 2018-05-20 Bioconductor
#> knitr 1.20 2018-02-20 CRAN (R 3.5.0)
#> magrittr 1.5 2014-11-22 CRAN (R 3.5.0)
#> matrixStats * 0.53.1 2018-02-11 CRAN (R 3.5.0)
#> memoise 1.1.0 2017-04-21 CRAN (R 3.5.0)
#> methods * 3.5.0 2018-04-24 local
#> parallel * 3.5.0 2018-04-24 local
#> Rcpp 0.12.17 2018-05-18 CRAN (R 3.5.0)
#> rhdf5 * 2.25.0 2018-05-01 Bioconductor
#> Rhdf5lib 1.3.1 2018-05-17 Bioconductor
#> rmarkdown 1.9 2018-03-01 CRAN (R 3.5.0)
#> rprojroot 1.3-2 2018-01-03 CRAN (R 3.5.0)
#> S4Vectors * 0.19.5 2018-05-20 Bioconductor
#> snow 0.4-2 2016-10-14 CRAN (R 3.5.0)
#> stats * 3.5.0 2018-04-24 local
#> stats4 * 3.5.0 2018-04-24 local
#> stringi 1.2.2 2018-05-02 CRAN (R 3.5.0)
#> stringr 1.3.1 2018-05-10 CRAN (R 3.5.0)
#> tools 3.5.0 2018-04-24 local
#> utils * 3.5.0 2018-04-24 local
#> withr 2.1.2 2018-03-15 CRAN (R 3.5.0)
#> yaml 2.1.19 2018-05-01 CRAN (R 3.5.0)
- Ngôn ngữ chính
- R
- Star
- 29
- Fork
- 12
- Chỉ số merge pull request
- Không có pull request nào được merge trong 30 ngày
Hướng dẫn đóng góp
Chưa lập chỉ mục được hướng dẫn đóng góp cho kho mã nguồn này
Bắt đầu từ đâu
- Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
- Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
- Fork repository và làm thay đổi trên một nhánh.
- Mở pull request có tham chiếu số hiệu của issue.
Issue khác của Bioconductor/DelayedArray
-
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 38/100
Bioconductor/DelayedArray#129 · 11 bình luận ·
-
Custom delayed operations Đang mở
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 45/100
Bioconductor/DelayedArray#127 · 1 bình luận ·
-
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 25/100
Bioconductor/DelayedArray#125 · 1 bình luận ·
-
Độ khó 3/5 1-2 ngày Mức phù hợp với người mới 45/100
Bioconductor/DelayedArray#123 ·
-
Độ khó 5/5 Hơn một tuần Mức phù hợp với người mới 20/100
Bioconductor/DelayedArray#122 ·
Tất cả issue của Bioconductor/DelayedArray
Issue tương tự
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 75/100
robjhyndman/forecast#1220 ·
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 65/100
JamesHWade/deputy#192 ·
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 75/100
-
bug triage_needed
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 75/100
-
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 72/100
pharmaverse/rtables#1123 · 1 bình luận · 1 reaction ·