Repository navigation
fread support for parquet #2505
Description
Activity
The latest seems to be that Romain Francois is working to build out the R bindings to Apache Arrow:
https://github.com/romainfrancois/arrow/
According to my read of the discussion here, this is by far the simplest approach to building a local (i.e. spark-context-free) parquet reader:
So, closing this pending updates there. Ideally, speed from
arrow::read_parquetwill be fast enough & we can justsetDT.Seems to be more logical to me than expending substantial engineering effort here building our own reader.
Reacted by Jan Gorecki and Oliver Per MadsenI would love to reopen this request. As mentioned in #2026 (comment), setDT is a bottleneck when working with big datasets. I'm not sure if this makes sense, architecturally, but maybe optimizing setDT to work directly on arrow tables could be a good solution?
Reacted by Jonathan Alzetta, Jakob Weickmann, isaacnorwich, Dan Simonet, Indrajeet Patil, Nir Grinberg, Elise Hellwig, Apoorva Lal and Grant McDermottI looked into this again recently especially given how performant (and otherwise self-sufficient: single binary)
duckdbis at read multiple compressed files (see below). We can get an Arrow memory representation from it, but I can for the life of me not an efficient data.table instantiation. I (more or less) know my way aroundnanoarrow, as does @eitsupi.I wrote a naive 'from Arrow via nanoarrow to Rcpp to data.frame to setDT' but it takes forever. The question is what to put into the middle: we can get the two
void*for an Arrow table object (or the batched RecordStream), at that point we are still lightweight. I now need the missing middle to end up at 'profit'. I am aware that materializing has a cost, the question is what the minimal cost is. I could live with a few seconds for something large enough.Small arrow demo below. The csv.zstd files are here if you need them; I cannot get anywhere near this with
data.tablefor a directory full ofcsv.xzfiles.edd@paul:~/git/r2u-logs(master)$ /usr/bin/time duckdb temp "create temp table r as select * from read_csv('r*/*.csv.zst'); select count(*) as n from r;" 100% ▕████████████████████████████████████████████████████████████▏ ┌─────────────────┐ │ n │ │ int64 │ ├─────────────────┤ │ 48862271 │ │ (48.86 million) │ └─────────────────┘ 26.02user 3.33system 0:03.14elapsed 933%CPU (0avgtext+0avgdata 9730840maxresident)k 423256inputs+24outputs (171major+2660019minor)pagefaults 0swaps edd@paul:~/git/r2u-logs(master)$ edd@paul:~/git/r2u-logs(master)$ ls -1 */*csv.zst r2u/r2u_r2u-2022.csv.zst r2u/r2u_r2u-2023.csv.zst r2u/r2u_r2u-2024.csv.zst r2u/r2u_r2u-2025-07.csv.zst r2u/r2u_r2u-2025-q1.csv.zst r2u/r2u_r2u-2025-q2.csv.zst rob/r2u_rob-2022.csv.zst rob/r2u_rob-2023.csv.zst rob/r2u_rob-2024.csv.zst rob/r2u_rob-2025-07.csv.zst rob/r2u_rob-2025-q1.csv.zst rob/r2u_rob-2025-q2.csv.zst edd@paul:~/git/r2u-logs(master)$
Three seconds is pretty good for a dozen zstd compressed csv files, read and counted.
group byis also effin fast.edd@paul:~/git/r2u-logs(master)$ time duckdb temp "create temp table r as select * from read_csv('r*/*.csv.zst'); select count(*) as n, strftime(date, '%Y-%m') as ym from r group by ym order by ym limit 10;" 100% ▕████████████████████████████████████████████████████████████▏ ┌────────┬─────────┐ │ n │ ym │ │ int64 │ varchar │ ├────────┼─────────┤ │ 124504 │ 2022-05 │ │ 67952 │ 2022-06 │ │ 64357 │ 2022-07 │ │ 92528 │ 2022-08 │ │ 147136 │ 2022-09 │ │ 155306 │ 2022-10 │ │ 301534 │ 2022-11 │ │ 347006 │ 2022-12 │ │ 981534 │ 2023-01 │ │ 898208 │ 2023-02 │ ├────────┴─────────┤ │ 10 rows │ └──────────────────┘ real 0m3.450s user 0m41.018s sys 0m3.435s edd@paul:~/git/r2u-logs(master)$
setDT is a bottleneck when working with big datasets
This is not true. When you call setDT, what is most likely happening, all parquet data area being pulled into data.frame memory layout. If you first materialize parquet table as a data.frame, then it will take no time to call setDT.
It is not possible to call setDT directly on parquet because it's memory layout is not compatible to data.frame, and there has to be an in-memory copy from parquet format to data.frame format.Follow Dirk on his efforts mentioned above as he is basically trying to minimize the overhead of that copy.
@eddelbuettel I am unfortunately not able to help with it, @aitap any ideas?
Reacted by Dirk Eddelbuettel@jangorecki A good first step may be an 'arrow2altrep' helper function but sadly I am also not all that fluent in building altrep vectors. @aitap @eitsupi : Thoughts? 1
Footnotes
-
@jangorecki you replied to a post from 2023. I just came here because I didn't want to open a redundant issue; this seems like a 'close enough' issue for arrow (or feather, parquet, ...) to data.table questions. ↩
-
I don't fully understand the context, but I think altrep is used when reading using functions from the arrow package.
(This must be happening when converting an Arrow Array to an R vector.)tmpf <- withr::local_tempfile() arrow::write_parquet(penguins, tmpf) df <- arrow::read_parquet(tmpf, as_data_frame = FALSE) |> as.data.frame() .Internal(inspect(df)) #> @560b8098fc38 19 VECSXP g0c4 [OBJ,REF(2),ATT] (len=8, tl=0) #> @560b80908100 13 INTSXP g0c0 [OBJ,REF(65535),ATT] arrow::array_factor<0x7f7758004420, dictionary<values=string, indices=int8, ordered=0>, 1 chunks, 0 nulls> len=344 #> ATTRIB: #> @560b80907ae0 02 LISTSXP g0c0 [REF(1)] #> TAG: @560b79b5de60 01 SYMSXP g1c0 [MARK,REF(65535),LCK,gp=0x4000] "levels" (has value) #> @560b80907b18 16 STRSXP g0c0 [REF(65535)] arrow::array_string_vector<0x560b80093cd0, string, 1 chunks, 0 nulls> len=3 #> TAG: @560b79b5e2c0 01 SYMSXP g1c0 [MARK,REF(41569),LCK,gp=0x6000] "class" (has value) #> @560b7d4ca310 16 STRSXP g0c1 [REF(65535)] (len=1, tl=0) #> @560b79bf4a90 09 CHARSXP g1c1 [MARK,REF(371),gp=0x61] [ASCII] [cached] "factor" #> @560b8090c050 13 INTSXP g0c0 [OBJ,REF(65535),ATT] arrow::array_factor<0x7f775c004a50, dictionary<values=string, indices=int8, ordered=0>, 1 chunks, 0 nulls> len=344 #> ATTRIB: #> @560b8090baa0 02 LISTSXP g0c0 [REF(1)] #> TAG: @560b79b5de60 01 SYMSXP g1c0 [MARK,REF(65535),LCK,gp=0x4000] "levels" (has value) #> @560b8090bad8 16 STRSXP g0c0 [REF(65535)] arrow::array_string_vector<0x560b7f5a9f00, string, 1 chunks, 0 nulls> len=3 #> TAG: @560b79b5e2c0 01 SYMSXP g1c0 [MARK,REF(41569),LCK,gp=0x6000] "class" (has value) #> @560b7d4ca310 16 STRSXP g0c1 [REF(65535)] (len=1, tl=0) #> @560b79bf4a90 09 CHARSXP g1c1 [MARK,REF(371),gp=0x61] [ASCII] [cached] "factor" #> @560b8090b720 14 REALSXP g0c0 [REF(65535)] arrow::array_dbl_vector<0x7f77640032a0, double, 1 chunks, 2 nulls> len=344 #> @560b8090b3d8 14 REALSXP g0c0 [REF(65535)] arrow::array_dbl_vector<0x7f77600038b0, double, 1 chunks, 2 nulls> len=344 #> @560b8090b090 13 INTSXP g0c0 [REF(65535)] arrow::array_int_vector<0x7f776c0039d0, int32, 1 chunks, 2 nulls> len=344 #> ... #> ATTRIB: #> @560b809d25a0 02 LISTSXP g0c0 [REF(1)] #> TAG: @560b79b5dd10 01 SYMSXP g1c0 [MARK,REF(65535),LCK,gp=0x4000] "names" (has value) #> @560b80897678 16 STRSXP g0c4 [REF(65535)] (len=8, tl=0) #> @560b7f3b0228 09 CHARSXP g0c1 [REF(17),gp=0x60] [ASCII] [cached] "species" #> @560b7f3b01f0 09 CHARSXP g0c1 [REF(17),gp=0x60] [ASCII] [cached] "island" #> @560b7f372e08 09 CHARSXP g0c2 [REF(17),gp=0x60] [ASCII] [cached] "bill_len" #> @560b7f372dc8 09 CHARSXP g0c2 [REF(17),gp=0x60] [ASCII] [cached] "bill_dep" #> @560b7f372d88 09 CHARSXP g0c2 [REF(17),gp=0x60] [ASCII] [cached] "flipper_len" #> ... #> TAG: @560b79b5dae0 01 SYMSXP g1c0 [MARK,REF(65535),LCK,gp=0x4000] "row.names" (has value) #> @560b808fcc60 13 INTSXP g0c1 [REF(65535)] (len=2, tl=0) -2147483648,-344 #> TAG: @560b79b5e2c0 01 SYMSXP g1c0 [MARK,REF(41569),LCK,gp=0x6000] "class" (has value) #> @560b809cff98 16 STRSXP g0c1 [REF(1)] (len=1, tl=0) #> @560b79bff9a8 09 CHARSXP g1c2 [MARK,REF(621),gp=0x61,ATT] [ASCII] [cached] "data.frame"
Created on 2025-08-28 with reprex v2.1.1
@eitsupi Sorry I did not make this very clear: The quest is more about a generic facility to go from Arrow's memory representation as efficiently as possible to data.table, and less about parquet 1. Preferably using only nanoarrow, but if it is done via a mediating package arrow could be used 2. My starting point is duckdb: get (stream of) RecordBatch records, make it data.table. I have yet to find a way that does not involve materialization, and while that step may in fact unavoidable for data.table I am effectively looking for fastest / lightest way from arrow memory to data.table representation.
Footnotes
could you clarify what you mean by "data.table representation" (as opposed to, say, "data.frame representation")?
Ah, I understand.
Indeed, if we could achieve a fast conversion from Arrow Array Stream to data.table, we would benefit greatly from the nanoarrow ecosystem (the nanoarrow_array_stream class).My understanding is that the arrow package is currently the fastest for converting from an Arrow memory representation to an R data.frame.
could you clarify what you mean by "data.table representation" (as opposed to, say, "data.frame representation")?
I take data.frame as intermediary as my goal (as for many of these arrow/feather/parquet/... issues here) is being in data.table. Hence my trying to restart a discussion in the data.table repository.
Part of the problem is
data.tablecurrently being allergic toALTREPvectors (cf. the second-to-last call insetDT). They are sources of data sharing, don't have aTRUELENGTHand cannot be directly resized. The latter two problems we have to somehow solve anyway (#6180), but the former one may be even harder (#4783).So then materializing is the only way to go.
From (nano)arrow we get chunks, but know total length so it would be (at least) two nested loops of memcpy (over all the chunks, within each chunk over all columns). I wrote some (nano)arrow to R converters in the past, part of the problem of course also is that we not always get memcpy as the set of types in arrow is larger and we sometimes need to cast. Which when done element by element is not exactly fast.
pandashaspandas.read_parquet('file_or_dir')and it's pretty spot on. Nosparkdependency, just requirespyarrowin the background.My admittedly cursory search only turned up
SparkR::read.parquetand asparklyrversion, neither of which are well-documented... better, IMO, to justfread('file_or_dir').