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

Scala下DataFrame按列值删行与分组筛选及字段修改方案咨询

Hey there! Let's break down your two Scala Spark DataFrame questions clearly, with actionable code examples you can use right away.


1. How to delete rows from a DataFrame based on column values in Scala?

In Spark, "deleting" rows is essentially filtering out the ones that match your unwanted criteria. You’ve got two go-to methods here: filter() and where() (they work almost identically). Use them with reverse conditions or the ~ (negation) operator for Column expressions.

Here are common use cases with examples:
Suppose you have a DataFrame named df:

  • Remove rows where count is less than 9:
import org.apache.spark.sql.functions._

val filteredDF = df.filter(col("count") >= 9)
// Or use where—same result
val filteredDF = df.where(col("count") >= 9)
  • Remove rows where summary equals "equal":
val filteredDF = df.filter(col("summary") =!= "equal")
// Alternatively, use negation
val filteredDF = df.filter(!col("summary").equalTo("equal"))
  • Remove rows matching multiple conditions (e.g., date before 20170101 AND count equals 9):
val filteredDF = df.filter(!(col("date") < 20170101 && col("count") === 9))

2. Group by (id, count) to get rows with the minimum date, plus update "equal" to "more" in summary

For your specific DataFrame, we can solve this in a couple of ways—let’s start with a straightforward step-by-step approach, then show a more efficient window function method.

First, let’s recap your input data:

id secid count date summary
1 2 9 20170608 equal
1 3 9 20160608 equal
2 3 8 20170608 less
3 3 9 20160608 equal

Method 1: Group + Join + Update

Step 1: Find the minimum date per (id, count) group

First, aggregate to get the smallest date for each group:

val minDatePerGroup = df.groupBy("id", "count")
                        .agg(min("date").alias("min_date"))

Step 2: Join back to the original DataFrame to filter rows

Join the aggregated DataFrame with the original to keep only rows where date matches the group’s minimum date:

val minDateRows = df.join(minDatePerGroup,
  (df("id") === minDatePerGroup("id")) &&
  (df("count") === minDatePerGroup("count")) &&
  (df("date") === minDatePerGroup("min_date"))
).drop("min_date") // Clean up the extra column

Step 3: Update the summary column

Use when() to replace "equal" with "more", leave other values as-is:

val finalResult = minDateRows.withColumn("summary",
  when(col("summary") === "equal", "more")
    .otherwise(col("summary"))
)

Method 2: Window Function (More Efficient)

For larger datasets, window functions are better—no join needed, which saves performance:

import org.apache.spark.sql.expressions.Window

// Define a window partitioned by (id, count), ordered by date ascending
val windowSpec = Window.partitionBy("id", "count").orderBy("date")

// Add a row number, filter to keep only the first row (smallest date), then update summary
val finalResult = df.withColumn("row_num", row_number().over(windowSpec))
                    .filter(col("row_num") === 1)
                    .drop("row_num")
                    .withColumn("summary",
                      when(col("summary") === "equal", "more")
                        .otherwise(col("summary"))
                    )

Full Test Code

Here’s a complete snippet you can run to verify:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

object DataFrameTask {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("SparkDataFrameExample")
      .master("local[*]")
      .getOrCreate()
      
    import spark.implicits._
    
    // Create test DataFrame
    val inputData = Seq(
      (1, 2, 9, 20170608, "equal"),
      (1, 3, 9, 20160608, "equal"),
      (2, 3, 8, 20170608, "less"),
      (3, 3, 9, 20160608, "equal")
    ).toDF("id", "secid", "count", "date", "summary")
    
    // Using window function method (recommended)
    val windowSpec = Window.partitionBy("id", "count").orderBy("date")
    val finalDF = inputData.withColumn("row_num", row_number().over(windowSpec))
                          .filter(col("row_num") === 1)
                          .drop("row_num")
                          .withColumn("summary",
                            when(col("summary") === "equal", "more")
                              .otherwise(col("summary"))
                          )
    
    // Show the result
    finalDF.show()
  }
}

Running this will output exactly what you need:

+---+-----+-----+--------+-------+
| id|secid|count|    date|summary|
+---+-----+-----+--------+-------+
|  1|    2|    9|20170608|   more|
|  2|    3|    8|20170608|   less|
|  3|    3|    9|20160608|  equal|
+---+-----+-----+--------+-------+

内容的提问来源于stack exchange,提问作者Ravikumar Reddy Yeruva

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:55:42