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

基于Scala+Spark的银行交易DataFrame:取消交易对应前置行删除需求

Solution for Removing Preceding Rows for Cancellation Transactions in Spark DataFrame (Scala)

Hey there! Let's work through this problem step by step—your requirement hinges on correctly identifying cancellation patterns and targeting the right rows for removal, which we can tackle with Spark window functions in Scala.

First, let's align on the logic we'll implement based on your description:

  • For a single cancellation transaction: Remove the immediately preceding row (plus the cancellation row itself, since it's a cancel action we don't want to retain in the final dataset)
  • For three consecutive cancellation transactions: Remove the three rows immediately before this sequence, along with the three cancellation rows themselves

Step 1: Set Up Sample Data & Spark Session

We'll start with a sample DataFrame to simulate banking transactions, using a case class to define transaction structure (order is critical here, so we'll use a timestamp to maintain sequence).

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

// Initialize Spark Session
val spark = SparkSession.builder()
  .appName("BankTransactionCleanup")
  .master("local[*]")
  .getOrCreate()

import spark.implicits._

// Case class to model transaction data
case class Transaction(id: Int, transType: String, amount: Double, timestamp: Long)

// Sample initial DataFrame matching your scenario
val initialDF = Seq(
  Transaction(1, "DEPOSIT", 100.0, 1620000000),
  Transaction(2, "WITHDRAW", 50.0, 1620000100),
  Transaction(3, "CANCEL", 0.0, 1620000200), // Single cancel: target row 2 for deletion
  Transaction(4, "DEPOSIT", 200.0, 1620000300),
  Transaction(5, "WITHDRAW", 100.0, 1620000400),
  Transaction(6, "TRANSFER", 150.0, 1620000500),
  Transaction(7, "CANCEL", 0.0, 1620000600), // Start of 3 consecutive cancels
  Transaction(8, "CANCEL", 0.0, 1620000700),
  Transaction(9, "CANCEL", 0.0, 1620000800), // Target rows 4,5,6 for deletion
  Transaction(10, "DEPOSIT", 300.0, 1620000900)
).toDF()

Step 2: Add Helper Columns to Detect Cancellation Patterns

We'll use window functions to add metadata that helps us identify single vs. consecutive cancel sequences:

  1. Assign row numbers to maintain transaction order
  2. Flag whether a transaction is a cancellation
  3. Look ahead/behind to detect consecutive cancel groups
// Window ordered by timestamp to ensure we process transactions in sequence
val timeWindow = Window.orderBy("timestamp")

// Add helper columns for pattern detection
val withHelpersDF = initialDF
  .withColumn("row_num", row_number().over(timeWindow))
  .withColumn("is_cancel", when(col("transType") === "CANCEL", 1).otherwise(0))
  // Look ahead to detect 3 consecutive cancels
  .withColumn("next_1_cancel", lead(col("is_cancel"), 1).over(timeWindow))
  .withColumn("next_2_cancel", lead(col("is_cancel"), 2).over(timeWindow))
  // Look behind to confirm membership in a 3-cancel sequence
  .withColumn("prev_1_cancel", lag(col("is_cancel"), 1).over(timeWindow))
  .withColumn("prev_2_cancel", lag(col("is_cancel"), 2).over(timeWindow))

Step 3: Flag Rows for Deletion

Now we'll create a to_delete column that marks all rows needing removal:

  • Rows immediately before a single cancel
  • The three rows before a 3-consecutive cancel sequence
  • All cancellation rows themselves
val flaggedDF = withHelpersDF
  .withColumn("to_delete", 
    // Flag row before a single cancel
    when((col("is_cancel") === 0 && lead(col("is_cancel"), 1).over(timeWindow) === 1 
          && lead(col("next_1_cancel"), 1).over(timeWindow) === 0), 1)
    // Flag three rows before a 3-consecutive cancel sequence
    .when((col("is_cancel") === 0 && lead(col("is_cancel"), 1).over(timeWindow) === 1 
          && lead(col("next_1_cancel"), 1).over(timeWindow) === 1 
          && lead(col("next_2_cancel"), 1).over(timeWindow) === 1), 1)
    // Flag all cancellation rows
    .when(col("is_cancel") === 1, 1)
    .otherwise(0)
  )

Step 4: Filter to Get Target DataFrame

Finally, we'll filter out marked rows and clean up helper columns to get your desired result:

val targetDF = flaggedDF
  .filter(col("to_delete") === 0)
  .drop("row_num", "is_cancel", "next_1_cancel", "next_2_cancel", "prev_1_cancel", "prev_2_cancel", "to_delete")

// View the final result
targetDF.show()

Key Notes for Adjustment

  • Ordering: If your transactions use a different ordering key (like transaction ID instead of timestamp), just swap out the orderBy column in the window definition.
  • Rule Tweaks: If you don't want to remove the cancellation rows themselves, simply remove the .when(col("is_cancel") === 1, 1) line from the to_delete logic.
  • Edge Cases: For edge scenarios (like cancels at the start of the dataset), the lead/lag functions will return null, which the logic handles automatically (since null === 1 evaluates to false).

内容的提问来源于stack exchange,提问作者jose rivera

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:07:55