如何从pyspark中的数组中提取元素

Anm*_*ave 8 python apache-spark rdd pyspark

我有一个以下类型的数据框

col1|col2|col3|col4
xxxx|yyyy|zzzz|[1111],[2222]
Run Code Online (Sandbox Code Playgroud)

我希望我的输出是跟随类型

col1|col2|col3|col4|col5
xxxx|yyyy|zzzz|1111|2222
Run Code Online (Sandbox Code Playgroud)

我的col4是一个数组,我想将它转换为一个单独的列.需要做什么?

我在flatmap中看到很多答案,但是他们正在增加一行,我希望只将元组放在另一列但是在同一行中

以下是我的实际架构:

root
 |-- PRIVATE_IP: string (nullable = true)
 |-- PRIVATE_PORT: integer (nullable = true)
 |-- DESTINATION_IP: string (nullable = true)
 |-- DESTINATION_PORT: integer (nullable = true)
 |-- collect_set(TIMESTAMP): array (nullable = true)
 |    |-- element: string (containsNull = true)
Run Code Online (Sandbox Code Playgroud)

也可以请一些人帮我解释数据帧和RDD

Psi*_*dom 21

创建样本数据:

from pyspark.sql import Row
x = [Row(col1="xx", col2="yy", col3="zz", col4=[123,234])]
rdd = sc.parallelize([Row(col1="xx", col2="yy", col3="zz", col4=[123,234])])
df = spark.createDataFrame(rdd)
df.show()
#+----+----+----+----------+
#|col1|col2|col3|      col4|
#+----+----+----+----------+
#|  xx|  yy|  zz|[123, 234]|
#+----+----+----+----------+
Run Code Online (Sandbox Code Playgroud)

用于getItem从数组列中提取元素,在实际情况下替换col4为collect_set(TIMESTAMP):

df = df.withColumn("col5", df["col4"].getItem(1)).withColumn("col4", df["col4"].getItem(0))
df.show()
#+----+----+----+----+----+
#|col1|col2|col3|col4|col5|
#+----+----+----+----+----+
#|  xx|  yy|  zz| 123| 234|
#+----+----+----+----+----+
Run Code Online (Sandbox Code Playgroud)

  • @Lydia 请*非常小心*,并确保您在更改代码时知道自己在做什么:您的编辑破坏了一个非常好的答案,导致它抛出异常(将其恢复为 OP 的原始版本)... (4认同)

Myk*_*tko 8

您有 4 个选项来提取数组内的值:

df = spark.createDataFrame([[1, [10, 20, 30, 40]]], ['A', 'B'])
df.show()

+---+----------------+
|  A|               B|
+---+----------------+
|  1|[10, 20, 30, 40]|
+---+----------------+

from pyspark.sql import functions as F

df.select(
    "A",
    df.B[0].alias("B0"), # dot notation and index        
    F.col("B")[1].alias("B1"), # function col and index
    df.B.getItem(2).alias("B2"), # dot notation and method getItem
    F.col("B").getItem(3).alias("B3"), # function col and method getItem
).show()

+---+---+---+---+---+
|  A| B0| B1| B2| B3|
+---+---+---+---+---+
|  1| 10| 20| 30| 40|
+---+---+---+---+---+
Run Code Online (Sandbox Code Playgroud)

如果您有很多列,请使用列表理解:

df.select(
    'A', *[F.col('B')[i].alias(f'B{i}') for i in range(4)]
).show()

+---+---+---+---+---+
|  A| B0| B1| B2| B3|
+---+---+---+---+---+
|  1| 10| 20| 30| 40|
+---+---+---+---+---+
Run Code Online (Sandbox Code Playgroud)