Dependencies & cluster requirements
Dependencies & cluster requirements
Dependencies:
- ProphecySparkBasicsPython 0.0.1+
- ProphecySparkBasicsScala 0.0.1+
- UC dedicated clusters 14.3+ supported
- UC standard clusters 14.3+ supported
- Livy clusters 3.0.1+ supported
Parameters
| Parameter | Description |
|---|---|
| DataFrame | Input DataFrame |
| Row to keep |
|
| Deduplicate columns | Columns to consider while removing duplicate rows (not required for Distinct Rows) |
| Order columns | Columns to sort DataFrame on before de-duping in case of First and Last rows to keep. Sorting options:
|
Examples
Rows to keep: Any

def dedup(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.dropDuplicates(["tran_id"])
object dedup {
def apply(spark: SparkSession, in: DataFrame): DataFrame = {
in.dropDuplicates(List("tran_id"))
}
}
Rows to keep: First

def earliest_cust_order(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0\
.withColumn(
"row_number",
row_number()\
.over(Window\
.partitionBy("customer_id")\
.orderBy(col("order_dt").asc())
)\
.filter(col("row_number") == lit(1))\
.drop("row_number")
object earliest_cust_order {
def apply(spark: SparkSession, in: DataFrame): DataFrame = {
import org.apache.spark.sql.expressions.Window
in.withColumn(
"row_number",
row_number().over(
Window
.partitionBy("customer_id")
.orderBy(col("order_date").asc)
)
)
.filter(col("row_number") === lit(1))
.drop("row_number")
}
}
Rows to keep: Last

def latest_cust_order(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0\
.withColumn(
"row_number",
row_number()\
.over(Window\
.partitionBy("customer_id")\
.orderBy(col("order_dt").asc())
)\
.withColumn(
"count",
count("*")\
.over(Window\
.partitionBy("customer_id")
)\
.filter(col("row_number") == col("count"))\
.drop("row_number")\
.drop("count")
object latest_cust_order {
def apply(spark: SparkSession, in: DataFrame): DataFrame = {
import org.apache.spark.sql.expressions.Window
in.withColumn(
"row_number",
row_number().over(
Window
.partitionBy("customer_id")
.orderBy(col("order_date").asc)
)
)
.withColumn(
"count",
count("*").over(
Window
.partitionBy("customer_id")
)
)
.filter(col("row_number") === col("count"))
.drop("row_number")
.drop("count")
}
}
Rows to keep: Unique Only

def single_order_customers(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0\
.withColumn(
"count",
count("*")\
.over(Window\
.partitionBy("customer_id")
)\
.filter(col("count") == lit(1))\
.drop("count")
object single_order_customers {
def apply(spark: SparkSession, in: DataFrame): DataFrame = {
import org.apache.spark.sql.expressions.Window
in.withColumn(
"count",
count("*").over(
Window
.partitionBy("customer_id")
)
)
.filter(col("count") === lit(1))
.drop("count")
}
}
Rows to keep: Distinct Rows

def single_order_customers(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.distinct()
object single_order_customers {
def apply(spark: SparkSession, in: DataFrame): DataFrame = {
in.distinct()
}
}

