如何使用PySpark DataFrame筛选出符合交易发生条件的记录
PySpark筛选目标交易记录
根据你的需求和示例数据,这里提供两种常见场景的实现方案:
场景1:筛选金额发生变化的交易记录(含首次交易)
从示例输入输出来看,你需要排除用户连续交易中金额无变化的记录,只保留首次交易和金额变动后的记录。实现步骤如下:
代码实现
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, lag # 初始化Spark会话 spark = SparkSession.builder.appName("FilterTransactionRecords").getOrCreate() # 构造示例DataFrame data = [ ("Raju", "15-Jan-2022", 1000), ("Raju", "16-Jan-2022", 1000), ("Raju", "17-Jan-2022", 1500) ] df = spark.createDataFrame(data, ["Name", "transaction_date", "amount"]) # 转换日期格式,保证排序逻辑正确 df = df.withColumn("transaction_date", col("transaction_date").cast("date")) # 定义窗口:按用户分区,按交易日期升序排列 window_spec = Window.partitionBy("Name").orderBy("transaction_date") # 新增列存储上一次交易的金额 df_with_prev = df.withColumn("prev_amount", lag("amount").over(window_spec)) # 筛选首次交易(无上一次金额)或当前金额与上一次不同的记录 filtered_df = df_with_prev.filter( col("prev_amount").isNull() | (col("amount") != col("prev_amount")) ).drop("prev_amount") # 查看结果 filtered_df.show()
输出结果
+----+------------------+------+ |Name|transaction_date |amount| +----+------------------+------+ |Raju|2022-01-15 |1000 | |Raju|2022-01-17 |1500 | +----+------------------+------+
场景2:筛选存在实际交易的记录(排除金额为0的情况)
如果你的需求是排除无交易(金额为0)的记录,直接用过滤条件即可:
代码实现
# 筛选金额大于0的记录 filtered_df = df.filter(col("amount") > 0) filtered_df.show()
内容的提问来源于stack exchange,提问作者Naveen
相关产品推荐
相关产品推荐

