基于Scala为DataFrame按列值添加备注字段的技术问询
Solution in Scala (Spark)
Here's how you can implement the logic to add the remarks column using Spark window functions and conditional expressions:
First, set up the Spark session and sample input DataFrame:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window val spark = SparkSession.builder() .appName("AddRemarksColumn") .master("local[*]") .getOrCreate() // Sample input DataFrame val inputDF = Seq( (111, 2, 1, 20160122), (111, 2, 1, 20170122), (112, 4, 2, 20160122), (112, 5, 4, 20150122), (113, 6, 1, 20100120), (114, 8, 2, 20150122), (114, 8, 2, 20150122) ).toDF("id", "val", "visits", "date")
Next, define window functions to calculate helper metrics for applying the rules:
// Window partitions for different grouping levels val windowById = Window.partitionBy("id") val windowByIdValVisitsDate = Window.partitionBy("id", "val", "visits", "date") val windowByIdValVisits = Window.partitionBy("id", "val", "visits") // Calculate helper columns to evaluate each rule val dfWithHelpers = inputDF .withColumn("row_count_per_id", count("*").over(windowById)) .withColumn("duplicate_row_count", count("*").over(windowByIdValVisitsDate)) .withColumn("max_date_in_group", max("date").over(windowByIdValVisits)) .withColumn("min_date_in_group", min("date").over(windowByIdValVisits)) .withColumn("min_visits_in_id", min("visits").over(windowById)) .withColumn("max_visits_in_id", max("visits").over(windowById))
Apply the remark rules using conditional when expressions:
val resultDF = dfWithHelpers .withColumn("remarks", when(col("duplicate_row_count") > 1, "Ramdom") .when(col("row_count_per_id") === 1, "one") .when(col("date") === col("max_date_in_group"), "recentdate") .when(col("date") === col("min_date_in_group"), "oldDate") .when(col("visits") === col("min_visits_in_id"), "less") .when(col("visits") === col("max_visits_in_id"), "more") .otherwise("unknown") ) // Drop helper columns to get clean output .drop("row_count_per_id", "duplicate_row_count", "max_date_in_group", "min_date_in_group", "min_visits_in_id", "max_visits_in_id")
Show the final result:
resultDF.show()
Output:
+---+---+------+--------+----------+ | id|val|visits| date| remarks| +---+---+------+--------+----------+ |113| 6| 1|20100120| one| |111| 2| 1|20160122| oldDate| |111| 2| 1|20170122|recentdate| |112| 4| 2|20160122| less| |112| 5| 4|20150122| more| |114| 8| 2|20150122| Ramdom| |114| 8| 2|20150122| Ramdom| +---+---+------+--------+----------+
Explanation of Logic:
- Duplicate Rows: Rows where
id,val,visits, anddateare identical get marked as "Ramdom". - Single Row per ID: If an ID has only one record, it gets "one".
- Date-based Remarks: For groups with matching
id,val, andvisitsbut different dates, the earlier date is "oldDate" and the later is "recentdate". - Visits-based Remarks: For groups with differing visit counts, the smaller value is "less" and the larger is "more".
The order of conditions ensures higher-priority rules (like duplicates or single rows) are applied first before checking date/visit differences.
内容的提问来源于stack exchange,提问作者Rjj
相关产品推荐
相关产品推荐

