--- title: "Performance and caching" output: rmarkdown::html_vignette vignette: > %\VignetteIndexEntry{Performance and caching} %\VignetteEngine{knitr::rmarkdown} %\VignetteEncoding{UTF-8} --- ```{r, include = FALSE} knitr::opts_chunk$set( collapse = TRUE, comment = "#>", eval = FALSE ) ``` This guide explains how `delta.sharing` downloads and caches shared files, and how concurrency, materializer choice, and batch size affect reads. Start with `vignette("delta-sharing")` if you have not yet created a client and read a table. A read has two distinct phases: **staging files** and **reading rows**. The package writes the files selected by the sharing server directly to its disk-backed session cache before opening an Arrow stream; complete downloaded files are not held in R memory. Eager materializers can still exhaust RAM by collecting the scan output. A lazy Arrow reader avoids that collection, but it does not avoid staging or stream remote Parquet files directly from their signed URLs. ## The read lifecycle Every snapshot or change read follows the same broad path: 1. R requests the relevant file actions from the Delta Sharing server. 2. Missing Parquet and deletion-vector files are downloaded into the table's session cache. 3. R writes a new local Delta log that points to the cached files. 4. Delta Kernel scans the local files and emits Arrow batches. 5. An eager materializer, when used, collects those batches into its result. The stages have different reuse behavior: | Work | Happens on another materializer call? | Reused from the session cache? | |---|---:|---:| | Sharing query and pagination | Yes | No | | Download selected files | Only when missing or invalid | Yes | | Build the local Delta log | Yes | No | | Delta Kernel scan | Yes | No | | Collect an eager result | Yes | No | Reusing cached files avoids downloading them again, but each materializer call still performs the Sharing request, local log construction, Delta Kernel scan, and result collection. ## What the cache contains Each table has a deterministic cache directory under R's session temporary directory. Its identity includes the profile endpoint and the table's share, schema, and name. New handles for the same table therefore share cached files, even when they use different download concurrency: ```{r, eval = FALSE} client <- sharing_client("~/config.share") orders <- client$table("sales.default.orders", concurrency = 4) faster_downloads <- client$table("sales.default.orders", concurrency = 8) identical(orders$cache_path, faster_downloads$cache_path) #> TRUE ``` Snapshot and change data feed reads for the same table use this shared cache. When the server selects the same underlying file for both reads, its local copy is reused. Files that do not overlap are downloaded normally. The cache contains local copies of selected Parquet and deletion-vector files. Files are identified using immutable IDs supplied by the sharing server, so they can be reused even after a signed download URL changes. When the server provides a file size, the package checks it before reusing the cached file. Downloads are completed in temporary files before becoming available to the cache, so interrupted transfers are not reused. Cache directories and files are created with user-only permissions. ## Cache lifetime The cache lasts for the R session, not for the lifetime of a table handle. It lives under `tempdir()`, R's per-session temporary directory. It is normally removed when the session ends, but it is not durable storage: abnormal termination may leave files behind, and system temporary-directory policies may remove them during a long-running session. See `?tempdir` for how R chooses the directory. The cache is not persistent across sessions and does not contain materialized tibbles or Arrow tables. The read-only `cache_path` property is available for inspection: ```{r, eval = FALSE} orders$cache_path fs::dir_info(orders$cache_path) ``` After all lazy readers and streams for the table have been closed, the directory can be removed to reclaim disk space: ```{r, eval = FALSE} fs::dir_delete(orders$cache_path) ``` The next read recreates the directory and downloads any required files. A long session that touches many large tables can retain substantial data, so `cache_path` is also useful for monitoring local disk usage. ## Download concurrency Set concurrency when creating a table handle: ```{r, eval = FALSE} orders <- client$table( "sales.default.orders", concurrency = 8 ) ``` The default is four. `concurrency` controls the maximum number of missing files downloaded at once. It does not affect Sharing API pagination, Delta Kernel scanning, Arrow batch size, or reads whose files are already cached. Higher concurrency is most useful when a query selects many files and network latency leaves the connection idle. It may make little difference when there are few files, a single file saturates the connection, or the provider throttles requests. Benchmark representative uncached reads before increasing it. Interactive reads show progress for missing downloads. When every file has a known size, the display also shows the total bytes to download. Cache hits are not counted as downloads. ## Shape the read before tuning it Query options affect different stages: - `columns` reduces local scan and materialization work, but not network transfer or cache disk usage. The sharing server selects files before Delta Kernel applies the projection, and the complete selected Parquet files are downloaded. - `limit` is sent to the server as a hint and is enforced exactly by Delta Kernel. A provider may use the hint to return fewer files, but this is not guaranteed. - `predicate` is a best-effort server hint. It may improve file pruning, but it is not an exact row filter; apply an exact filter in the downstream consumer. ## Choose a materializer Materializers differ in how they consume Arrow batches: | Method | Collection behavior | Typical use | |---|---|---| | `to_tibble()` | Eager, in R memory | Ordinary R analysis | | `to_data_frame()` | Eager, via the tibble path | Base data-frame consumers | | `to_arrow()` | Eager, in Arrow memory | Repeated Arrow or DuckDB scans | | `to_arrow_reader()` | Lazy after file staging | One-pass Arrow or DuckDB scan | | `to_arrow_stream()` | Low-level and lazy after staging | Arrow C Stream consumers | Arrow is a required dependency. Eager R materializers use Arrow's converter and return BIGINT columns as `bit64::integer64`. Close an Arrow reader when it is no longer needed, and release a low-level stream if its consumer does not take ownership. ### Query larger-than-memory results with DuckDB When the analysis can be expressed in SQL, DuckDB is usually simpler than processing Arrow batches manually. DuckDB can spill larger-than-memory operations to disk, while `delta.sharing` keeps the selected source files in its disk-backed session cache. This requires the optional `DBI`, `duckdb`, and `withr` packages. Register a lazy Arrow reader so the full table is not first collected in R: ```{r, eval = FALSE} reader <- orders$snapshot( columns = c("status", "amount") )$to_arrow_reader() con <- DBI::dbConnect(duckdb::duckdb()) duckdb::duckdb_register_arrow(con, "orders", reader) summary <- withr::with_options( list(arrow.use_threads = FALSE), DBI::dbGetQuery(con, " SELECT status, count(*) AS orders, sum(amount) AS total_amount FROM orders GROUP BY status ORDER BY status ") ) duckdb::duckdb_unregister_arrow(con, "orders") DBI::dbDisconnect(con) ``` Let Arrow manage the registered reader's lifetime. Calling `reader$Close()` directly can release its stream while the scanner is still reading ahead, including after a SQL query returns early (for example, with `LIMIT`). Arrow registration does not copy the full source into DuckDB. Only the final query result returned by `dbGetQuery()` is collected into R memory. Use SQL to filter or aggregate the data to a manageable result; collecting a very large final result can still exhaust R's memory. The local Arrow option disables parallel scan execution; background read-ahead can still occur. DuckDB can execute the rest of the query in parallel. See DuckDB's [larger-than-memory documentation](https://duckdb.org/docs/current/guides/performance/how_to_tune_workloads#larger-than-memory-workloads-out-of-core-processing) and its [Arrow registration reference](https://r.duckdb.org/reference/duckdb_register_arrow.html) for details and current limitations. ## Arrow batch size `batch_size` controls the maximum number of rows emitted in each Arrow batch. The default is 65,536 rows. It does not control the number or size of downloaded files. Smaller batches can reduce peak memory while a batch is being processed, at the cost of more conversion and boundary overhead. Larger batches can reduce that overhead but use more memory. An eager materializer still collects every batch, so reducing `batch_size` does not make `to_tibble()`, `to_data_frame()`, or `to_arrow()` safe for a result larger than memory. Change the default only after measuring a representative workload. ## A practical tuning order 1. Use `columns` to reduce local scan and materialization work. Use `limit` and predicates as server-side pruning hints, while respecting their exact and best-effort semantics. 2. Choose the execution path for the result size: eager R materialization when it fits comfortably in memory, or DuckDB for larger-than-memory SQL analytics. 3. Reuse an existing result when the same data is needed again. Let the session cache handle overlapping reads automatically. 4. If downloading uncached files is the bottleneck, compare a small number of concurrency values on a representative read. 5. Tune `batch_size` last and only when measurements show batch overhead or memory pressure. Measure end-to-end elapsed time, peak memory, and cache disk usage. Network throughput, latency, provider behavior, file count, compression, schema width, query shape, and local hardware all affect performance, so there is no universally optimal configuration.