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

Spark基于play/stop事件计算各视频总观看时长实现方法

Spark 视频总观看时长统计实现方案

基础信息说明

  • 源表:事件日志表 Table_event_log,共包含4个属性字段:
    • device_id:设备ID
    • video_id:视频ID
    • event_timestamp:事件时间戳
    • event_type:事件类型,本次计算仅用到play(播放)、stop(停止)两类事件
  • 计算规则:
    1. 按设备+视频维度匹配一一对应的play、stop事件对
    2. 每对匹配成功的事件的时间差,即为该设备观看对应视频的单次观看时长
    3. 最终按video_id维度汇总所有有效单次观看时长,得到每个视频的总观看时长
  • 计算逻辑示例:

以video_id=1的场景为例:Android设备从触发play事件到触发stop事件的时间差为1分钟,对应单次观看时长1分钟;Apple设备从触发play事件到触发stop事件的时间差为1分钟,对应单次观看时长1分钟,因此video_id=1的总观看时长为2分钟。

已完成的预处理代码

当前已完成测试数据集构造、表结构定义、时间戳格式转换,代码如下:

data1=[("Android",1,'2021-07-24 12:01:19.000',"play"),("Android",1,'2021-07-24 12:02:19.000',"stop"),
       ("Apple",1,'2021-07-24 12:03:19.000',"play"),("Apple",1,'2021-07-24 12:04:19.000',"stop"),]

schema1=StructType([StructField('device_id', StringType(),True),
                    StructField('video_id',IntegerType(),True),
                    StructField('event_timestamp',StringType(),True),
                    StructField('event_type',StringType(),True)
                       ])

transaction=spark.createDataFrame(data1,schema=schema1)
transaction=transaction.withColumn("Converted_timestamp",to_timestamp("event_timestamp"))

完整实现代码

核心实现思路为通过窗口函数,按设备、视频分组后按事件时间排序,将每个play事件和它之后紧邻的第一个stop事件配对,过滤无效配对后计算单次时长,最后按视频维度汇总即可。

from pyspark.sql.window import Window
from pyspark.sql.functions import col, lead, sum as spark_sum

# 定义窗口规则:同设备同视频下,按事件时间升序排序
match_window = Window.partitionBy("device_id", "video_id").orderBy("Converted_timestamp")

# 关联同组内下一个事件的时间、类型
event_pairs = transaction.withColumn("next_event_time", lead("Converted_timestamp", 1).over(match_window))\
                         .withColumn("next_event_type", lead("event_type", 1).over(match_window))

# 过滤出play后紧邻stop的有效配对,计算单次观看时长(单位:秒,可按需转换为分钟/小时)
single_watch_record = event_pairs.filter(
    (col("event_type") == "play") & (col("next_event_type") == "stop")
).withColumn(
    "single_duration",
    (col("next_event_time").cast("long") - col("Converted_timestamp").cast("long"))
)

# 按video_id汇总总观看时长
total_watch_res = single_watch_record.groupBy("video_id")\
                                     .agg(spark_sum("single_duration").alias("total_watch_seconds"))

# 输出查看结果
total_watch_res.show()

结果验证

针对提供的测试数据集运行代码,输出结果为video_id=1对应总观看时长120秒(即2分钟),和需求示例的计算结果完全一致。

内容的提问来源于stack exchange,提问作者Narendra Inamdar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 04:54:36