PySpark apply function on 2 dataframes and write to csv for billions of rows on small hardware

Viewed 386

I am trying to apply a levenshtein function for each string in dfs against each string in dfc and write the resulting dataframe to csv. The issue is that I'm creating so many rows by using the cross join and then applying the function, that my machine is struggling to write anything (taking forever to execute).

Trying to improve write performance:

  • I'm filtering out a few things on the result of the cross join i.e. rows where the LevenshteinDistance is less than 15% of the target word's.
  • Using bucketing on the first letter of each target word i.e. a, b, c, etc. still no luck (i.e. job runs for hours and doesn't generate any results).
from datetime import datetime
from config import config
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
from pyspark.sql.functions import col
from pyspark.sql import Window

def fuzzy_match(dfs, dfc, path_summary):
    """
    Implements the Levenshtein and Soundex algorithms and returns a fuzzy matched DataFrame.
    Filters out those where resulting LS distance is less than 15% of SF name length.
    """
    # Apply Levenshtein and Soundex functions
    dfs = dfs.withColumn("OrganisationNameKeyLen", F.length("OrganisationNameKey"))
    df = dfc\
        .crossJoin(dfs)\
        .withColumn( "LevenshteinDistance", F.levenshtein( F.lower("OrganisationNameKey") , F.lower("CompanyNameKey") ) )\
        .withColumn( "HasSameSoundex", F.soundex("OrganisationNameKey") == F.soundex("CompanyNameKey") )\
        .where("LevenshteinDistance < OrganisationNameKeyLen * 0.15")\
        .orderBy("OrganisationName", "CompanyName")
    

def fuzzy_match_approve(df, path_fuzzy_match_approved, path_fuzzy_match_rejected, path_summary):
    """
    Filters fuzzy matching DataFrame results on approved/rejected based on set of conditions:
    - If there is only 1 match against the SF name
    - If more than 1 match then take that with LS distance of 1
    - If more than 1 match and more multiple LS distances of 1, then take the one where Soundex codes are the same
    Writes results and summary to CSV.
    """

    def write_with_bucket(df, bucket_col, path):
        df.write\
            .mode("overwrite")\
            .bucketBy(26, bucket_col)\
            .option("path", path)\
            .option("header", True)\
            .saveAsTable("bucket", format="csv")
    

    # Add window function columns:
    #   OrganisationNameMatchCount: Count AccountID per OrganisationName
    #   LevenshteinDistance1Count: Count AccountID per OrganisationName where LevenshteinDistance = 1
    windowSpec = Window.partitionBy("OrganisationName")
    df = df\
        .select("AccountID", "OrganisationName", "OrganisationNameKey", "CompanyNumber", "CompanyName", "LevenshteinDistance", "HasSameSoundex")\
        .withColumn("OrganisationNameMatchCount", F.count("AccountID").over(windowSpec))\
        .withColumn("LevenshteinDistance1Count", F.count(F.when(F.col("LevenshteinDistance")==1, F.col("AccountID"))).over(windowSpec))
    
    # Add bucket key column
    df = df.withColumn( "OrganisationNameBucketKey", F.substring( col("OrganisationNameKey"),0,1) )

    # Define fuzzy match approved condition
    is_approved_1 = ( F.col("OrganisationNameMatchCount") == 1 )
    is_approved_2 = ( (F.col("OrganisationNameMatchCount") > 1) & (F.col("LevenshteinDistance1Count") == 1) & (F.col("LevenshteinDistance") == 1) )
    is_approved_3 = ( (F.col("OrganisationNameMatchCount") > 1) & (F.col("LevenshteinDistance1Count") > 1) & (F.col("HasSameSoundex") == 'true') )
    is_approved = is_approved_1 | is_approved_2 | is_approved_3
    
    # Split fuzzy match results into approved and rejected
    df_approved = df.filter(is_approved)
    df_rejected = df.filter(~is_approved)

    # Export results
    # df_approved.write.csv(path_fuzzy_match_approved, mode="overwrite", header=True, quoteAll=True)
    # df_rejected.write.csv(path_fuzzy_match_rejected, mode="overwrite", header=True, quoteAll=True)
    write_with_bucket(df_approved, "OrganisationNameBucketKey", path_fuzzy_match_approved)
    write_with_bucket(df_rejected, "OrganisationNameBucketKey", path_fuzzy_match_rejected)


def main():
    spark = SparkSession...

    # Apply fuzzy match
    dfs = spark.read...
    dfc = spark.read...
    path_summary = ...
    df_fuzzy_match = fuzzy_match(dfs, dfc, path_summary)

    # Export results
    path_fuzzy_match_approved = ...
    path_fuzzy_match_rejected = ...
    fuzzy_match_approve(df_fuzzy_match, path_fuzzy_match_approved, path_fuzzy_match_rejected, path_summary)


main()

Other info:

  • df.rdd.getNumPartitions() is 2
  • dfs.count() is 12,515
  • dfc.count() is 5,110,430

Jobs: enter image description here

How can I improve performance here and get the results into a CSV successfully?

1 Answers

There are a couple of things you can do to improve your computation:

Improve parallelism

As Nithish mentioned in the comments, you don't have enough partitions in your input data frames to make use of all your CPU cores. You're not using all your CPU capability and this will slow you down.

To increase your parallelism, repartition dfc to at least your number of cores:

dfc = dfc.repartition(dfc.sql_ctx.sparkContext.defaultParallelism)

You need to do this because your crossJoin is run as a BroadcastNestedLoopJoin which doesn't reshuffle your large input dataframe.

Separate your computation stages

A Spark dataframe/RDD is conceptually just a directed action graph (DAG) of operations to run on your input data but it does not hold data. One consequence of this behavior is that, by default, you'll rerun your computations as many times as you reuse your dataframe.

In your fuzzy_match_approve function, you run 2 separate filters on your df, this means you rerun the whole cross-join operations twice. You really don't want this !

One easy way to avoid this is to use cache() on your fuzzy_match result which should be fairly small given your inputs and matching criteria.

def fuzzy_match_running(dfs, dfc, path_summary):
    """
    Implements the Levenshtein and Soundex algorithms and returns a fuzzy matched DataFrame.
    Filters out those where resulting LS distance is less than 15% of SF name length.
    """
    # Apply Levenshtein and Soundex functions
    dfs = dfs.withColumn("OrganisationNameKeyLen", F.length("OrganisationNameKey")).cache()
    dfc = dfc.repartition(dfc.sql_ctx.sparkContext.defaultParallelism).cache()
    
    df = dfc.crossJoin(dfs) \
        .withColumn( "LevenshteinDistance", F.levenshtein( F.lower("OrganisationNameKey") , F.lower("CompanyNameKey") ) ) \
        .withColumn( "HasSameSoundex", F.soundex("OrganisationNameKey") == F.soundex("CompanyNameKey") ) \
        .where("LevenshteinDistance < OrganisationNameKeyLen * 0.15") \
        .orderBy("OrganisationName", "CompanyName") \
        .cache()
    return df

If I run my fuzzy_match_running on some example data frames on my 8 core/16 threads I9-9980HK laptop (spark in local[*] mode with 8GB driver memory):

dfc rowcount : 572494
dfs rowcount : 17728
fuzzy_match rowcount: 7228499
Duration: 679.5572581291199 seconds
Matches/core/sec: 933436.210726889

The job takes about 12 min doing 572494*17728 ~ 10 billion row comparisons at 933k comparisons/seconds/core. Since your job does 64 billions row comparisons I would expect it to take about 80 min on my laptop.

You should run a similar experiment on your computer with a smaller sample to get an idea of your actual computing speed.

Going further: maximizing matches/sec

To go faster, we need to adjust the computation and increase the number of comparisons that can be done per seconds. A few things stand out in the function:

  • you filter your output by comparing the levenshtein distance, an integer, to a decimal calculation. This means spark will cast your integer to a decimal and operate on decimal. Comparing decimals is much slower than integers and it's unnecessary here, you can cast the bound to an int beforehand.
  • your levenshtein operates on the lower versions of your keys, this means, for each row comparison, Spark will convert the column values to lower again and again, wasting CPU cycles for redundant stuff. You can preprocess this before your join.

I update the function like this:

def fuzzy_match(dfs: DataFrame, dfc: DataFrame, path_summary: str) -> DataFrame:
    dfs = dfs.withColumn("OrganisationNameKeyLower", F.lower("OrganisationNameKey"))\
        .withColumn("MatchingTolerance", F.ceil(F.length("OrganisationNameKey") * 0.15).cast("int"))\
        .cache()

    dfc = dfc.repartition(dfc.sql_ctx.sparkContext.defaultParallelism)\
        .withColumn("CompanyNameKeyLower", F.lower("CompanyNameKey"))\
        .cache()

    df = dfc.crossJoin(dfs)\
        .withColumn("LevenshteinDistance", F.levenshtein(F.col("OrganisationNameKeyLower"), F.col("CompanyNameKeyLower")).cast("int")) \
        .where("LevenshteinDistance < MatchingTolerance")\
        .drop("MatchingTolerance")\
        .cache()

    # clean unnecessary caches before returning
    dfs.unpersist()
    dfc.unpersist()
    return df

When running the updated version on the same inputs as before and on the same computer I get nearly twice the performance as the first implementation

dfc rowcount : 572494
dfs rowcount : 17728
fuzzy_match rowcount: 7228499
Duration: 356.23311281204224 seconds
Matches/core/sec: 1780641.1846241967

If that is still too slow for your needs, you'll need to find conditions on your data that you can use as a join condition but that's highly data and use case specific.

Related