如何基于多DataFrame实现Dask map_partitions的Lambda函数?
Got it, let's work through how to parallelize your fuzzy matching with Dask. The main trick here is making sure each partition of your emb1 Dask DataFrame can access the full set of strings from emb2—since every row in emb1 needs to compare against all entries in emb2, not just a subset partition. Here's a step-by-step solution:
Step 1: Prepare the Full emb2 String List
First, we'll extract the complete set of strings from emb2 as a Pandas Series. Since every partition of emb1 needs to compare against all of emb2, we can safely compute this once (it'll be passed to each worker process):
# Get the full strings list from emb2 (compute to Pandas since we need all of it) emb2_strings = emb2['strings'].compute()
Step 2: Write a Partition-Level Matching Function
We'll create a function that handles an entire Pandas DataFrame partition (from emb1) and runs the fuzzy matching against our prepped emb2_strings list. This replaces your row-wise apply with partition-wise processing:
from fuzzywuzzy import process, fuzz def fuzzy_match_partition(df_partition, choices, scorer, cutoff): # Apply fuzzy matching to every row in the partition matches = df_partition['strings'].apply( lambda x: process.extractOne(x, choices=choices, scorer=scorer, score_cutoff=cutoff) ) # Return a DataFrame with results (preserve the original index to align data) return pd.DataFrame({'fuzzy_result': matches}, index=df_partition.index)
Step 3: Run Parallel Processing with map_partitions
Use Dask's map_partitions to apply our function across every partition of emb1. We need to specify meta so Dask knows the structure of the output DataFrame:
import dask.dataframe as dd # Run the parallel fuzzy match dask_fuzzy_results = emb1.map_partitions( fuzzy_match_partition, choices=emb2_strings, scorer=fuzz.ratio, cutoff=90, # Define metadata for the output DataFrame (critical for Dask to work correctly) meta=pd.DataFrame({'fuzzy_result': pd.Series(dtype='object')}) ) # Convert Dask results to a Pandas DataFrame when ready to compute final_results = dask_fuzzy_results.compute()
Key Optimizations to Speed Things Up More
- Swap to RapidFuzz: FuzzyWuzzy is slow for large datasets—replace it with
rapidfuzz(a faster, drop-in replacement written in C). It supports batch operations that will drastically cut down per-partition time:from rapidfuzz import process, fuzz def fuzzy_match_partition(df_partition, choices, scorer, cutoff): # Batch processing with RapidFuzz (much faster than row-wise apply) matches = process.extractOne( df_partition['strings'].tolist(), choices=choices.tolist(), scorer=scorer, score_cutoff=cutoff ) return pd.DataFrame({'fuzzy_result': matches}, index=df_partition.index) - Tune Partition Count: Instead of using
numCoresdirectly, trynumCores * 2for partitions—this can balance load better if some partitions have more rows than others. - Memory Management: If
emb2is extremely large, use Dask'sbroadcastmechanism to share it across workers instead of passing copies (though for most cases, thecompute()approach works fine).
内容的提问来源于stack exchange,提问作者guru

