PySpark中如何基于登录时间匹配操作员对应登出时间?
解决PySpark中卡车操作员与登出时间匹配的问题
需求说明
需要为每个卡车操作员匹配对应的登出时间,规则为:同一卡车下,操作员的登出时间是其登录(firstlogin)到下一位操作员登录(lead_login)时间段内的最晚登出时间。
现有两个PySpark DataFrame:
df_op_login:存储操作员登录信息,字段包括truck、operator、firstlogin、lead_login、firstlogin_unix、leadlogin_unixdf_op_logout:存储卡车登出时间,未关联操作员,字段包括truck、log_out_ts、log_out_unix_ts
初始化代码如下:
import pyspark from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf import pyspark.sql.functions as sf from pyspark.sql.types import StringType from pyspark.sql.window import Window from pyspark.sql.types import (DoubleType, StringType, StructField, StructType, ArrayType, MapType, TimestampType, LongType) spark = SparkSession.builder.appName('SparkByExamples.com').getOrCreate() columns1 = ["truck","operator","firstlogin","lead_login","firstlogin_unix","leadlogin_unix"] data1 = [("t1","1", "2023-06-02T00:17:02.095+0000", "2023-06-02T01:57:31.868+0000",1685665022,1685671051), ("t1","2", "2023-06-02T01:57:31.868+0000","2023-06-02T02:25:55.484+0000",1685671051,1685672755), ("t1","3", "2023-06-02T02:25:55.484+0000","2023-06-02T13:56:47.373+0000",1685672755,1685714207), ("t1","1", "2023-06-02T13:56:47.373+0000","2023-06-02T23:53:39.829+0000",1685714207,1685750019), ("t1","4", "2023-06-02T23:53:39.829+0000",None,1685750019,None), ("t2","A", "2023-06-02T14:00:00.373+0000","2023-06-02T23:53:39.829+0000",1685714418,1685750019)] df_op_login = spark.createDataFrame(data=data1,schema=columns1) df_op_login = df_op_login.select( col('truck'),col('operator'),col('firstlogin'),col('lead_login'), col('firstlogin_unix').cast(LongType()),col('leadlogin_unix').cast(LongType()) ) columns2 = ["truck","log_out_ts","log_out_unix_ts"] data2 = [("t1","2023-06-02T00:02:50.242000",1685664170), ("t1","2023-06-02T01:50:42.436000",1685670642), ("t1","2023-06-02T01:53:31.231000",1685670811), ("t1","2023-06-02T02:00:27.855000",1685671227), ("t1","2023-06-02T02:46:20.058000",1685673980), ("t1","2023-06-02T05:57:15.370000",1685685435), ("t1","2023-06-02T06:06:30.526000",1685685990), ("t1","2023-06-02T06:15:16.062000",1685686516), ("t1","2023-06-02T10:26:53.795000",1685701613), ("t1","2023-06-02T10:33:43.012000",1685702023), ("t2","2023-06-02T15:35:43.012000",1685720143)] df_op_logout = spark.createDataFrame(data=data2,schema=columns2) df_op_logout = df_op_logout.select( col('truck'),col('log_out_ts'),col('log_out_unix_ts').cast(LongType()) )
问题与报错
编写了get_logout UDF尝试匹配时间段内的最晚登出时间,运行时报错:cannot resolve 'firstlogin_unix' given input columns: [log_out_ts, log_out_unix_ts, truck];
UDF代码如下:
w = Window.partitionBy("truck") def get_logout(x, y, z): df_op_logout_tmp = ( df_op_logout .filter(col('truck') == x) .filter(col('log_out_unix_ts') >= y) .filter(col('log_out_unix_ts') <= z) .withColumn('max_ts', sf.max('log_out_unix_ts').over(w)) .filter(sf.col('log_out_unix_ts') == sf.col('max_ts') ) ) if df_op_logout_tmp.count() > 0: return df_op_logout_tmp.collect()[0]['log_out_ts'] else: return None df_op_final = ( df_op_login .withColumn('op_logout',sf.lit(get_logout(col('truck'), col('firstlogin_unix'), col('leadlogin_unix')))) ) display(df_op_final)
报错原因分析
- UDF内部操作DataFrame违反Spark执行模型:UDF内调用
count()、collect()等action会触发局部计算,破坏Spark的分布式执行逻辑,导致上下文混乱。 - Window函数引用错误:UDF内部的
df_op_logout_tmp仅包含truck、log_out_ts、log_out_unix_ts列,但Window定义依赖的逻辑在该上下文无法关联到firstlogin_unix等列。 - UDF调用方式错误:用
sf.lit()包裹UDF调用是错误的,lit用于创建常量列,UDF需用sf.udf()包装后作为列函数使用。
解决方案:用Spark原生DSL实现
放弃UDF,改用区间关联+窗口函数的方式,符合Spark分布式执行逻辑,步骤如下:
- 处理登录数据,将
leadlogin_unix为None的记录替换为极大值(确保能匹配后续登出时间); - 按
truck关联登录和登出数据,筛选登出时间在操作员登录区间内的记录; - 用窗口函数按操作员分组,提取每组内的最晚登出时间;
- 去重后得到最终匹配结果。
完整代码
# 处理登录数据:替换None的leadlogin_unix为极大值,避免区间截断 df_login_processed = df_op_login.withColumn( "leadlogin_unix", sf.coalesce(col("leadlogin_unix"), sf.lit(9999999999).cast(LongType())) ) # 关联登录与登出数据,筛选登出时间在当前操作员的登录区间内 joined_df = df_login_processed.join( df_op_logout, on=[ df_login_processed.truck == df_op_logout.truck, df_op_logout.log_out_unix_ts >= df_login_processed.firstlogin_unix, df_op_logout.log_out_unix_ts <= df_login_processed.leadlogin_unix ], how="left" # 左关联确保没有匹配登出时间的操作员也能保留 ) # 定义窗口:按卡车、操作员、登录时间分组,取组内最大的登出时间戳 window_spec = Window.partitionBy("truck", "operator", "firstlogin_unix") # 提取每个操作员对应的最晚登出时间 df_op_final = joined_df.withColumn( "max_log_out_unix", sf.max("log_out_unix_ts").over(window_spec) ).filter( # 保留最大登出时间的记录,或无匹配登出的记录 col("log_out_unix_ts") == col("max_log_out_unix") | col("log_out_unix_ts").isNull() ).select( df_login_processed.truck, df_login_processed.operator, df_login_processed.firstlogin, df_login_processed.lead_login, df_login_processed.firstlogin_unix, df_login_processed.leadlogin_unix, col("log_out_ts").alias("op_logout") ).distinct() # 去重避免重复记录 # 查看结果 df_op_final.show(truncate=False)
预期结果
| truck | operator | firstlogin | lead_login | firstlogin_unix | leadlogin_unix | op_logout |
|---|---|---|---|---|---|---|
| t1 | 1 | 2023-06-02T00:17:02.095+0000 | 2023-06-02T01:57:31.868+0000 | 1685665022 | 1685671051 | 2023-06-02T01:53:31.231000 |
| t1 | 2 | 2023-06-02T01:57:31.868+0000 | 2023-06-02T02:25:55.484+0000 | 1685671051 | 1685672755 | 2023-06-02T02:00:27.855000 |
| t1 | 3 | 2023-06-02T02:25:55.484+0000 | 2023-06-02T13:56:47.373+0000 | 1685672755 | 1685714207 | 2023-06-02T10:33:43.012000 |
| t1 | 1 | 2023-06-02T13:56:47.373+0000 | 2023-06-02T23:53:39.829+0000 | 1685714207 | 1685750019 | null |
| t1 | 4 | 2023-06-02T23:53:39.829+0000 | null | 1685750019 | 9999999999 | null |
| t2 | A | 2023-06-02T14:00:00.373+0000 | 2023-06-02T23:53:39.829+0000 | 1685714418 | 1685750019 | 2023-06-02T15:35:43.012000 |
内容的提问来源于stack exchange,提问作者Ravi Teja
相关产品推荐
相关产品推荐

