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

Spark:基于关联数据集多列组合键及金额条件过滤数据集

Filter ds1 Using ds2's Composite Key with AMT=0

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:51:11