How to bin in PySpark?

Viewed 32155

For example, I'd like to classify a DataFrame of people into the following 4 bins according to age.

age_bins = [0, 6, 18, 60, np.Inf]
age_labels = ['infant', 'minor', 'adult', 'senior']

I would use pandas.cut() to do this in pandas. How do I do this in PySpark?

3 Answers

You could also write a PySpark UDF:

def categorizer(age):
  if age < 6:
    return "infant"
  elif age < 18:
    return "minor"
  elif age < 60:
    return "adult"
  else: 
    return "senior"

Then:

bucket_udf = udf(categorizer, StringType() )
bucketed = df.withColumn("bucket", bucket_udf("age"))

In my case I had to randomly bucket a string value column, so it required me some extra steps:

from pyspark.sql.types import LongType, IntegerType
import pyspark.sql.functions as F


buckets_number = 4    # number of buckets desired

df.withColumn("sub", F.substring(F.md5('my_col'), 0, 16)) \
  .withColumn("translate", F.translate("sub", "abcdefghijklmnopqrstuvwxyz", "01234567890123456789012345").cast(LongType())) \
  .select("my_col",
         (F.col("translate") % (buckets_number + 1)).cast(IntegerType()).alias("bucket_my_col"))
  1. hash it with MD5
  2. substring the result to 16 characters (otherwise would have a too big number in following steps)
  3. translate letters generated by MD5 in numbers
  4. apply modulo function based on the number of desired buckets
Related