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
countis 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
summaryequals "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.,
datebefore 20170101 ANDcountequals 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

