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

如何用PySpark访问Spark DataFrame嵌套数组列中的元素

Spark DataFrame访问嵌套数组元素的方法

给定的DataFrame Schema

root
 |-- CONTRATO: long (nullable = true)
 |-- FECHA_FIN: date (nullable = true)
 |-- IMPORTE_FIN: double (nullable = true)
 |-- MOVIMIENTOS: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- FECHA: date (nullable = true)
 |    |    |-- IMPORTE: double (nullable = true)

数据示例

[Row(CONTRATO=1, FECHA_FIN=datetime.date(2022, 10, 31), IMPORTE_FIN=895.83, MOVIMIENTOS=[Row(FECHA=datetime.date(2020, 9, 14), IMPORTE=10), Row(FECHA=datetime.date(2020, 9, 15), IMPORTE=20)])]

[Row(CONTRATO=2, FECHA_FIN=datetime.date(2022, 9, 30), IMPORTE_FIN=5.83, MOVIMIENTOS=[Row(FECHA=datetime.date(2021, 9, 14), IMPORTE=30), Row(FECHA=datetime.date(2020, 7, 15), IMPORTE=40)])]

需求说明

熟悉Pandas但刚接触Spark,希望实现类似Pandas的索引访问操作,比如:

df['MOVIMIENTOS'][df['CONTRATO'] == 1][0][0] --> 14/09/2020
df['MOVIMIENTOS'][df['CONTRATO'] == 1][0][1] --> 10
df['MOVIMIENTOS'][df['CONTRATO'] == 1][1][0] --> 15/09/2020
df['MOVIMIENTOS'][df['CONTRATO'] == 1][1][1] --> 20
df['MOVIMIENTOS'][df['CONTRATO'] == 2][1][0] --> 14/09/2021
df['MOVIMIENTOS'][df['CONTRATO'] == 2][1][1] --> 30

解决方案

Spark是分布式计算框架,DataFrame底层基于RDD,不能像Pandas那样直接通过链式索引访问本地数据,需要用Spark的列操作语法来实现:

1. 获取单个指定值

如果只是需要提取某一条记录的特定数组元素,可以先过滤行,再通过数组[索引].字段名提取,最后用first()或collect()拿到本地值:

# 获取CONTRATO=1的第一个MOVIMIENTOS的FECHA
fecha_val = df.filter(df.CONTRATO == 1) \
              .select(df.MOVIMIENTOS[0].FECHA) \
              .first()[0]
print(fecha_val.strftime('%d/%m/%Y'))  # 输出 14/09/2020

# 获取CONTRATO=1的第一个MOVIMIENTOS的IMPORTE
importe_val = df.filter(df.CONTRATO == 1) \
                .select(df.MOVIMIENTOS[0].IMPORTE) \
                .first()[0]
print(importe_val)  # 输出 10

# 获取CONTRATO=2的第二个MOVIMIENTOS的FECHA
fecha_val_2 = df.filter(df.CONTRATO == 2) \
                .select(df.MOVIMIENTOS[1].FECHA) \
                .first()[0]
print(fecha_val_2.strftime('%d/%m/%Y'))  # 输出 15/07/2020

2. 批量展开数组处理(推荐)

如果需要处理所有数组元素,更符合Spark的分布式场景,使用explode函数将数组展开为多行,之后可以直接访问struct字段:

from pyspark.sql.functions import explode

# 展开MOVIMIENTOS数组,生成每行对应一个movimiento元素
exploded_df = df.select(
    'CONTRATO',
    explode('MOVIMIENTOS').alias('movimiento')
).select(
    'CONTRATO',
    'movimiento.FECHA',
    'movimiento.IMPORTE'
)

# 筛选CONTRATO=1的所有记录
exploded_df.filter(exploded_df.CONTRATO == 1).show()
# 输出结果:
# +---------+----------+------+
# |CONTRATO|     FECHA|IMPORTE|
# +---------+----------+------+
# |        1|2020-09-14|    10|
# |        1|2020-09-15|    20|
# +---------+----------+------+

3. 生成包含指定数组元素的新列

如果需要把数组中指定位置的元素作为新列添加到DataFrame中:

# 添加第一个和第二个movimiento的FECHA、IMPORTE作为新列
df_expanded = df.withColumn('FECHA_1', df.MOVIMIENTOS[0].FECHA) \
                .withColumn('IMPORTE_1', df.MOVIMIENTOS[0].IMPORTE) \
                .withColumn('FECHA_2', df.MOVIMIENTOS[1].FECHA) \
                .withColumn('IMPORTE_2', df.MOVIMIENTOS[1].IMPORTE)

df_expanded.select('CONTRATO', 'FECHA_1', 'IMPORTE_1', 'FECHA_2', 'IMPORTE_2').show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 06:10:28