Spark Scala中基于窗口函数的重复记录过滤及场景适配问题
Let's tackle both of your scenarios one by one to get the expected results:
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 validO|!|rows with non-nullAuditorID. - The
|| $"FFAction|!|" =!= "O|!|"clause allows all other action types (likeD|!|,I|!|) to pass through even ifAuditorIDis null, as long as they're the latest in their partition.
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:
- Including
UpdateReason_updateReasonIdin the partition split your target group into two separate partitions (one withUpdateReason_updateReasonId=6, one with null), so the latest timestamp wasn't being picked correctly. - The filter condition was inverted - it rejected
O|!|rows with nullUpdateReason_updateReasonIdinstead 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_updateReasonIdfrom 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 ofUpdateReason_updateReasonIdbeing null), OR - It's a non-
O|!|action with a non-nullUpdateReason_updateReasonId.
- It's an
This will return the row with the latest timestamp (2018-05-10T10:11:15+00:00) as expected.
内容的提问来源于stack exchange,提问作者Sudarshan kumar

