PySpark提取结构体数组奇偶索引元素并计算时序差值
纯PySpark管道函数解决方案
步骤拆解与实现代码
假设你的DataFrame已加载完成,以下是完整的纯PySpark函数实现流程:
from pyspark.sql import functions as F # 1. 提取moves中的other字段数组 df = df.withColumn("other_arr", F.col("moves.other")) # 2. 拆分偶数/奇数索引的元素(数组索引从0开始) # 提取偶数索引(玩家1的走棋序列) df = df.withColumn("player1_arr", F.expr(""" aggregate( posexplode(other_arr), array(), (acc, x) -> if(x.pos % 2 == 0, array_union(acc, array(x.col)), acc) ) """)) # 提取奇数索引(玩家2的走棋序列) df = df.withColumn("player2_arr", F.expr(""" aggregate( posexplode(other_arr), array(), (acc, x) -> if(x.pos % 2 == 1, array_union(acc, array(x.col)), acc) ) """)) # 3. 将数组元素转换为null或clk字段的第0个元素 df = df.withColumn("player1_clk", F.transform( "player1_arr", lambda elem: F.when(elem.isNotNull(), elem["clk"][0]).otherwise(F.lit(None)) )) df = df.withColumn("player2_clk", F.transform( "player2_arr", lambda elem: F.when(elem.isNotNull(), elem["clk"][0]).otherwise(F.lit(None)) )) # 4. 计算非null相邻元素的时间差 def get_time_diffs(arr_col): # 先过滤数组中的null值 clean_arr = F.filter(arr_col, lambda x: x.isNotNull()) # 用zip_with配对原数组与"去掉第一个元素的数组",计算差值 return F.expr(f""" zip_with( {clean_arr}, slice({clean_arr}, 2, size({clean_arr})), (prev_time, curr_time) -> curr_time - prev_time ) """) df = df.withColumn("player1_time_diffs", get_time_diffs("player1_clk")) df = df.withColumn("player2_time_diffs", get_time_diffs("player2_clk"))
关键逻辑说明
- 奇偶索引拆分:利用
posexplode将数组展开为(位置索引, 元素)的键值对,再通过aggregate聚合符合奇偶位置条件的元素,替代Python中带步长的切片操作。 - 元素转换:用
transform遍历数组,通过when/otherwise判断元素是否为null,非null时直接提取clk[0]。 - 时间差计算:
- 先用
filter清理数组中的null值,避免无效计算; - 用
slice获取去掉第一个元素的子数组,再通过zip_with将原数组与子数组配对,计算后一个元素与前一个的差值,得到相邻时间差数组。
- 先用
示例输出
假设输入数据的moves列包含序列:[{"other": {"clk": [100]}}, None, {"other": {"clk": [300]}}, {"other": {"clk": [450]}}, None, {"other": {"clk": [600]}}],最终输出的关键列结果如下:
+---------------+-------------------+---------------+-------------------+ |player1_clk |player1_time_diffs |player2_clk |player2_time_diffs | +---------------+-------------------+---------------+-------------------+ |[100, 300, null]|[200] |[null, 450, 600]|[150] | +---------------+-------------------+---------------+-------------------+
内容的提问来源于stack exchange,提问作者Buzz Moschetti
相关产品推荐
相关产品推荐

