Dask DataFrame map_partitions分组键错误:KeyError 3.0 解决求助
Let’s break down what’s causing that frustrating KeyError 3.0 and rebuild your fuzzy matching workflow to be more robust—no groupby required.
Why You’re Getting the KeyError
Your current code uses groupby('phase').get_group(3.0) inside each partition. The problem? If even one partition in your large Dask DataFrame doesn’t contain any rows where phase == 3.0, get_group() will throw a KeyError. Since Dask runs partitions in parallel (and you’re seeing the error after hours), it’s likely hitting a partition without phase 3.0 data only after processing many others.
The Solution: Split DataFrames with Boolean Filtering
Instead of grouping in each partition, we’ll explicitly filter for the phase groups we need directly in the partition function. This avoids get_group() entirely and handles empty groups gracefully.
Here’s the revised approach:
Step 1: Define a Robust Partition Matching Function
First, create a function that operates on a single Dask partition, filters for your source (phase 3.0) and target (phase 1.5/2.0) groups, and runs the fuzzy match only if both groups exist:
import pandas as pd from fuzzywuzzy import fuzz def fuzzy_match_partition(df, target_phase, scorer, threshold): # Filter source (phase 3.0) and target groups directly source_drugs = df[df['phase'] == 3.0]['drugs'] target_drugs = df[df['phase'] == target_phase]['drugs'] # Skip empty groups to avoid errors if source_drugs.empty or target_drugs.empty: return pd.Series([], dtype='object', name='results') # Helper to match a single drug against the target list def match_drug(drug): # Get all target drugs that meet the threshold matches = target_drugs[target_drugs.apply(lambda x: scorer(drug, x) >= threshold)] return matches.tolist() # Apply matching to all source drugs return source_drugs.apply(match_drug)
Step 2: Run the Matching with Dask map_partitions
Now loop through your target phases, using the new function with map_partitions. We’ll also specify a clear meta parameter to help Dask handle the output correctly:
# Assuming `emb` is your original Dask DataFrame for target_phase in [2.0, 1.5]: fuzzy_results = emb.map_partitions( fuzzy_match_partition, target_phase=target_phase, scorer=fuzz.token_sort_ratio, threshold=85, meta=pd.Series(dtype='object', name='results') # Explicit meta for stability ).compute() # Add your result handling here (e.g., save to file, merge with other data) # print(f"Completed matching for phase {target_phase}")
Key Improvements Over Your Original Code
- No more KeyErrors: Boolean filtering returns empty Series instead of throwing errors when a phase is missing from a partition.
- Faster execution: Skipping
groupbyreduces overhead in each partition, especially for large datasets. - Explicit meta: Defining the output metadata helps Dask manage task scheduling and result aggregation more reliably.
Bonus: Optional Optimization for Large Target Groups
If your target phase groups are very large, consider precomputing token-sorted versions of the drug names to speed up the token_sort_ratio checks—this cuts down on redundant string processing during matching.
内容的提问来源于stack exchange,提问作者guru

