You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于多DataFrame实现Dask map_partitions的Lambda函数?

Parallel Fuzzy Matching with Dask for Your DataFrames

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 numCores directly, try numCores * 2 for partitions—this can balance load better if some partitions have more rows than others.
  • Memory Management: If emb2 is extremely large, use Dask's broadcast mechanism to share it across workers instead of passing copies (though for most cases, the compute() approach works fine).

内容的提问来源于stack exchange,提问作者guru

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 06:40:45