Spark:基于关联数据集多列组合键及金额条件过滤数据集
Got it, let's tackle this problem step by step. You need to keep only those records in ds1 that match the composite key (ID, NAME, ACCID, CURR) from ds2 where the corresponding AMT value is 0. Here's the most efficient way to do this with Spark (Scala syntax, since your ds1 is defined using Scala):
Step 1: Prepare the Filter Keys from ds2
First, we'll extract just the composite key columns from ds2 where AMT equals 0. This gives us a focused dataset of keys we want to match against ds1:
// Filter ds2 for AMT=0, then select only the composite key columns val filterKeys = ds2.filter($"AMT" === 0).select("ID", "NAME", "ACCID", "CURR")
Step 2: Perform a Left Semi Join to Filter ds1
A left semi join is ideal here because it only returns records from ds1 that have a matching key in our filterKeys dataset—no extra columns, no duplicate rows, and it’s far more efficient than a regular join followed by cleanup:
// Join ds1 with filterKeys using the composite key, keeping only matching ds1 records val filteredDs1 = ds1.join(filterKeys, Seq("ID", "NAME", "ACCID", "CURR"), "leftsemi")
Optional: Optimize with Broadcast Join (Small Filter Dataset)
If your filterKeys dataset is small (e.g., thousands of rows or less), use a broadcast join to avoid shuffling data across the cluster—this can drastically speed up the operation:
import org.apache.spark.sql.functions.broadcast val filteredDs1 = ds1.join(broadcast(filterKeys), Seq("ID", "NAME", "ACCID", "CURR"), "leftsemi")
Example Verification
Let’s say ds2 has these relevant records:
val ds2 = Seq( ("T1002","WELLINGTON","8787","CAD",0), ("T1003","WELLINGTON","654","USD",0), ("T1001","WELLINGTON","991","CAD",100) // This AMT isn't 0, so it won't be included in filterKeys )
Running the code above would leave filteredDs1 with these rows from your original ds1:
- ("T1002","WELLINGTON","8787","CAD",1450)
- ("T1003","WELLINGTON","654","USD",200)
That’s exactly what you need—only the ds1 records that match the valid composite keys from ds2.
内容的提问来源于stack exchange,提问作者1pluszara

