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

Spark Scala中基于窗口函数的重复记录过滤及场景适配问题

Let's tackle both of your scenarios one by one to get the expected results:


Scenario 1

Problem Statement

You need to filter records where rank=1, with the rule that only when FFAction|!|="O|!" does AuditorID need to be non-null.

Sample Input DataFrame

DataPartition TimeStamp OrganizationID SourceID AuditorID AuditorEnumerationId AuditorOpinionCode AuditorOpinionId IsPlayingAuditorRole IsPlayingCSRAuditorRole IsPlayingTaxAdvisorRole FFAction|!| AuditorOpinionOnInternalControlCode AuditorOpinionOnGoingConcernCode AuditorOpinionOnInternalControlsId AuditorOpinionOnGoingConcernId rank 
Japan 2018-05-03T09:52:48+00:00 4295876589 194 null null null null null null null O|!| null null null null 1 
Japan 2018-05-03T09:52:48+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O|!| null null null null 1 
Japan 2018-05-03T09:52:48+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O|!| null null null null 1 
Japan 2018-05-03T07:36:47+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O|!| null null null null 1 
Japan 2018-05-03T07:36:47+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O|!| null null null null 1 
Japan 2018-05-02T10:37:50+00:00 4295876589 156 null null null null null null null O|!| null null null null 1 
Japan 2018-05-02T10:37:50+00:00 4295876589 157 null null null null null null null O|!| null null null null 1 
Japan 2018-05-02T10:37:56+00:00 4295876589 193 null null null null null null null O|!| null null null null 1 
Japan 2018-05-03T08:10:19+00:00 4295876589 196 null null null null null null null D|!| null null null null 1 
Japan 2018-05-03T09:52:48+00:00 4295876589 195 null null null null null null null O|!| null null null null 1 

Expected Output

DataPartition TimeStamp OrganizationID SourceID AuditorID AuditorEnumerationId AuditorOpinionCode AuditorOpinionId IsPlayingAuditorRole IsPlayingCSRAuditorRole IsPlayingTaxAdvisorRole FFAction|!| AuditorOpinionOnInternalControlCode AuditorOpinionOnGoingConcernCode AuditorOpinionOnInternalControlsId AuditorOpinionOnGoingConcernId rank 
Japan 2018-05-03T09:52:48+00:00 4295876589 194 2719 3023331 AOP 3010542 true false true O|!| null null null null 1 
Japan 2018-05-03T09:52:48+00:00 4295876589 195 16157 1002485247 UWE 3010547 true false false O|!| null null null null 1 
Japan 2018-05-03T07:36:47+00:00 4295876589 196 3252 3024053 ONC 3020538 true false true O|!| null null null null 1 
Japan 2018-05-03T07:36:47+00:00 4295876589 195 5937 3026578 NOP 3010543 true false true O|!| null null null null 1 
Japan 2018-05-02T10:37:56+00:00 4295876589 193 null null null null null null null I|!| null null null null 1 
Japan 2018-05-03T08:10:19+00:00 4295876589 196 null null null null null null null D|!| null null null null 1 

Existing Code

import org.apache.spark.sql.expressions._ 
val windowSpec = Window.partitionBy("OrganizationID", "SourceID", "AuditorID").orderBy(unix_timestamp($"TimeStamp", "yyyy-MM-dd'T'HH:mm:ss").cast("timestamp").desc) 
val latestForEachKey1 = finaldf.withColumn("rank", row_number().over(windowSpec)) 
 .filter($"rank" === 1 && $"AuditorID" =!= "null") 

Corrected Code & Explanation

Your original code filters out all rows with a null AuditorID, which is too strict. We need to only enforce non-null AuditorID for O|!| actions.

import org.apache.spark.sql.expressions._

val windowSpec = Window.partitionBy("OrganizationID", "SourceID", "AuditorID")
  .orderBy(unix_timestamp($"TimeStamp", "yyyy-MM-dd'T'HH:mm:ss").cast("timestamp").desc)

val latestForEachKey1 = finaldf.withColumn("rank", row_number().over(windowSpec))
  .filter(
    $"rank" === 1 && 
    (($"FFAction|!|" === "O|!|" && $"AuditorID".isNotNull) || $"FFAction|!|" =!= "O|!|")
  )
  • The condition ($"FFAction|!|" === "O|!|" && $"AuditorID".isNotNull) ensures we only keep valid O|!| rows with non-null AuditorID.
  • The || $"FFAction|!|" =!= "O|!|" clause allows all other action types (like D|!|, I|!|) to pass through even if AuditorID is null, as long as they're the latest in their partition.

Scenario 2

Problem Statement

You want to retain the latest (rank=1) records, where O|!| actions are allowed to have null UpdateReason_updateReasonId, and you need the most recent timestamp row for each core group.

Sample Input DataFrame

uniqueFundamentalSet PeriodId SourceId StatementTypeCode StatementCurrencyId UpdateReason_updateReasonId UpdateReasonComment UpdateReasonComment_languageId UpdateReasonEnumerationId FFAction|!| DataPartition PartitionYear TimeStamp 
192730230775 297 182 INC 500186 6 UpdateReasonToDelete 505074 3019685 I|!| Japan 2017 2018-05-10T09:57:29+00:00 
192730230775 297 182 INC 500186 6 UpdateReasonToDelete 505074 3019685 I|!| Japan 2017 2018-05-10T10:00:40+00:00 
192730230775 297 182 INC 500186 null null null null O|!| Japan 2017 2018-05-10T10:11:15+00:00 
192730230775 310 182 INC 500186 null null null null O|!| Japan 2018 2018-05-10T08:30:53+00:00 

Expected Output

192730230775 297 182 INC 500186 null null null null O|!| Japan 2017 2018-05-10T10:11:15+00:00 

Existing Code

val windowSpec = Window.partitionBy("uniqueFundamentalSet", "PeriodId", "SourceId", "StatementTypeCode", "StatementCurrencyId", "UpdateReason_updateReasonId").orderBy(unix_timestamp($"TimeStamp", "yyyy-MM-dd'T'HH:mm:ss").cast("timestamp").desc) 
val latestForEachKey1 = tempReorder.withColumn("rank", row_number().over(windowSpec)) 
 .filter($"rank" === 1 && (($"UpdateReason_updateReasonId" =!= "null" && $"FFAction|!|" === "O|!|") || $"FFAction|!|" =!= "O|!|")).drop("rank") 

Corrected Code & Explanation

Two key issues broke your original code:

  1. Including UpdateReason_updateReasonId in the partition split your target group into two separate partitions (one with UpdateReason_updateReasonId=6, one with null), so the latest timestamp wasn't being picked correctly.
  2. The filter condition was inverted - it rejected O|!| rows with null UpdateReason_updateReasonId instead of allowing them.

Here's the fixed code:

import org.apache.spark.sql.expressions._

// Remove UpdateReason_updateReasonId from partitionBy to group all core records together
val windowSpec = Window.partitionBy("uniqueFundamentalSet", "PeriodId", "SourceId", "StatementTypeCode", "StatementCurrencyId")
  .orderBy(unix_timestamp($"TimeStamp", "yyyy-MM-dd'T'HH:mm:ss").cast("timestamp").desc)

val latestForEachKey1 = tempReorder.withColumn("rank", row_number().over(windowSpec))
  .filter(
    $"rank" === 1 && 
    // Keep latest row: allow O|!| with null UpdateReason, or non-O|!| with non-null UpdateReason
    (($"FFAction|!|" === "O|!|") || ($"FFAction|!|" =!= "O|!|" && $"UpdateReason_updateReasonId".isNotNull))
  )
  .drop("rank")
  • Partition Adjustment: By removing UpdateReason_updateReasonId from the partition, we group all rows belonging to the same core set together, so we can correctly select the latest timestamp row.
  • Filter Logic: The condition keeps the latest row, and either:
    • It's an O|!| action (regardless of UpdateReason_updateReasonId being null), OR
    • It's a non-O|!| action with a non-null UpdateReason_updateReasonId.

This will return the row with the latest timestamp (2018-05-10T10:11:15+00:00) as expected.


内容的提问来源于stack exchange,提问作者Sudarshan kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:03:24