Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
214 changes: 67 additions & 147 deletions versions/4.x/entity-orm.md
Original file line number Diff line number Diff line change
Expand Up @@ -2427,25 +2427,24 @@ Working with large datasets can be challenging in terms of performance and memor
This module introduces a variety of data handling strategies tailored to different use cases, including:

- Chunk-based processing to break large datasets into manageable segments.
- Keyset chunking (`chunkById()`) for huge tables and for jobs that modify the rows they read.
- Cursor-based streaming for low-memory iteration over millions of records.
- Generators to provide lazy, memory-efficient traversal of datasets.
- Fiber-based concurrency for advanced parallel data processing and streaming.

All these approaches are designed with memory safety in mind, incorporating garbage collection and careful resource cleanup to avoid memory leaks or excessive overhead during execution.
All these approaches are designed with memory safety in mind and release each page of models as soon as it has been processed.

Whether you're building scalable ETL pipelines or simply need to process records in a loop without crashing your application, Doppar’s dataset handling capabilities provide a robust and flexible foundation.

### chunk
When dealing with large datasets, loading all records into memory at once can cause significant performance degradation or even out-of-memory errors. The `chunk()` method addresses this by breaking the dataset into smaller chunks and processing each chunk independently. This allows you to iterate through the entire dataset efficiently, without overwhelming system resources.
When dealing with large datasets, loading all records into memory at once can cause significant performance degradation or even out-of-memory errors. The `chunk()` method addresses this by breaking the dataset into smaller chunks and processing each chunk independently.

The `chunk()` method accepts three parameters:
- `$chunkSize` – the number of records to retrieve per batch.
- `$processor` – a callback function that receives each chunk for processing.
- `$total (optional)` – the total number of records expected, useful for tracking progress.
- `$chunkSize` – the number of records to retrieve per page. Must be greater than zero.
- `$processor` – a callback that receives each chunk (a `Collection`), the number of records processed so far (including the current chunk), and `$total`.
- `$total (optional)` – the total number of records expected, passed through to the callback for progress tracking.

Each chunk is fetched using pagination via `LIMIT` and `OFFSET`. After processing a chunk, the memory is freed, and garbage collection is invoked to minimize leaks.
Each chunk is fetched with `LIMIT` and `OFFSET`. If your query has no `ORDER BY`, Doppar orders by the primary key so pages never overlap or skip rows. A `limit()` / `offset()` already set on the query is respected, and iteration stops as soon as a short page is returned (no extra empty query).

Below is an example where we process active users in chunks of 100 records:
```php
use App\Models\User;
use Phaseolies\Support\Collection;
Expand All @@ -2458,136 +2457,87 @@ User::query()
}
});
```
This approach keeps memory usage low and is ideal for operations such as batch updates, exporting data, or applying transformations to large tables.

### Tracking Progress with chunk()
In addition to memory-efficient processing, the `chunk()` method in Doppar also supports progress tracking. This is especially useful when you want to provide feedback during long-running operations, such as logging progress, displaying status bars, or estimating time to completion.
> ⚠️ Do not update or delete rows inside `chunk()` in a way that changes which rows match your `WHERE` clause. Later pages shift and rows get skipped. Use `chunkById()` for that.

By passing a third parameter, `$total`, to the `chunk()` method, your processor callback will receive not only the current chunk of data but also:
### Tracking Progress with chunk()
By passing a third parameter, `$total`, to `chunk()`, your processor callback receives the running count and the total:
```php
User::query()
->where('status', true)
->chunk(
chunkSize: 100,
processor: function (Collection $users, $processed, $total) {
foreach ($users as $user) {
// Perform operations on each user
echo "Processed $processed of $total\n";
}
echo "Processed $processed of $total\n";
},
total: 500
);
```
This pattern is ideal when you're running background tasks or CLI scripts and want real-time feedback or logging.

### fchunk() – Concurrent Chunk
When speed is critical and your system can handle parallel operations, Doppar provides the `fchunk()` method—a fiber-based parallel chunk processor. This allows multiple chunks to be processed concurrently, taking advantage of lightweight PHP Fibers to speed up the total processing time.
### chunkById() – Keyset Chunking
`chunkById()` pages with `WHERE id > :last ORDER BY id LIMIT n` instead of `OFFSET`. Two benefits:

Unlike traditional `chunk()` processing, which is sequential, `fchunk()` runs multiple chunks in parallel, each in its own fiber. This can significantly reduce the time needed to process large datasets, especially when the per-record logic is I/O-bound or computationally light.
- It stays fast on very large tables. `OFFSET` makes the database scan and discard every earlier row on each page, so `chunk()` gets slower the deeper it goes.
- It is safe to update or delete the rows you are processing, even if they stop matching your `WHERE` clause.

Parameters
- **chunkSize:** Number of records per chunk.
- **processor:** Callback to process each chunk of records.
- **concurrency:** Number of fibers (chunks) to process in parallel.
- **chunkSize:** Number of records per page.
- **processor:** Callback receiving the chunk, the running processed count and `$total`.
- **column (optional):** The unique, sortable column to page on. Defaults to the model's primary key.
- **total (optional):** Passed through to the callback.

> ⚠️ Note: Fiber-based concurrency is cooperative, not true multithreading. It still runs on a single PHP thread, but enables interleaving tasks efficiently.

Example with concurrently processing active users
```php
User::query()
->where('status', true)
->fchunk(
chunkSize: 100,
processor: function (Collection $users) {
foreach ($users as $user) {
// Handle user logic
}
},
concurrency: 4
);
->where('status', 'pending')
->chunkById(500, function (Collection $users) {
foreach ($users as $user) {
$user->update(['status' => 'processed']);
}
});
```
In this example:
- Users are fetched in chunks of 100.
- Up to 4 chunks are processed in parallel
- Each chunk is passed to the callback where individual records are handled.

This technique is ideal for high-throughput batch jobs where database or API operations can be overlapped efficiently. It's especially helpful in CLI scripts or daemons where time-to-completion matters.
> `chunkById()` ignores any `orderBy()`, `limit()` and `offset()` on the query, because it controls the order and paging itself.

### cursor() – Stream Records
The `cursor()` method is designed for ultra-efficient, low-memory iteration over large datasets. Instead of loading all records or chunks into memory, `cursor()` uses a forward-only PDO cursor to stream records one at a time directly from the database.

This method is ideal when dealing with millions of records, performing exports, or executing operations where memory usage must stay constant.
The `cursor()` method iterates over a result set one model at a time through a forward-only PDO cursor, so only one hydrated model is alive at once.

Parameters
- **processor:** A callback function that is invoked for each record.
- **processor:** Callback invoked for each record. It receives the model, the 1-based index of the record, and `$total`.
- **total (optional):** Total number of records expected, useful for progress tracking.
- **unbuffered (optional):** Stream rows from the server instead of buffering the result set (MySQL only, see below). Defaults to `false`.

### How It Works
Internally, `cursor()` prepares the query, binds parameters, and fetches each row as an associative array. Each row is immediately converted into a model and passed to the callback function. After each record is processed, the memory is cleaned up and garbage collection is triggered.
```php
User::query()
->where('status', true)
->cursor(function ($user) {
// Handle user logic
});
```
The operation is memory-stable regardless of dataset size.

- *Best for:* Iterating huge datasets with low memory usage.

> ⚠️ Avoid using cursor() if your processing logic requires sorting, seeking, or buffering large sets of rows.

You can also handle progress in cursor like this way
```php
User::query()
->where('status', true)
->cursor(function ($user, $processed, $total) {
// Handle user logic
echo "Processed $processed of $total\n";
}, 500);
```

### fcursor() – Hybrid Cursor
The `fcursor()` method blends the memory efficiency of a traditional cursor with the performance benefits of Fibers. This hybrid approach lets you stream records from the database in batches using a cursor, but buffer and process them efficiently using Fiber-based suspension.
Rows are hydrated the same way as `get()` (connection name, encrypted attributes, and so on). Eager loading is not applied because rows are handled one at a time.

Unlike `cursor()` which processes one row at a time, `fcursor()` buffers a group of records (using a configurable bufferSize) and yields them together to the processor. This reduces the overhead of per-row context switching and provides a balance between memory usage and performance.
#### Buffered vs. unbuffered (MySQL)
By default the MySQL PDO driver downloads the **entire** result set into PHP memory when the query runs, so a plain `cursor()` limits the number of *models* in memory, not the raw rows. For truly constant memory on huge result sets pass `unbuffered: true`:

Parameters
- **processor:** A callback function invoked for each record in the buffered batch.
- **bufferSize:** The number of records to accumulate before passing them to the processor (default is 1000).

#### How It Works
- A PDO cursor fetches records from the database.
- Records are buffered into an array.
- Once the buffer reaches the defined bufferSize, the Fiber suspends and yields control to the processor.
- Processing resumes after the callback is done, and continues until all records are exhausted.

Example with user list streaming
```php
User::query()
->where('status', true)
->fcursor(function ($user) {
//
});
User::query()->cursor(function ($user) {
// ...
}, unbuffered: true);
```
This method is particularly useful when:

- You want better performance than single-row cursors.
- Your logic can handle medium-sized memory buffers.
- You’re building CLI tools or background workers that must process large volumes reliably.
While an unbuffered cursor is open, the connection cannot run any other query. Do not read from or write to the database on the same connection inside the callback; collect IDs and act on them afterwards, or use `chunkById()` instead. The previous buffering mode is restored when the cursor finishes. SQLite and PostgreSQL ignore this option.

> ⚡ Efficient like cursor(), but faster for many real-world workloads due to fiber-based batching.
- *Best for:* Read-only exports and reports over very large result sets.

### stream() – Lazy Generator
The `stream()` method offers a clean, generator-based approach to iterating over large datasets in a memory-efficient way. It loads records in small chunks `(like chunk())`, but yields one record at a time, allowing you to use it with foreach just like a native PHP array — without holding everything in memory.
> ⚠️ Avoid using cursor() if your processing logic requires sorting, seeking, or buffering large sets of rows.

Unlike `cursor()`, `stream()` handles records in small batches (e.g., 100 at a time) and is suitable when you prefer chunked I/O with on-the-fly transformation or mapping.
### stream() – Lazy Generator
The `stream()` method offers a clean, generator-based approach to iterating over large datasets in a memory-efficient way. It loads records in pages like `chunk()` (same ordering and `limit()` / `offset()` rules), but yields one record at a time, so you can use it with `foreach`.

Parameters
- **chunkSize:** The number of records to fetch per internal chunk.
- **transform (optional):** A callback function to transform each model before yielding.
- **chunkSize:** The number of records to fetch per internal page.
- **transform (optional):** A callback to transform each model before yielding.

Example with Lazy Streaming. Fetch 100 at a time lazily
```php
$users = User::query()
->where('status', true)
Expand All @@ -2598,10 +2548,7 @@ foreach ($users as $user) {
}
```

In this example Each user is yielded individually from the generator. Memory usage stays consistent regardless of dataset size.

## Transforming Models During Stream
You can also transform records before they are yielded using the optional `transform` parameter.
#### Transforming Models During Stream
```php
foreach (
User::query()
Expand All @@ -2612,78 +2559,51 @@ foreach (
dump($userName);
}
```
In this example
- Here, each user is mapped to an uppercase string before being yielded.
- Each user is mapped to an uppercase string before being yielded.
- The result is a generator of transformed values (not models).
- Useful for lightweight export pipelines, CSV dumps, or messaging queues.
- Useful for lightweight export pipelines, CSV dumps, or message queues.

> 💡 Tip: Use stream() when you want a foreach-friendly syntax with low memory use, especially in long-running CLI tasks or APIs.

### fstream() – Streaming with Backpressure Control
The `fstream()` method is an advanced Fiber-powered generator for streaming large datasets efficiently. It combines the lazy iteration style of `stream()` with the performance and concurrency control of Fibers.

Unlike `stream()`, `fstream()` is designed for high-throughput environments where you may want to buffer records in controlled amounts and suspend/resume execution more precisely.

Parameters
- **chunkSize:** Number of records to fetch per database request.
- **transform (optional):** Callback to transform each model before yielding.
- **bufferSize:** Maximum number of items to buffer before yielding (defaults to 1000).

Example: Streaming with Transformation
```php
foreach (
User::query()
->where('status', true)
->fstream(1000, fn($user) => strtoupper($user->name))
as $userName
) {
dump($userName);
}
```

In this example
- Records are fetched in chunks of 1000.
- Transformed immediately into uppercase strings.
- Buffer is filled (up to default 1000) before being flushed into the foreach loop.
- Memory remains stable, and Fiber context switching keeps performance high.

#### How It Works
- Internally, Fibers pull in small batches of records.
- Each record is passed through the transform (if given).
- Fibers suspend after each yield, keeping processing responsive.
- A buffer size (default 1000) prevents overloads while maximizing I/O efficiency.

### batch() – Grouped Batch Processing
The `batch()` method is tailored for scenarios where you want to process large numbers of records in grouped batches instead of one at a time. It provides an efficient way to collect data in memory-limited chunks and flush it to a processor in meaningful units (e.g., every 1000 records).

This is particularly useful when writing to external systems, doing bulk inserts, or when you want to limit the number of operations per database or API call.
### batch() – Grouped Batch Processing
The `batch()` method collects records into groups and hands each group to a processor. It is useful for bulk inserts, or when you want to limit the number of calls made to an external system.

Parameters
- **chunkSize:** Number of records to fetch from the database per query.
- **batchProcessor:** A callback that receives each batch as a Collection.
- **batchSize:** Number of records per batch to accumulate before processing (default: 1000).
- **batchSize:** Number of records per batch (default: 1000).

Example: Process Users in Batches of 1000
```php
User::query()
->where('status', true)
->batch(
chunkSize: 500,
batchProcessor: function ($batch) {
foreach ($batch as $user) {
//
}

echo "Processed batch of {$batch->count()} users\n";
},
batchSize: 1000
);
```

In this example
- This example fetches records in chunks of 500 from the database.
- Records are accumulated into a batch of 1000 before calling the processor.
- The batchProcessor is only triggered once the batch size is met (or at the end of the stream).
- Records are fetched in pages of 500.
- They are accumulated until the batch holds 1000 records, then the processor is called.
- A final, smaller batch with the remaining records is flushed at the end.

### Removed: fchunk(), fcursor(), fstream()
The Fiber-based `fchunk()`, `fcursor()` and `fstream()` methods have been removed.

PHP Fibers are cooperative and everything runs on one thread and one PDO connection, so these methods never processed anything in parallel, and the Fiber buffering gave no backpressure (a page is always fully loaded before it is processed). They also had data-loss bugs: `fstream()` dropped the first record of every chunk and yielded `null` after each chunk, and `fcursor()` skipped the first `$bufferSize` records. Replace them as follows:

| Removed | Use instead |
| --- | --- |
| `fchunk($size, $processor, $concurrency)` | `chunk($size, $processor)` |
| `fstream($size, $transform, $bufferSize)` | `stream($size, $transform)` |
| `fcursor($processor, $bufferSize)` | `cursor($processor)` |

#### Other behavior changes
- `chunk()`, `stream()`, `batch()` now throw `InvalidArgumentException` when the chunk size is below 1.
- `cursor()` no longer forces garbage collection after every row, which was a large per-row cost. PHP's own collector handles cycles.
- Database errors from `cursor()` keep the original `PDOException` as the exception's previous exception.

## Handling Multiple Database Connection
Doppar allows you to manage multiple database connections seamlessly. This is especially useful when working with different databases for various modules (e.g., separating reporting, user data, or third-party integrations).
Expand Down
Loading