Python parallel processing library for building data pipelines. Supports thread/process/async execution with batch and streaming modes. Zero dependencies, type-safe, perfect for ETL and single-machine workloads.
Project description
FastPipe
Minimal parallel pipeline library for Python with thread/process/async execution, batch and streaming modes
Build parallel data processing pipelines with clean API. Zero dependencies, type-safe, perfect for ETL and single-machine parallel computing.
Perfect for: ETL pipelines • Data processing • API batch operations • Stream processing
Quick Start
import fastpipe
# Batch processing - collect all results
results = (fastpipe.create()
.map(download, workers=4)
.map(process, workers=4)
.run(data))
# Iterator mode - memory-efficient for large datasets
for result in fastpipe.create().map(process, workers=4).iter(huge_data):
save_to_db(result)
# Streaming mode - continuous put/get
stream = (fastpipe.create()
.map(process, workers=4)
.stream())
stream.put(item)
result = stream.get()
stream.close()
Installation
pip install fastpipe
Features
- Three execution modes:
thread,process,asyncwith automatic adapters - Three operation modes:
run()(batch),iter()(memory-efficient),stream()(continuous) - Built-in batching:
batch()/unbatch()for performance optimization - Backpressure control: Bounded queues prevent OOM
- Zero dependencies: stdlib only
- Thread-safe: No locks needed
API
Operations
pipe = fastpipe.create() \
.map(func, workers=4, mode=fastpipe.Mode.THREAD) \
.flat_map(func, workers=4) \
.filter(predicate, workers=2) \
.each(side_effect, workers=1) \
.batch(size=10) \
.unbatch()
Available operations:
map(func, workers=1)- Transform each item (1:1)flat_map(func, workers=1)- Transform and flatten (1:N)filter(predicate, workers=1)- Keep matching itemseach(func, workers=1)- Side effects, return originalbatch(size)- Group into fixed-size batchesunbatch()- Flatten batches back to items
Execution Modes
# Thread mode (I/O-bound)
.map(download, workers=10, mode=fastpipe.Mode.THREAD)
# Process mode (CPU-bound)
.map(compute, workers=4, mode=fastpipe.Mode.PROCESS)
# Async mode (high-concurrency I/O)
.map(async_fetch, workers=100, mode=fastpipe.Mode.ASYNC)
# Mix freely - automatic adapters inserted
Mode selection:
Mode.THREAD- I/O-bound, 1-20 workersMode.PROCESS- CPU-bound, 4-16 workers (bypasses GIL)Mode.ASYNC- I/O-bound, 50-1000+ workers (requires async functions)
Running Pipelines
| Method | Use Case | Memory | Example |
|---|---|---|---|
run(data) |
Small datasets | O(n) | pipe.run(data) |
iter(data) |
Large datasets | O(1) | for x in pipe.iter(data): save(x) |
stream() |
Continuous | O(queue) | stream.put(x); stream.get() |
Stream API:
stream = pipe.stream(queue_size=100)
stream.put(item) # Add item
result = stream.get() # Get result
stream.close(drain=True) # Shutdown
Callable Classes
Pass class types with init_args/init_kwargs for efficient multiprocess execution:
class ModelProcessor:
def __init__(self, model_path, device='cpu'):
# Each worker loads model independently
self.model = load_model(model_path)
self.device = device
def __call__(self, data):
return self.model.predict(data)
# ✅ Efficient: Only class definition is serialized
results = (fastpipe.create()
.map(
ModelProcessor,
workers=4,
mode='process',
init_args=('model.pth',),
init_kwargs={'device': 'cuda'}
)
.run(data))
# ❌ Inefficient: Entire model serialized to each worker
processor = ModelProcessor('model.pth')
results = pipe.map(processor, workers=4, mode='process').run(data)
Benefits:
- Avoids serializing large objects (models, databases)
- Per-worker resource management
- Independent worker state
- Works with non-picklable resources
Supported on all operations: map(), filter(), flat_map(), each()
Examples
Async High-Concurrency
import fastpipe
import aiohttp
async def fetch(url):
async with aiohttp.ClientSession() as s:
async with s.get(url) as r:
return await r.text()
results = (fastpipe.create()
.map(fetch, workers=100, mode=fastpipe.Mode.ASYNC)
.map(parse, workers=4, mode=fastpipe.Mode.THREAD)
.run(urls))
Mixed CPU/IO Pipeline
pipe = fastpipe.create() \
.map(download, workers=10, mode=fastpipe.Mode.THREAD) \
.batch(size=100) \
.map(process_batch, workers=4, mode=fastpipe.Mode.PROCESS) \
.unbatch() \
.map(upload, workers=5, mode=fastpipe.Mode.THREAD)
results = pipe.run(urls)
Kafka Stream Processing
pipe = fastpipe.create() \
.map(parse, workers=4) \
.filter(validate, workers=2) \
.map(transform, workers=4)
stream = pipe.stream(queue_size=500)
# Producer thread
for event in kafka_consumer:
stream.put(event)
# Consumer thread
while running:
result = stream.get()
database.save(result)
Error Handling
Fail-fast: First exception stops entire pipeline immediately.
try:
results = fastpipe.create().map(may_fail, workers=4).run(data)
except ValueError as e:
print(f"Pipeline failed: {e}")
Error tolerance: Wrap your functions.
def safe_process(x):
try:
return risky_operation(x)
except Exception as e:
logger.warning(f"Failed: {e}")
return None
results = fastpipe.create() \
.map(safe_process, workers=4) \
.filter(lambda x: x is not None) \
.run(data)
vs pypeln
| Feature | pypeln | fastpipe |
|---|---|---|
| Streaming mode | ❌ | ✅ |
| Async support | Separate module | Unified API |
| Mode switching | Change imports | Change parameter |
| Auto adapters | Manual | Automatic |
| Maintenance | Abandoned | Active |
Architecture
Thread-safety by design:
- Stateful ops (BATCH, ADAPTER): Single worker (no races)
- Stateless ops (MAP, FILTER): Multiple workers (no shared state)
- Queues: Thread-safe (stdlib)
Automatic adapters: Inserted when mode OR worker count changes between stages.
SENTINEL propagation: End-of-stream markers coordinate worker shutdown across stages.
Design Philosophy
Provide primitives, not policies.
FastPipe focuses on core parallel processing. You add retry (tenacity), metrics (prometheus), and timeouts as needed.
Scope:
- ✅ Single-machine parallel processing (10-1000 workers)
- ❌ Distributed computing → Use Ray, Dask, Spark
- ❌ Task queues → Use Celery, RQ
Principles:
- Simple over clever - No magic, explicit API
- Reliable over feature-rich - Thread-safe by design, fail-fast
- Composable over monolithic - Small API, integrates with ecosystem
- No black magic - No
PyThreadState_SetAsyncExchacks, no hidden state - Performance without compromise - Fast but correct
Example: No built-in retry. You choose your library:
from tenacity import retry, stop_after_attempt
@retry(stop=stop_after_attempt(3))
def func(x):
return risky_operation(x)
pipe.map(func, workers=4) # Your policy, your choice
FastPipe vs Ray
Different scope, complementary tools:
| Ray | FastPipe | |
|---|---|---|
| Scope | Multi-machine cluster | Single machine |
| Execution | Process-based tasks | Thread/process/async |
| Data passing | Always serialized (object store) | Zero-copy in thread/async |
| Use case | Distributed ML, TB-scale | ETL, data pipelines, GB-scale |
Key difference: Serialization overhead
Ray requires all data to pass through its object store (serialization/deserialization). FastPipe's thread/async modes avoid this entirely with zero-copy reference passing—crucial for large numpy arrays, models, or non-picklable objects.
Complementary usage:
import ray, fastpipe
@ray.remote
def process_shard(shard):
# Ray: distribute across cluster
# FastPipe: optimize local processing
return fastpipe.create().map(f, workers=10, mode='thread').run(shard)
When to use: Ray for multi-node coordination, FastPipe for single-machine efficiency.
License
MIT
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file fastpipe-0.0.1.tar.gz.
File metadata
- Download URL: fastpipe-0.0.1.tar.gz
- Upload date:
- Size: 41.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.11.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
d778a12d8ae915c719a286331edb8679ac04ea29a18278355c6c917445a0a132
|
|
| MD5 |
ea9b876d2a4d0239e40215de6908936c
|
|
| BLAKE2b-256 |
c4fd30652263743318b51a287ff048258ff9e8d4b3b6025345d0ffcd04ed6611
|
File details
Details for the file fastpipe-0.0.1-py3-none-any.whl.
File metadata
- Download URL: fastpipe-0.0.1-py3-none-any.whl
- Upload date:
- Size: 25.0 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.11.11
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
36d7ddbdce1afc2ab060f2666c24575ccd1be0c6caf77b9b9e14a581a0e2c96f
|
|
| MD5 |
f71658ce28e507cbaab7b4dc08f80220
|
|
| BLAKE2b-256 |
fb7f498ffe15e0d34adc63f418a09c4ffb51b460b2d210d2dfd2329778f52c42
|