dplyr::collect() turns a query result into an R data
frame. Sometimes you want the result in Apache Arrow format instead:
- to keep working on it with the arrow package,
- to save it as Parquet files on your own disk,
- to process a result that is too large to hold in memory,
- or to hand it to another tool that reads Arrow data.
quak has two functions for this. collect_arrow() returns
the whole result at once. stream_arrow() returns it in
batches.
Setup
Both functions need the nanoarrow package. quak offers to install it the first time you call one of them. The arrow package is optional, but most of the examples below use it.
install.packages(c("nanoarrow", "arrow"))The examples use a Delta table on Azure:
library(quak)
conn <- az_conn()
az_set_chain_secret(conn, chain = "cli")
sales <- tbl_delta(conn, "abfss://container@account/path/sales")Choose a function
| You want | Use | You get |
|---|---|---|
| An R data frame | dplyr::collect() |
A tibble |
| Arrow data, all in memory | collect_arrow() |
A nanoarrow_array_stream, already read |
| Arrow data, one batch at a time | stream_arrow() |
A nanoarrow_array_stream
|
All three run the query in DuckDB. Filter, select and summarise before you collect, so that less data leaves DuckDB:
Both Arrow functions check the connection before they run, the same
way collect() does. They fail with a clear error if the
connection is closed or the azure extension is not
loaded.
Collect the whole result
collect_arrow() runs the query and reads the whole
result into memory. It returns the rows as an Arrow stream whose batches
are already loaded. The query has finished by the time it returns, so
the connection is free for other queries.
arr <- collect_arrow(large_sales)
#> ℹ Collecting data from Azure...
#> ✔ Done. 82500 rows collected in 93 ms.Like every Arrow stream, the result can be read only once. Convert it straight away and keep the converted object.
For an R data frame, use as.data.frame() or
tibble::as_tibble(). Neither needs the arrow package:
df <- as.data.frame(arr)
head(df, 3)
#> id region amount
#> 1 67 south 100.5
#> 2 68 east 102.0
#> 3 69 north 103.5The column types mostly match what collect() returns.
INTERVAL columns are the main exception. nanoarrow cannot
convert them, so cast them in the query or convert through arrow. See
?collect_arrow for the full list.
Continue with arrow
To keep working in Arrow instead, convert the result with
arrow::as_arrow_table(). Unlike the stream, a Table can be
used as often as you like. You can use arrow’s own dplyr support on it,
which runs in Arrow’s query engine rather than in R:
tab <- arrow::as_arrow_table(collect_arrow(large_sales))
tab |>
dplyr::group_by(region) |>
dplyr::summarise(total = sum(amount)) |>
dplyr::collect()
#> # A tibble: 3 × 2
#> region total
#> <chr> <dbl>
#> 1 south 3423734.
#> 2 east 3423750
#> 3 north 3423766.Or save it as a local Parquet file:
arrow::write_parquet(tab, "sales_large.parquet")Stream the result
stream_arrow() runs the query and returns a stream. The
stream produces batches of at most chunk_size rows, one at
a time, as you read them. Only the current batch is in memory, so the
whole result never has to fit.
stream <- stream_arrow(sales, chunk_size = 100000)As with collect_arrow(), a stream can be read only once.
Once a batch has been read, it cannot be read from the stream again.
Hand the stream to arrow
arrow::as_record_batch_reader() turns the stream into an
arrow RecordBatchReader. Most arrow functions accept one,
and they read it batch by batch.
Copy a large table from Azure to a partitioned Parquet dataset on your own disk, without holding it in memory:
reader <- arrow::as_record_batch_reader(stream_arrow(sales))
arrow::write_dataset(
reader,
"data/sales",
format = "parquet",
partitioning = "region"
)
list.files("data/sales", recursive = TRUE)
#> [1] "region=east/part-0.parquet" "region=north/part-0.parquet"
#> [3] "region=south/part-0.parquet"Then open the local copy with arrow::open_dataset() and
query it without going back to Azure.
You can also run dplyr verbs on the reader. arrow runs them on each batch as it arrives:
reader <- arrow::as_record_batch_reader(stream_arrow(sales))
reader |>
dplyr::filter(amount > 140) |>
dplyr::summarise(n = dplyr::n()) |>
dplyr::collect()
#> # A tibble: 1 × 1
#> n
#> <int>
#> 1 15000When the same filter can run in DuckDB, apply it to
sales before streaming instead. DuckDB then reads less data
from Azure.
Process batches in R
To work through the batches yourself, call
stream$get_next(). It returns the next batch as a
nanoarrow_array, or NULL when the stream is
finished:
stream <- stream_arrow(sales, chunk_size = 100000)
total <- 0
while (!is.null(batch <- stream$get_next())) {
df <- as.data.frame(batch)
total <- total + sum(df$amount)
}Stop reading early
If you stop before the end, call stream$release(). It
frees the query and its memory straight away, instead of whenever R next
collects garbage:
stream <- stream_arrow(sales)
first <- stream$get_next()
stream$release()Keep the connection free while a stream is open
Read a stream to the end, or release it, before you run any other query on the same connection. DuckDB ends an open stream as soon as its connection runs another query. The stream then acts as if it had finished normally, without an error, so the rows it had not yet returned are lost:
stream <- stream_arrow(sales, chunk_size = 10)
stream$get_next()$length
#> [1] 10
DBI::dbGetQuery(conn, "SELECT 42") # another query on the same connection
stream$get_next()
#> NULLThis also applies to a RecordBatchReader made from a
stream, because arrow reads it lazily. Finish with the reader, for
example by letting arrow::write_dataset() or
dplyr::collect() complete, before you use the connection
again. The behaviour is reported upstream as duckdb/duckdb-r#2772.
If you need to query while a stream is open, use a second connection.
Queries on another connection do not affect the stream.
az_conn() opens a new, empty database, so register a secret
on it again. Tables and views that you created with
load_delta() and similar functions on the first connection
are not visible there.
Choose a batch size
chunk_size is the largest number of rows in one batch.
The default is one million. Smaller batches use less memory each. Larger
batches mean fewer steps and less overhead per row. Lower it when rows
are wide, for example when they hold long strings or nested columns.
DuckDB allocates each batch for chunk_size rows before
filling it, even when the result has fewer rows. Keep
chunk_size modest: a value in the hundreds of millions can
reserve gigabytes of memory.