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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 23:19:01