Lesson reference: PySpark and Scala
Parse signup text
The setup builds customer_signups from the directory and assigns September 1–6 at 10:30 to its six IDs for this exercise. These dates are not fields in customers. to_date keeps the calendar date; to_timestamp also keeps the time. The session timezone is UTC. All inputs here are valid; malformed input needs a deliberate parsing policy.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
spark.conf.set("spark.sql.session.timeZone", "UTC")
customer_signups = customers.select(
"customer_id",
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
result = customer_signups.select(
col("customer_id"),
to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss").alias("signed_up_on")
).orderBy("customer_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
spark.conf.set("spark.sql.session.timeZone", "UTC")
val customer_signups = customers.select(
col("customer_id"),
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
val result = customer_signups.select(
col("customer_id"),
to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss").alias("signed_up_on")
).orderBy("customer_id")
result.show()
Filter signup dates
Parse text before comparing dates. This example keeps September 4 and later, including the boundary. Compare against a date literal instead of depending on text ordering. Filtering this working DataFrame does not remove customers from the directory.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
spark.conf.set("spark.sql.session.timeZone", "UTC")
customer_signups = customers.select(
"customer_id",
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
result = customer_signups.withColumn(
"signed_up_on", to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss")
).filter(col("signed_up_on") >= lit("2026-09-04").cast("date")).select("customer_id", "signed_up_on").orderBy("customer_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
spark.conf.set("spark.sql.session.timeZone", "UTC")
val customer_signups = customers.select(
col("customer_id"),
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
val result = customer_signups.withColumn(
"signed_up_on", to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss")
).filter(col("signed_up_on") >= lit("2026-09-04").cast("date")).select("customer_id", "signed_up_on").orderBy("customer_id")
result.show()
Days since signup
datediff(end, start) counts calendar-day boundaries, not elapsed 24-hour periods. September 10 is a fixed reporting date, making the example repeatable. Reversing the arguments reverses the sign.
PySpark example
from pyspark.sql.functions import (
col, lit, concat, trim, lower, upper, regexp_extract,
format_string, to_date, to_timestamp, datediff
)
spark.conf.set("spark.sql.session.timeZone", "UTC")
customer_signups = customers.select(
"customer_id",
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
result = customer_signups.select(
col("customer_id"),
datediff(lit("2026-09-10").cast("date"),
to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss")).alias("days_elapsed")
).orderBy("customer_id")
result.show()
Scala example
import org.apache.spark.sql.functions._
spark.conf.set("spark.sql.session.timeZone", "UTC")
val customer_signups = customers.select(
col("customer_id"),
format_string("2026-09-%02d 10:30:00", col("customer_id") - 100).alias("signup_text")
)
val result = customer_signups.select(
col("customer_id"),
datediff(lit("2026-09-10").cast("date"),
to_date(col("signup_text"), "yyyy-MM-dd HH:mm:ss")).alias("days_elapsed")
).orderBy("customer_id")
result.show()