You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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:

  1. Duplicate Rows: Rows where id, val, visits, and date are identical get marked as "Ramdom".
  2. Single Row per ID: If an ID has only one record, it gets "one".
  3. Date-based Remarks: For groups with matching id, val, and visits but different dates, the earlier date is "oldDate" and the later is "recentdate".
  4. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 09:52:47