A lightning-fast, zero-copy, cross-process data store for Python using Apache Arrow and shared memory.
Project description
๐น ArrowShelf
High-Performance Shared Memory Data Exchange for Python
ArrowShelf is a cutting-edge Python library that enables lightning-fast shared memory data exchange between processes using Apache Arrow's columnar format. Perfect for high-performance computing, machine learning pipelines, and distributed data processing.
๐ Key Features
- ๐ Zero-Copy Operations: Direct memory access without serialization overhead
- ๐ง Process-Safe: Thread and multiprocess safe data sharing
- ๐ Columnar Efficiency: Optimized for analytical workloads with Apache Arrow
- ๐ฏ FAISS Integration: Built-in support for approximate nearest neighbor search
- ๐ Automatic Cleanup: Smart memory management with reference counting
- ๐ก๏ธ Production Ready: Robust error handling and connection management
๐ฆ Installation
pip install arrowshelf
For FAISS integration (optional):
pip install faiss-cpu # or faiss-gpu for GPU support
๐ Starting the ArrowShelf Server
ArrowShelf requires a server daemon to manage shared memory. Start it before running your applications:
# Start the ArrowShelf server
arrowshelf-server
# Or run in background (Linux/Mac)
arrowshelf-server &
# Windows background (using PowerShell)
Start-Process arrowshelf-server -WindowStyle Hidden
The server will run on localhost:50051 by default.
๐ Quick Start
Setting Up ArrowShelf
- Start the server (required):
arrowshelf-server
- Basic Usage (in a separate terminal/process):
import arrowshelf
import pandas as pd
import numpy as np
# Create sample data
df = pd.DataFrame({
'x': np.random.rand(10000),
'y': np.random.rand(10000),
'z': np.random.rand(10000)
})
# Store in shared memory
key = arrowshelf.put(df)
# Access from any process
retrieved_df = arrowshelf.get(key)
print(f"Retrieved {len(retrieved_df)} rows")
# Cleanup
arrowshelf.delete(key)
Advanced Zero-Copy Access
import arrowshelf
import numpy as np
# Store data
key = arrowshelf.put(df)
# Get Arrow table for zero-copy operations
table = arrowshelf.get_arrow(key)
x_column = table.column("x").chunk(0).to_numpy(zero_copy_only=True)
# Direct NumPy operations without copying
result = np.mean(x_column)
๐ฏ Real-World Example: Parallel Nearest Neighbor Search
This example demonstrates how ArrowShelf enables efficient parallel processing with FAISS for approximate nearest neighbor search.
Prerequisites
- Start ArrowShelf server:
arrowshelf-server
- Install dependencies:
pip install arrowshelf faiss-cpu pandas numpy
Complete Example
import multiprocessing as mp
import threading
import pandas as pd
import numpy as np
import time
import arrowshelf
import faiss
from multiprocessing.pool import ThreadPool
# Enable thread-based multiprocessing
threading.Pool = ThreadPool
def worker_faiss_search(task_data):
"""Worker function for parallel FAISS nearest neighbor search"""
key, start_index, end_index = task_data
# Zero-copy access to shared data
table = arrowshelf.get_arrow(key).combine_chunks()
x = table.column("x").chunk(0).to_numpy(zero_copy_only=True)
y = table.column("y").chunk(0).to_numpy(zero_copy_only=True)
z = table.column("z").chunk(0).to_numpy(zero_copy_only=True)
# Stack coordinates for FAISS
all_points = np.stack([x, y, z], axis=1).astype(np.float32)
query_chunk = all_points[start_index:end_index]
# Configure FAISS IVF index
d = 3 # 3D points
nlist = 100 # Voronoi cells
quantizer = faiss.IndexFlatL2(d)
index = faiss.IndexIVFFlat(quantizer, d, nlist, faiss.METRIC_L2)
# Train and populate index
index.train(all_points)
index.add(all_points)
index.nprobe = 10 # Search cells
# Perform approximate k-NN search
_, distances = index.search(query_chunk, 11) # k=11 (excluding self)
avg_distance = np.mean(np.sqrt(distances[:, 1:])) # Exclude self-distance
return avg_distance
def parallel_nearest_neighbor_demo():
"""Demonstrate parallel processing with ArrowShelf + FAISS"""
# Check ArrowShelf connection
try:
arrowshelf.list_keys()
print("โ
ArrowShelf server connection OK")
except arrowshelf.ConnectionError:
print("โ ERROR: ArrowShelf server not running!")
print("Please start the server first: arrowshelf-server")
return
# Generate sample 3D points
num_points = 100_000
num_cores = 6
print(f"๐ Running parallel k-NN search on {num_points:,} 3D points")
# Create dataset
df = pd.DataFrame(np.random.rand(num_points, 3), columns=['x', 'y', 'z'])
print(f"๐ Dataset size: {df.memory_usage(deep=True).sum() / 1024**2:.2f} MB")
# Store in ArrowShelf
key = arrowshelf.put(df)
# Create tasks for parallel processing
chunk_size = num_points // num_cores
tasks = [
(key, i * chunk_size, (i + 1) * chunk_size)
for i in range(num_cores)
]
# Execute parallel search
print(f"โก Processing with {num_cores} cores...")
start_time = time.perf_counter()
with ThreadPool(processes=num_cores) as pool:
results = pool.map(worker_faiss_search, tasks)
duration = time.perf_counter() - start_time
avg_distance = np.mean(results)
# Results
print(f"โ
Average 10-NN distance: {avg_distance:.6f}")
print(f"๐ Processing time: {duration:.4f} seconds")
print(f"๐ฅ Throughput: {num_points/duration:,.0f} points/second")
# Cleanup
arrowshelf.delete(key)
print("๐งน Cleanup completed")
if __name__ == "__main__":
mp.set_start_method('spawn', force=True)
parallel_nearest_neighbor_demo()
Running the Example
- Terminal 1 - Start the server:
arrowshelf-server
- Terminal 2 - Run the example:
python nearest_neighbor_demo.py
Expected Output:
โ
ArrowShelf server connection OK
๐ Running parallel k-NN search on 100,000 3D points
๐ Dataset size: 2.29 MB
โก Processing with 6 cores...
โ
Average 10-NN distance: 210.789151
๐ Processing time: 1.0017 seconds
๐ฅ Throughput: 99,830 points/second
๐งน Cleanup completed
๐ Project Evolution
ArrowShelf has evolved from a simple data sharing concept to a high-performance computing powerhouse. Here's the journey of optimization:
Performance Evolution Timeline
| Benchmark | Architecture | Algorithm | Time | Improvement |
|---|---|---|---|---|
| Pickle + Brute Force | Slow Data Transfer | Brute Force O(nยฒ) | 16.7s | Baseline |
| ArrowShelf + Brute Force | Fast Data Transfer | Brute Force O(nยฒ) | 14.5s | 13% faster |
| ArrowShelf + FAISS IndexFlatL2 | Fast Data Transfer | Optimized Exact Search | 1.84s | 87% faster |
| ArrowShelf + FAISS IndexIVFFlat | Fast Data Transfer | Approximate Search | 1.00s | 94% faster |
๐ Performance Visualization
Traditional Pickle Approach โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ 16.7s
ArrowShelf + Brute Force โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ 14.5s
ArrowShelf + FAISS Exact โโโโ 1.84s
ArrowShelf + FAISS Approx โโ 1.00s โก
๐ฏ Key Milestones
- ๐๏ธ Phase 1: Foundation - Basic shared memory with Apache Arrow
- โก Phase 2: Optimization - Zero-copy operations and efficient data transfer
- ๐ Phase 3: Intelligence - FAISS integration for similarity search
- ๐ Phase 4: Approximation - IVF indexing for ultimate performance
The evolution demonstrates a 16.7x performance improvement from traditional pickle-based approaches to our current FAISS-optimized implementation.
๐ Performance
ArrowShelf delivers exceptional performance for data-intensive applications:
FAISS Integration Benchmark
- Dataset: 100,000 3D points (2.29 MB)
- Operation: Approximate 10-NN search with IVF index
- Hardware: 6-core parallel processing
- Result: 1.00 seconds processing time
- Throughput: ~100K points/second
Key Performance Benefits
- Zero-Copy Access: Direct memory mapping eliminates serialization overhead
- Columnar Storage: Optimized for analytical operations and vectorized computations
- Parallel Processing: Efficient multi-core scaling with shared memory
- Memory Efficiency: Reference counting prevents memory leaks
๐ง API Reference
Core Functions
# Store data in shared memory
key = arrowshelf.put(data)
# Retrieve data as pandas DataFrame
df = arrowshelf.get(key)
# Retrieve data as Arrow Table (zero-copy)
table = arrowshelf.get_arrow(key)
# List all stored keys
keys = arrowshelf.list_keys()
# Delete data from shared memory
arrowshelf.delete(key)
# Close connection
arrowshelf.close()
Advanced Operations
# Batch operations
arrowshelf.delete_all() # Clear all data
# Connection management
arrowshelf.is_connected() # Check connection status
# Memory statistics
arrowshelf.memory_usage() # Get usage statistics
๐ ๏ธ Use Cases
๐ค Machine Learning
- Feature Engineering: Share preprocessed datasets across training processes
- Model Serving: Cache model predictions and intermediate results
- Hyperparameter Tuning: Efficient data sharing in parallel optimization
๐ Data Analytics
- ETL Pipelines: Zero-copy data transformations
- Distributed Computing: Shared memory for map-reduce operations
- Real-time Analytics: High-throughput data processing
๐ฌ Scientific Computing
- Numerical Simulations: Share large arrays between simulation processes
- Image Processing: Efficient pixel data sharing
- Geospatial Analysis: Fast coordinate and geometry operations
๐๏ธ Architecture
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
โ Process A โ โ ArrowShelf โ โ Process B โ
โ โ โ Server โ โ โ
โ put(data) โโโโโโโโโถโ โโโโโโโโโโ get(key) โ
โ โ โ Apache Arrow โ โ โ
โ โ โ Shared Memory โ โ โ
โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโ
๐ค Contributing
We welcome contributions! Please see our Contributing Guide for details.
- Fork the repository
- Create your feature branch (
git checkout -b feature/amazing-feature) - Commit your changes (
git commit -m 'Add amazing feature') - Push to the branch (
git push origin feature/amazing-feature) - Open a Pull Request
๐ Requirements
- Python 3.7+
- Apache Arrow
- pandas
- numpy
Optional dependencies:
- FAISS (for nearest neighbor search)
- multiprocessing support
๐ License
This project is licensed under the MIT License - see the LICENSE file for details.
๐ Acknowledgments
- Built on Apache Arrow columnar memory format
- Optimized for FAISS similarity search
- Inspired by modern high-performance computing needs
๐ Support
- ๐ง Issues: GitHub Issues
- ๐ฌ Discussions: GitHub Discussions
- ๐ Documentation: Wiki
โญ Star this repository if ArrowShelf helps accelerate your data processing workflows!
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 Distributions
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 arrowshelf-2.3.1-py3-none-win_amd64.whl.
File metadata
- Download URL: arrowshelf-2.3.1-py3-none-win_amd64.whl
- Upload date:
- Size: 832.5 kB
- Tags: Python 3, Windows x86-64
- Uploaded using Trusted Publishing? No
- Uploaded via: maturin/1.8.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3a41834069bae8094c98071dde0af4a2268116d8953522af7b88a3dfe479330e
|
|
| MD5 |
790f09a406fd29087dc21de90e0c5830
|
|
| BLAKE2b-256 |
f9d8aa4f99b2c8dab003266a453f92a92233b4ffc157126f1041bf10b73d90f5
|