---
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.
