Python多进程处理大数据:Shared Memory Dictionary方案问询
Got it, let's work through how to solve this problem—handling that massive 86GB pickle dictionary (1.75 billion keys!) across 32+ processes while crunching a 2TB dataset is definitely a tough one, but there are proven approaches to make this feasible without draining your system's resources. The core issue here is avoiding loading the entire dictionary into every process's memory (32 processes × 86GB = 2.7TB of RAM—way too much for most systems). Here are your best options:
1. Use LMDB (Lightning Memory-Mapped Database)
I’ve used LMDB for exactly this kind of scenario: large read-only key-value stores that need shared access across multiple processes. It’s blazingly fast because it uses memory mapping under the hood—your entire dictionary lives in system memory once, and all processes can access it concurrently without duplicating data.
Step-by-Step Implementation:
First, convert your pickle dict to LMDB (one-time setup):
import lmdb import pickle # Load your original pickle dictionary with open("large_dict.pkl", "rb") as f: large_dict = pickle.load(f) # Initialize LMDB environment—set map_size to be slightly bigger than your dict (100GB here) env = lmdb.open("large_dict_lmdb", map_size=100 * 1024 * 1024 * 1024) # Write the dict to LMDB with env.begin(write=True) as txn: for key, value in large_dict.items(): # Serialize key/value—swap pickle for msgpack/pyarrow for faster speeds if needed txn.put(pickle.dumps(key), pickle.dumps(value)) env.close()
Then, in each worker process, query LMDB in read-only mode:
import lmdb import pickle def worker_process(task_chunk): # Open LMDB in read-only mode—lock=False skips unnecessary locking for read-only access env = lmdb.open("large_dict_lmdb", readonly=True, lock=False) with env.begin() as txn: for item in task_chunk: target_key = item["key"] # Adjust this to match your dataset's key format # Look up the value value_bytes = txn.get(pickle.dumps(target_key)) if value_bytes: value = pickle.loads(value_bytes) # Run your data processing logic here process_data_item(item, value) env.close()
2. Use a SQLite Database on a RAM Disk
SQLite supports concurrent read access out of the box, and if you host the database on a RAM disk (like /dev/shm on Linux), query speeds are almost as fast as in-memory structures. It’s a great option if you prefer standard library tools.
Step-by-Step Implementation:
First, convert your pickle dict to SQLite:
import sqlite3 import pickle conn = sqlite3.connect("large_dict.db") cursor = conn.cursor() # Create a table with blob columns for serialized keys/values cursor.execute("CREATE TABLE IF NOT EXISTS mappings (key BLOB PRIMARY KEY, value BLOB)") with open("large_dict.pkl", "rb") as f: large_dict = pickle.load(f) # Batch insert to speed up the conversion cursor.executemany( "INSERT INTO mappings VALUES (?, ?)", [(pickle.dumps(k), pickle.dumps(v)) for k, v in large_dict.items()] ) conn.commit() conn.close()
Move the database to a RAM disk (Linux example—adjust for your OS):
mv large_dict.db /dev/shm/
Then, in worker processes, query the read-only database:
import sqlite3 import pickle def worker_process(task_chunk): # Connect to the RAM-disk database in read-only mode conn = sqlite3.connect("file:/dev/shm/large_dict.db?mode=ro", uri=True) cursor = conn.cursor() for item in task_chunk: target_key = item["key"] cursor.execute( "SELECT value FROM mappings WHERE key = ?", (pickle.dumps(target_key),) ) result = cursor.fetchone() if result: value = pickle.loads(result[0]) process_data_item(item, value) conn.close()
3. Shared Memory Hash Table (Standard Library Only)
If you want to avoid third-party libraries, you can serialize your dictionary into shared memory arrays. This requires a bit more manual work, but it’s doable with Python’s multiprocessing.Array.
Core Idea:
- Serialize all keys and values into a single byte array stored in shared memory.
- Store offsets and lengths of each key/value pair in another shared array, so processes can calculate where to read data.
- For faster lookups, precompute a hash map of keys to their index in the offset array (store this hash map in shared memory too).
Example outline (simplified):
from multiprocessing import Pool, Array import pickle import struct def preprocess_dict(): with open("large_dict.pkl", "rb") as f: large_dict = pickle.load(f) serialized_data = b"" offsets = [] key_to_index = {} for idx, (key, value) in enumerate(large_dict.items()): k_bytes = pickle.dumps(key) v_bytes = pickle.dumps(value) offsets.append((len(serialized_data), len(k_bytes), len(v_bytes))) serialized_data += k_bytes + v_bytes key_to_index[key] = idx # Create shared arrays shared_data = Array('c', serialized_data) shared_offsets = Array('Q', [val for trio in offsets for val in trio]) # Flatten offsets # You'd also need to serialize key_to_index into shared memory here return shared_data, shared_offsets, key_to_index # Worker process would use these shared arrays to look up values by key index
Critical Performance Tips
- Swap pickle for faster serialization: Use
msgpackorpyarrowinstead of pickle to cut down on serialization/deserialization time. - Batch process your dataset: Split the 2TB dataset into chunks to minimize I/O overhead and keep workers busy.
- Tune process count: Don’t just set 32 processes—match it to your CPU core count (or slightly higher, but watch for I/O bottlenecks if your dataset is on disk).
- Avoid data copies: All the above approaches ensure the dictionary is loaded into memory once, not per process—this is non-negotiable for your scale.
内容的提问来源于stack exchange,提问作者Jon Deaton

