chunk_execution
chunk_execution ¶
Async Execution Utilities
This module provides utilities for async function execution with features like: - Automatic retry logic with customizable error handling - Concurrency control with semaphores - Sequential execution of async functions - list chunking for batch processing
Functions:
| Name | Description |
|---|---|
run_with_retry |
Decorator for automatic retry logic on async functions |
gather_with_concurrency |
Execute multiple coroutines with concurrency limits |
run_sequence |
Execute async functions sequentially |
chunk_list |
Split a list into smaller chunks for batch processing |
Example
@run_with_retry(max_retry=3) async def fetch_data(): ... # Function that might fail ... return await api_call()
Limit concurrent operations¶
results = await gather_with_concurrency(*coroutines, n=10)
Process data in chunks¶
chunks = chunk_list(large_list, chunk_size=100)
chunk_list ¶
chunk_list(obj_ls: list, chunk_size: int)
Split a list into smaller chunks of specified size.
Divides a large list into smaller sublists for batch processing. The last chunk may be smaller than chunk_size if the list length is not evenly divisible by chunk_size.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
obj_ls
|
list[Any]
|
list of objects to split into chunks |
required |
chunk_size
|
int
|
Maximum number of items per chunk |
required |
Returns:
| Type | Description |
|---|---|
|
list[list[Any]]: list of chunks, where each chunk is a list of objects |
Raises:
| Type | Description |
|---|---|
ValueError
|
If chunk_size is less than 1 |
Example
data = [1, 2, 3, 4, 5, 6, 7, 8, 9] chunks = chunk_list(data, chunk_size=3) print(chunks) # [[1, 2, 3], [4, 5, 6], [7, 8, 9]]
chunks = chunk_list(data, chunk_size=4) print(chunks) # [[1, 2, 3, 4], [5, 6, 7, 8], [9]]
Note
This is useful for processing large datasets in smaller batches to manage memory usage or API rate limits.
Source code in src/crew_dcs/utils/chunk_execution.py
203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 | |
gather_with_concurrency
async
¶
gather_with_concurrency(*coros, n: int = 60)
Execute multiple coroutines with concurrency control.
Limits the number of concurrently running coroutines using a semaphore, preventing overwhelming of system resources or external APIs.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
*coros
|
Variable number of coroutines to execute |
()
|
|
n
|
int
|
Maximum number of concurrent coroutines (default: 60) |
60
|
Returns:
| Type | Description |
|---|---|
|
list[T]: Results from all coroutines in the same order as input |
Example
async def fetch_url(url): ... async with httpx.AsyncClient() as client: ... return await client.get(url)
urls = ["http://example.com/1", "http://example.com/2"] coroutines = [fetch_url(url) for url in urls] results = await gather_with_concurrency(*coroutines, n=5)
Note
This is particularly useful when making many API calls or I/O operations where you want to limit concurrent connections.
Source code in src/crew_dcs/utils/chunk_execution.py
131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 | |
run_sequence
async
¶
run_sequence(*functions)
Execute a sequence of async functions sequentially.
Executes each async function in order, waiting for each to complete before starting the next one. Useful when functions have dependencies or when you need to preserve execution order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
*functions
|
Variable number of awaitable functions to execute |
()
|
Returns:
| Type | Description |
|---|---|
|
list[Any]: Results from all functions in execution order |
Example
async def step1(): ... return "first"
async def step2(): ... return "second"
results = await run_sequence(step1(), step2()) print(results) # ["first", "second"]
Note
Unlike asyncio.gather(), this executes functions sequentially, not concurrently.
Source code in src/crew_dcs/utils/chunk_execution.py
170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | |
run_with_retry ¶
run_with_retry(
max_retry: int = 1,
errors_to_retry_tp: tuple[type, ...] | None = None,
)
Decorator that adds automatic retry logic to async functions.
This decorator will retry the decorated function if it raises specified exceptions, with special handling for httpx.ConnectTimeout errors (if httpx is available).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
max_retry
|
int
|
Maximum number of retry attempts (default: 1) |
1
|
errors_to_retry_tp
|
Tuple[type, ...]
|
Tuple of exception types to retry on. If None, retries on any Exception except ConnectTimeout. |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
Callable |
Decorated function with retry logic |
Example
@run_with_retry(max_retry=3, errors_to_retry_tp=(ConnectionError,)) async def fetch_data(): ... return await some_api_call()
Function will retry up to 3 times on ConnectionError¶
result = await fetch_data()
Note
- ConnectTimeout errors (if httpx available) include a 2-second delay between retries
- Other errors retry immediately
- All retry attempts are logged to stdout
Source code in src/crew_dcs/utils/chunk_execution.py
44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 | |