基于Scala+Spark的银行交易DataFrame:取消交易对应前置行删除需求
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:
- Assign row numbers to maintain transaction order
- Flag whether a transaction is a cancellation
- 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
orderBycolumn 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 theto_deletelogic. - Edge Cases: For edge scenarios (like cancels at the start of the dataset), the
lead/lagfunctions will returnnull, which the logic handles automatically (sincenull === 1evaluates to false).
内容的提问来源于stack exchange,提问作者jose rivera

