update the doc structure
This commit is contained in:
@@ -0,0 +1,58 @@
|
||||
---
|
||||
layout: default
|
||||
title: "(Advanced) Async"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 5
|
||||
---
|
||||
|
||||
# (Advanced) Async
|
||||
|
||||
**Async** Nodes implement `prep_async()`, `exec_async()`, `exec_fallback_async()`, and/or `post_async()`. This is useful for:
|
||||
|
||||
1. **prep_async()**: For *fetching/reading data (files, APIs, DB)* in an I/O-friendly way.
|
||||
|
||||
2. **exec_async()**: Typically used for async LLM calls.
|
||||
|
||||
3. **post_async()**: For *awaiting user feedback*, *coordinating across multi-agents* or any additional async steps after `exec_async()`.
|
||||
|
||||
|
||||
**Note**: `AsyncNode` must be wrapped in `AsyncFlow`. `AsyncFlow` can also include regular (sync) nodes.
|
||||
|
||||
### Example
|
||||
|
||||
```python
|
||||
class SummarizeThenVerify(AsyncNode):
|
||||
async def prep_async(self, shared):
|
||||
# Example: read a file asynchronously
|
||||
doc_text = await read_file_async(shared["doc_path"])
|
||||
return doc_text
|
||||
|
||||
async def exec_async(self, prep_res):
|
||||
# Example: async LLM call
|
||||
summary = await call_llm_async(f"Summarize: {prep_res}")
|
||||
return summary
|
||||
|
||||
async def post_async(self, shared, prep_res, exec_res):
|
||||
# Example: wait for user feedback
|
||||
decision = await gather_user_feedback(exec_res)
|
||||
if decision == "approve":
|
||||
shared["summary"] = exec_res
|
||||
return "approve"
|
||||
return "deny"
|
||||
|
||||
summarize_node = SummarizeThenVerify()
|
||||
final_node = Finalize()
|
||||
|
||||
# Define transitions
|
||||
summarize_node - "approve" >> final_node
|
||||
summarize_node - "deny" >> summarize_node # retry
|
||||
|
||||
flow = AsyncFlow(start=summarize_node)
|
||||
|
||||
async def main():
|
||||
shared = {"doc_path": "document.txt"}
|
||||
await flow.run_async(shared)
|
||||
print("Final Summary:", shared.get("summary"))
|
||||
|
||||
asyncio.run(main())
|
||||
```
|
||||
@@ -0,0 +1,107 @@
|
||||
---
|
||||
layout: default
|
||||
title: "Batch"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 4
|
||||
---
|
||||
|
||||
# Batch
|
||||
|
||||
**Batch** makes it easier to handle large inputs in one Node or **rerun** a Flow multiple times. Example use cases:
|
||||
- **Chunk-based** processing (e.g., splitting large texts).
|
||||
- **Iterative** processing over lists of input items (e.g., user queries, files, URLs).
|
||||
|
||||
## 1. BatchNode
|
||||
|
||||
A **BatchNode** extends `Node` but changes `prep()` and `exec()`:
|
||||
|
||||
- **`prep(shared)`**: returns an **iterable** (e.g., list, generator).
|
||||
- **`exec(item)`**: called **once** per item in that iterable.
|
||||
- **`post(shared, prep_res, exec_res_list)`**: after all items are processed, receives a **list** of results (`exec_res_list`) and returns an **Action**.
|
||||
|
||||
|
||||
### Example: Summarize a Large File
|
||||
|
||||
```python
|
||||
class MapSummaries(BatchNode):
|
||||
def prep(self, shared):
|
||||
# Suppose we have a big file; chunk it
|
||||
content = shared["data"]
|
||||
chunk_size = 10000
|
||||
chunks = [content[i:i+chunk_size] for i in range(0, len(content), chunk_size)]
|
||||
return chunks
|
||||
|
||||
def exec(self, chunk):
|
||||
prompt = f"Summarize this chunk in 10 words: {chunk}"
|
||||
summary = call_llm(prompt)
|
||||
return summary
|
||||
|
||||
def post(self, shared, prep_res, exec_res_list):
|
||||
combined = "\n".join(exec_res_list)
|
||||
shared["summary"] = combined
|
||||
return "default"
|
||||
|
||||
map_summaries = MapSummaries()
|
||||
flow = Flow(start=map_summaries)
|
||||
flow.run(shared)
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 2. BatchFlow
|
||||
|
||||
A **BatchFlow** runs a **Flow** multiple times, each time with different `params`. Think of it as a loop that replays the Flow for each parameter set.
|
||||
|
||||
|
||||
### Example: Summarize Many Files
|
||||
|
||||
```python
|
||||
class SummarizeAllFiles(BatchFlow):
|
||||
def prep(self, shared):
|
||||
# Return a list of param dicts (one per file)
|
||||
filenames = list(shared["data"].keys()) # e.g., ["file1.txt", "file2.txt", ...]
|
||||
return [{"filename": fn} for fn in filenames]
|
||||
|
||||
# Suppose we have a per-file Flow (e.g., load_file >> summarize >> reduce):
|
||||
summarize_file = SummarizeFile(start=load_file)
|
||||
|
||||
# Wrap that flow into a BatchFlow:
|
||||
summarize_all_files = SummarizeAllFiles(start=summarize_file)
|
||||
summarize_all_files.run(shared)
|
||||
```
|
||||
|
||||
### Under the Hood
|
||||
1. `prep(shared)` returns a list of param dicts—e.g., `[{filename: "file1.txt"}, {filename: "file2.txt"}, ...]`.
|
||||
2. The **BatchFlow** loops through each dict. For each one:
|
||||
- It merges the dict with the BatchFlow’s own `params`.
|
||||
- It calls `flow.run(shared)` using the merged result.
|
||||
3. This means the sub-Flow is run **repeatedly**, once for every param dict.
|
||||
|
||||
---
|
||||
|
||||
## 3. Nested or Multi-Level Batches
|
||||
|
||||
You can nest a **BatchFlow** in another **BatchFlow**. For instance:
|
||||
- **Outer** batch: returns a list of diretory param dicts (e.g., `{"directory": "/pathA"}`, `{"directory": "/pathB"}`, ...).
|
||||
- **Inner** batch: returning a list of per-file param dicts.
|
||||
|
||||
At each level, **BatchFlow** merges its own param dict with the parent’s. By the time you reach the **innermost** node, the final `params` is the merged result of **all** parents in the chain. This way, a nested structure can keep track of the entire context (e.g., directory + file name) at once.
|
||||
|
||||
```python
|
||||
|
||||
class FileBatchFlow(BatchFlow):
|
||||
def prep(self, shared):
|
||||
directory = self.params["directory"]
|
||||
# e.g., files = ["file1.txt", "file2.txt", ...]
|
||||
files = [f for f in os.listdir(directory) if f.endswith(".txt")]
|
||||
return [{"filename": f} for f in files]
|
||||
|
||||
class DirectoryBatchFlow(BatchFlow):
|
||||
def prep(self, shared):
|
||||
directories = [ "/path/to/dirA", "/path/to/dirB"]
|
||||
return [{"directory": d} for d in directories]
|
||||
|
||||
# MapSummaries have params like {"directory": "/path/to/dirA", "filename": "file1.txt"}
|
||||
inner_flow = FileBatchFlow(start=MapSummaries())
|
||||
outer_flow = DirectoryBatchFlow(start=inner_flow)
|
||||
```
|
||||
@@ -0,0 +1,130 @@
|
||||
---
|
||||
layout: default
|
||||
title: "Communication"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 3
|
||||
---
|
||||
|
||||
# Communication
|
||||
|
||||
Nodes and Flows **communicate** in two ways:
|
||||
|
||||
1. **Shared Store (recommended)**
|
||||
|
||||
- A global data structure (often an in-mem dict) that all nodes can read and write by `prep()` and `post()`.
|
||||
- Great for data results, large content, or anything multiple nodes need.
|
||||
- You shall design the data structure and populate it ahead.
|
||||
|
||||
|
||||
2. **Params (only for [Batch](./batch.md))**
|
||||
- Each node has a local, ephemeral `params` dict passed in by the **parent Flow**, used as an identifier for tasks. Parameter keys and values shall be **immutable**.
|
||||
- Good for identifiers like filenames or numeric IDs, in Batch mode.
|
||||
|
||||
If you know memory management, think of the **Shared Store** like a **heap** (shared by all function calls), and **Params** like a **stack** (assigned by the caller).
|
||||
|
||||
> Use `Shared Store` for almost all cases. It's flexible and easy to manage. It separates *Data Schema* from *Compute Logic*, making the code easier to maintain. `Params` is more a syntax sugar for [Batch](./batch.md).
|
||||
{: .best-practice }
|
||||
|
||||
---
|
||||
|
||||
## 1. Shared Store
|
||||
|
||||
### Overview
|
||||
|
||||
A shared store is typically an in-mem dictionary, like:
|
||||
```python
|
||||
shared = {"data": {}, "summary": {}, "config": {...}, ...}
|
||||
```
|
||||
|
||||
It can also contain local file handlers, DB connections, or a combination for persistence. We recommend deciding the data structure or DB schema first based on your app requirements.
|
||||
|
||||
### Example
|
||||
|
||||
```python
|
||||
class LoadData(Node):
|
||||
def post(self, shared, prep_res, exec_res):
|
||||
# We write data to shared store
|
||||
shared["data"] = "Some text content"
|
||||
return None
|
||||
|
||||
class Summarize(Node):
|
||||
def prep(self, shared):
|
||||
# We read data from shared store
|
||||
return shared["data"]
|
||||
|
||||
def exec(self, prep_res):
|
||||
# Call LLM to summarize
|
||||
prompt = f"Summarize: {prep_res}"
|
||||
summary = call_llm(prompt)
|
||||
return summary
|
||||
|
||||
def post(self, shared, prep_res, exec_res):
|
||||
# We write summary to shared store
|
||||
shared["summary"] = exec_res
|
||||
return "default"
|
||||
|
||||
load_data = LoadData()
|
||||
summarize = Summarize()
|
||||
load_data >> summarize
|
||||
flow = Flow(start=load_data)
|
||||
|
||||
shared = {}
|
||||
flow.run(shared)
|
||||
```
|
||||
|
||||
Here:
|
||||
- `LoadData` writes to `shared["data"]`.
|
||||
- `Summarize` reads from `shared["data"]`, summarizes, and writes to `shared["summary"]`.
|
||||
|
||||
---
|
||||
|
||||
## 2. Params
|
||||
|
||||
**Params** let you store *per-Node* or *per-Flow* config that doesn't need to live in the shared store. They are:
|
||||
- **Immutable** during a Node's run cycle (i.e., they don't change mid-`prep->exec->post`).
|
||||
- **Set** via `set_params()`.
|
||||
- **Cleared** and updated each time a parent Flow calls it.
|
||||
|
||||
|
||||
> Only set the uppermost Flow params because others will be overwritten by the parent Flow.
|
||||
>
|
||||
> If you need to set child node params, see [Batch](./batch.md).
|
||||
{: .warning }
|
||||
|
||||
Typically, **Params** are identifiers (e.g., file name, page number). Use them to fetch the task you assigned or write to a specific part of the shared store.
|
||||
|
||||
### Example
|
||||
|
||||
```python
|
||||
# 1) Create a Node that uses params
|
||||
class SummarizeFile(Node):
|
||||
def prep(self, shared):
|
||||
# Access the node's param
|
||||
filename = self.params["filename"]
|
||||
return shared["data"].get(filename, "")
|
||||
|
||||
def exec(self, prep_res):
|
||||
prompt = f"Summarize: {prep_res}"
|
||||
return call_llm(prompt)
|
||||
|
||||
def post(self, shared, prep_res, exec_res):
|
||||
filename = self.params["filename"]
|
||||
shared["summary"][filename] = exec_res
|
||||
return "default"
|
||||
|
||||
# 2) Set params
|
||||
node = SummarizeFile()
|
||||
|
||||
# 3) Set Node params directly (for testing)
|
||||
node.set_params({"filename": "doc1.txt"})
|
||||
node.run(shared)
|
||||
|
||||
# 4) Create Flow
|
||||
flow = Flow(start=node)
|
||||
|
||||
# 5) Set Flow params (overwrites node params)
|
||||
flow.set_params({"filename": "doc2.txt"})
|
||||
flow.run(shared) # The node summarizes doc2, not doc1
|
||||
```
|
||||
|
||||
---
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
layout: default
|
||||
title: "Core Abstraction"
|
||||
nav_order: 2
|
||||
has_children: true
|
||||
---
|
||||
@@ -0,0 +1,181 @@
|
||||
---
|
||||
layout: default
|
||||
title: "Flow"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 2
|
||||
---
|
||||
|
||||
# Flow
|
||||
|
||||
A **Flow** orchestrates a graph of Nodes. You can chain Nodes in a sequence or create branching depending on the **Actions** returned from each Node's `post()`.
|
||||
|
||||
## 1. Action-based Transitions
|
||||
|
||||
Each Node's `post()` returns an **Action** string. By default, if `post()` doesn't return anything, we treat that as `"default"`.
|
||||
|
||||
You define transitions with the syntax:
|
||||
|
||||
1. **Basic default transition**: `node_a >> node_b`
|
||||
This means if `node_a.post()` returns `"default"`, go to `node_b`.
|
||||
(Equivalent to `node_a - "default" >> node_b`)
|
||||
|
||||
2. **Named action transition**: `node_a - "action_name" >> node_b`
|
||||
This means if `node_a.post()` returns `"action_name"`, go to `node_b`.
|
||||
|
||||
It's possible to create loops, branching, or multi-step flows.
|
||||
|
||||
## 2. Creating a Flow
|
||||
|
||||
A **Flow** begins with a **start** node. You call `Flow(start=some_node)` to specify the entry point. When you call `flow.run(shared)`, it executes the start node, looks at its returned Action from `post()`, follows the transition, and continues until there's no next node.
|
||||
|
||||
### Example: Simple Sequence
|
||||
|
||||
Here's a minimal flow of two nodes in a chain:
|
||||
|
||||
```python
|
||||
node_a >> node_b
|
||||
flow = Flow(start=node_a)
|
||||
flow.run(shared)
|
||||
```
|
||||
|
||||
- When you run the flow, it executes `node_a`.
|
||||
- Suppose `node_a.post()` returns `"default"`.
|
||||
- The flow then sees `"default"` Action is linked to `node_b` and runs `node_b`.
|
||||
- `node_b.post()` returns `"default"` but we didn't define `node_b >> something_else`. So the flow ends there.
|
||||
|
||||
### Example: Branching & Looping
|
||||
|
||||
Here's a simple expense approval flow that demonstrates branching and looping. The `ReviewExpense` node can return three possible Actions:
|
||||
|
||||
- `"approved"`: expense is approved, move to payment processing
|
||||
- `"needs_revision"`: expense needs changes, send back for revision
|
||||
- `"rejected"`: expense is denied, finish the process
|
||||
|
||||
We can wire them like this:
|
||||
|
||||
```python
|
||||
# Define the flow connections
|
||||
review - "approved" >> payment # If approved, process payment
|
||||
review - "needs_revision" >> revise # If needs changes, go to revision
|
||||
review - "rejected" >> finish # If rejected, finish the process
|
||||
|
||||
revise >> review # After revision, go back for another review
|
||||
payment >> finish # After payment, finish the process
|
||||
|
||||
flow = Flow(start=review)
|
||||
```
|
||||
|
||||
Let's see how it flows:
|
||||
|
||||
1. If `review.post()` returns `"approved"`, the expense moves to the `payment` node
|
||||
2. If `review.post()` returns `"needs_revision"`, it goes to the `revise` node, which then loops back to `review`
|
||||
3. If `review.post()` returns `"rejected"`, it moves to the `finish` node and stops
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
review[Review Expense] -->|approved| payment[Process Payment]
|
||||
review -->|needs_revision| revise[Revise Report]
|
||||
review -->|rejected| finish[Finish Process]
|
||||
|
||||
revise --> review
|
||||
payment --> finish
|
||||
```
|
||||
|
||||
### Running Individual Nodes vs. Running a Flow
|
||||
|
||||
- `node.run(shared)`: Just runs that node alone (calls `prep->exec->post()`), returns an Action.
|
||||
- `flow.run(shared)`: Executes from the start node, follows Actions to the next node, and so on until the flow can't continue.
|
||||
|
||||
|
||||
> `node.run(shared)` **does not** proceed to the successor.
|
||||
> This is mainly for debugging or testing a single node.
|
||||
>
|
||||
> Always use `flow.run(...)` in production to ensure the full pipeline runs correctly.
|
||||
{: .warning }
|
||||
|
||||
## 3. Nested Flows
|
||||
|
||||
A **Flow** can act like a Node, which enables powerful composition patterns. This means you can:
|
||||
|
||||
1. Use a Flow as a Node within another Flow's transitions.
|
||||
2. Combine multiple smaller Flows into a larger Flow for reuse.
|
||||
3. Node `params` will be a merging of **all** parents' `params`.
|
||||
|
||||
### Flow's Node Methods
|
||||
|
||||
A **Flow** is also a **Node**, so it will run `prep()` and `post()`. However:
|
||||
|
||||
- It **won't** run `exec()`, as its main logic is to orchestrate its nodes.
|
||||
- `post()` always receives `None` for `exec_res` and should instead get the flow execution results from the shared store.
|
||||
|
||||
|
||||
### Basic Flow Nesting
|
||||
|
||||
Here's how to connect a flow to another node:
|
||||
|
||||
```python
|
||||
# Create a sub-flow
|
||||
node_a >> node_b
|
||||
subflow = Flow(start=node_a)
|
||||
|
||||
# Connect it to another node
|
||||
subflow >> node_c
|
||||
|
||||
# Create the parent flow
|
||||
parent_flow = Flow(start=subflow)
|
||||
```
|
||||
|
||||
When `parent_flow.run()` executes:
|
||||
1. It starts `subflow`
|
||||
2. `subflow` runs through its nodes (`node_a->node_b`)
|
||||
3. After `subflow` completes, execution continues to `node_c`
|
||||
|
||||
### Example: Order Processing Pipeline
|
||||
|
||||
Here's a practical example that breaks down order processing into nested flows:
|
||||
|
||||
```python
|
||||
# Payment processing sub-flow
|
||||
validate_payment >> process_payment >> payment_confirmation
|
||||
payment_flow = Flow(start=validate_payment)
|
||||
|
||||
# Inventory sub-flow
|
||||
check_stock >> reserve_items >> update_inventory
|
||||
inventory_flow = Flow(start=check_stock)
|
||||
|
||||
# Shipping sub-flow
|
||||
create_label >> assign_carrier >> schedule_pickup
|
||||
shipping_flow = Flow(start=create_label)
|
||||
|
||||
# Connect the flows into a main order pipeline
|
||||
payment_flow >> inventory_flow >> shipping_flow
|
||||
|
||||
# Create the master flow
|
||||
order_pipeline = Flow(start=payment_flow)
|
||||
|
||||
# Run the entire pipeline
|
||||
order_pipeline.run(shared_data)
|
||||
```
|
||||
|
||||
This creates a clean separation of concerns while maintaining a clear execution path:
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
subgraph order_pipeline[Order Pipeline]
|
||||
subgraph paymentFlow["Payment Flow"]
|
||||
A[Validate Payment] --> B[Process Payment] --> C[Payment Confirmation]
|
||||
end
|
||||
|
||||
subgraph inventoryFlow["Inventory Flow"]
|
||||
D[Check Stock] --> E[Reserve Items] --> F[Update Inventory]
|
||||
end
|
||||
|
||||
subgraph shippingFlow["Shipping Flow"]
|
||||
G[Create Label] --> H[Assign Carrier] --> I[Schedule Pickup]
|
||||
end
|
||||
|
||||
paymentFlow --> inventoryFlow
|
||||
inventoryFlow --> shippingFlow
|
||||
end
|
||||
```
|
||||
|
||||
@@ -0,0 +1,110 @@
|
||||
---
|
||||
layout: default
|
||||
title: "Node"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 1
|
||||
---
|
||||
|
||||
# Node
|
||||
|
||||
A **Node** is the smallest building block. Each Node has 3 steps `prep->exec->post`:
|
||||
|
||||
<div align="center">
|
||||
<img src="https://github.com/the-pocket/PocketFlow/raw/main/assets/node.png?raw=true" width="400"/>
|
||||
</div>
|
||||
|
||||
|
||||
1. `prep(shared)`
|
||||
- **Read and preprocess data** from `shared` store.
|
||||
- Examples: *query DB, read files, or serialize data into a string*.
|
||||
- Return `prep_res`, which is used by `exec()` and `post()`.
|
||||
|
||||
2. `exec(prep_res)`
|
||||
- **Execute compute logic**, with optional retries and error handling (below).
|
||||
- Examples: *(mostly) LLM calls, remote APIs, tool use*.
|
||||
- ⚠️ This shall be only for compute and **NOT** access `shared`.
|
||||
- ⚠️ If retries enabled, ensure idempotent implementation.
|
||||
- Return `exec_res`, which is passed to `post()`.
|
||||
|
||||
3. `post(shared, prep_res, exec_res)`
|
||||
- **Postprocess and write data** back to `shared`.
|
||||
- Examples: *update DB, change states, log results*.
|
||||
- **Decide the next action** by returning a *string* (`action = "default"` if *None*).
|
||||
|
||||
|
||||
|
||||
> **Why 3 steps?** To enforce the principle of *separation of concerns*. The data storage and data processing are operated separately.
|
||||
>
|
||||
> All steps are *optional*. E.g., you can only implement `prep` and `post` if you just need to process data.
|
||||
{: .note }
|
||||
|
||||
|
||||
### Fault Tolerance & Retries
|
||||
|
||||
You can **retry** `exec()` if it raises an exception via two parameters when define the Node:
|
||||
|
||||
- `max_retries` (int): Max times to run `exec()`. The default is `1` (**no** retry).
|
||||
- `wait` (int): The time to wait (in **seconds**) before next retry. By default, `wait=0` (no waiting).
|
||||
`wait` is helpful when you encounter rate-limits or quota errors from your LLM provider and need to back off.
|
||||
|
||||
```python
|
||||
my_node = SummarizeFile(max_retries=3, wait=10)
|
||||
```
|
||||
|
||||
When an exception occurs in `exec()`, the Node automatically retries until:
|
||||
|
||||
- It either succeeds, or
|
||||
- The Node has retried `max_retries - 1` times already and fails on the last attempt.
|
||||
|
||||
You can get the current retry times (0-based) from `self.cur_retry`.
|
||||
|
||||
```python
|
||||
class RetryNode(Node):
|
||||
def exec(self, prep_res):
|
||||
print(f"Retry {self.cur_retry} times")
|
||||
raise Exception("Failed")
|
||||
```
|
||||
|
||||
### Graceful Fallback
|
||||
|
||||
To **gracefully handle** the exception (after all retries) rather than raising it, override:
|
||||
|
||||
```python
|
||||
def exec_fallback(self, shared, prep_res, exc):
|
||||
raise exc
|
||||
```
|
||||
|
||||
By default, it just re-raises exception. But you can return a fallback result instead, which becomes the `exec_res` passed to `post()`.
|
||||
|
||||
### Example: Summarize file
|
||||
|
||||
```python
|
||||
class SummarizeFile(Node):
|
||||
def prep(self, shared):
|
||||
return shared["data"]
|
||||
|
||||
def exec(self, prep_res):
|
||||
if not prep_res:
|
||||
return "Empty file content"
|
||||
prompt = f"Summarize this text in 10 words: {prep_res}"
|
||||
summary = call_llm(prompt) # might fail
|
||||
return summary
|
||||
|
||||
def exec_fallback(self, shared, prep_res, exc):
|
||||
# Provide a simple fallback instead of crashing
|
||||
return "There was an error processing your request."
|
||||
|
||||
def post(self, shared, prep_res, exec_res):
|
||||
shared["summary"] = exec_res
|
||||
# Return "default" by not returning
|
||||
|
||||
summarize_node = SummarizeFile(max_retries=3)
|
||||
|
||||
# node.run() calls prep->exec->post
|
||||
# If exec() fails, it retries up to 3 times before calling exec_fallback()
|
||||
action_result = summarize_node.run(shared)
|
||||
|
||||
print("Action returned:", action_result) # "default"
|
||||
print("Summary stored:", shared["summary"])
|
||||
```
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
---
|
||||
layout: default
|
||||
title: "(Advanced) Parallel"
|
||||
parent: "Core Abstraction"
|
||||
nav_order: 6
|
||||
---
|
||||
|
||||
# (Advanced) Parallel
|
||||
|
||||
**Parallel** Nodes and Flows let you run multiple **Async** Nodes and Flows **concurrently**—for example, summarizing multiple texts at once. This can improve performance by overlapping I/O and compute.
|
||||
|
||||
> Because of Python’s GIL, parallel nodes and flows can’t truly parallelize CPU-bound tasks (e.g., heavy numerical computations). However, they excel at overlapping I/O-bound work—like LLM calls, database queries, API requests, or file I/O.
|
||||
{: .warning }
|
||||
|
||||
> - **Ensure Tasks Are Independent**: If each item depends on the output of a previous item, **do not** parallelize.
|
||||
>
|
||||
> - **Beware of Rate Limits**: Parallel calls can **quickly** trigger rate limits on LLM services. You may need a **throttling** mechanism (e.g., semaphores or sleep intervals).
|
||||
>
|
||||
> - **Consider Single-Node Batch APIs**: Some LLMs offer a **batch inference** API where you can send multiple prompts in a single call. This is more complex to implement but can be more efficient than launching many parallel requests and mitigates rate limits.
|
||||
{: .best-practice }
|
||||
|
||||
|
||||
## AsyncParallelBatchNode
|
||||
|
||||
Like **AsyncBatchNode**, but run `exec_async()` in **parallel**:
|
||||
|
||||
```python
|
||||
class ParallelSummaries(AsyncParallelBatchNode):
|
||||
async def prep_async(self, shared):
|
||||
# e.g., multiple texts
|
||||
return shared["texts"]
|
||||
|
||||
async def exec_async(self, text):
|
||||
prompt = f"Summarize: {text}"
|
||||
return await call_llm_async(prompt)
|
||||
|
||||
async def post_async(self, shared, prep_res, exec_res_list):
|
||||
shared["summary"] = "\n\n".join(exec_res_list)
|
||||
return "default"
|
||||
|
||||
node = ParallelSummaries()
|
||||
flow = AsyncFlow(start=node)
|
||||
```
|
||||
|
||||
## AsyncParallelBatchFlow
|
||||
|
||||
Parallel version of **BatchFlow**. Each iteration of the sub-flow runs **concurrently** using different parameters:
|
||||
|
||||
```python
|
||||
class SummarizeMultipleFiles(AsyncParallelBatchFlow):
|
||||
async def prep_async(self, shared):
|
||||
return [{"filename": f} for f in shared["files"]]
|
||||
|
||||
sub_flow = AsyncFlow(start=LoadAndSummarizeFile())
|
||||
parallel_flow = SummarizeMultipleFiles(start=sub_flow)
|
||||
await parallel_flow.run_async(shared)
|
||||
```
|
||||
Reference in New Issue
Block a user