Skip to content

fread support for parquet #2505

Description

@MichaelChirico

pandas has pandas.read_parquet('file_or_dir') and it's pretty spot on. No spark dependency, just requires pyarrow in the background.

My admittedly cursory search only turned up SparkR::read.parquet and a sparklyr version, neither of which are well-documented... better, IMO, to just fread('file_or_dir').

Activity

  1. DavidArenburg commented on Dec 6, 2017

    @DavidArenburg
  2. MichaelChirico commented on Aug 16, 2018

    @MichaelChirico
    MemberAuthor

    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:

    jimhester/rarrow#1

    So, closing this pending updates there. Ideally, speed from arrow::read_parquet will be fast enough & we can just setDT.

    Seems to be more logical to me than expending substantial engineering effort here building our own reader.

  3. DrMaphuse commented on Jan 26, 2023

    @DrMaphuse

    I 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?

  4. eddelbuettel commented on Aug 27, 2025

    @eddelbuettel
    Contributor

    I looked into this again recently especially given how performant (and otherwise self-sufficient: single binary) duckdb is 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 around nanoarrow, 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.table for a directory full of csv.xz files.

    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 by is 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)$ 
  5. jangorecki commented on Aug 27, 2025

    @jangorecki
    Member

    @DrMaphuse

    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?

  6. eddelbuettel commented on Aug 27, 2025

    @eddelbuettel
    Contributor

    @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

    1. @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. ↩

  7. eitsupi commented on Aug 28, 2025

    @eitsupi
    Contributor

    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

  8. eddelbuettel commented on Aug 28, 2025

    @eddelbuettel
    Contributor

    @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

    1. The parquet use case at the beginning of this thread may be addressed by nanoparquet. ↩

    2. Especially as arrow as you note deals already with ALTREP. ↩

  9. MichaelChirico commented on Aug 28, 2025

    @MichaelChirico
    MemberAuthor

    could you clarify what you mean by "data.table representation" (as opposed to, say, "data.frame representation")?

  10. eitsupi commented on Aug 28, 2025

    @eitsupi
    Contributor

    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.

  11. eddelbuettel commented on Aug 28, 2025

    @eddelbuettel
    Contributor

    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.

  12. aitap commented on Aug 28, 2025

    @aitap
    Member

    Part of the problem is data.table currently being allergic to ALTREP vectors (cf. the second-to-last call in setDT). They are sources of data sharing, don't have a TRUELENGTH and cannot be directly resized. The latter two problems we have to somehow solve anyway (#6180), but the former one may be even harder (#4783).

  13. eddelbuettel commented on Aug 28, 2025

    @eddelbuettel
    Contributor

    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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions