首次使用PySpark计算连续两日净新增独立URL数量
用PySpark计算连续两日净新增独立URL
核心思路
要实现“当日出现的URL若前一日已存在,则不计入次日净新增”的需求,核心是先提取每日的独立URL,再判断每个URL是否为当日首次出现或前一日未出现,最后按日期统计符合条件的URL数量。
具体实现步骤
假设你的原始数据集DataFrame名为df,包含Date(日期类型)和Url_path(URL路径字符串)两列:
- 提取每日独立URL
先对同一天的重复URL去重,确保每个URL在单日只统计一次:
daily_unique_urls = df.select("Date", "Url_path").distinct()
- 判断URL是否为当日净新增
这里提供两种高效实现方式:
方式一:自连接匹配前一日记录
通过左连接关联前一日的URL数据,未匹配到的即为当日净新增URL:
from pyspark.sql.functions import date_add, col, count, when # 自连接,匹配当前URL在前一天的记录 joined_df = daily_unique_urls.alias("curr") .join( daily_unique_urls.alias("prev"), (col("curr.Url_path") == col("prev.Url_path")) & (col("curr.Date") == date_add(col("prev.Date"), 1)), how="left_outer" ) # 按日期统计净新增数量 net_new_result = joined_df.groupBy("curr.Date") .agg( count(when(col("prev.Url_path").isNull(), col("curr.Url_path"))) .alias("net_new_unique_urls") ) .orderBy("curr.Date")
方式二:窗口函数追踪上一次出现日期
对每个URL按日期排序,通过lag函数获取上一次出现的日期,判断间隔是否为1天:
from pyspark.sql.window import Window from pyspark.sql.functions import lag, datediff, sum, when # 按URL分区、日期排序的窗口 url_window = Window.partitionBy("Url_path").orderBy("Date") # 标记是否为净新增:首次出现 或 上一次出现间隔≠1天 marked_df = daily_unique_urls.withColumn( "prev_appear_date", lag("Date").over(url_window) ).withColumn( "is_net_new", when( (col("prev_appear_date").isNull()) | (datediff(col("Date"), col("prev_appear_date")) != 1), 1 ).otherwise(0) ) # 按日期求和得到净新增数量 net_new_result = marked_df.groupBy("Date") .agg(sum("is_net_new").alias("net_new_unique_urls")) .orderBy("Date")
- 查看结果
执行net_new_result.show()即可得到每日的净新增独立URL数量。
关键说明
- 两种方式本质都是先确保单日URL唯一,再筛选出前一日未出现的URL进行统计
- 若你的
Date列是字符串类型,需要先转换为日期类型:df = df.withColumn("Date", to_date(col("Date"), "yyyy-MM-dd"))(根据实际日期格式调整)
内容的提问来源于stack exchange,提问作者Khushboo
相关产品推荐
相关产品推荐

