【发布时间】:2021-12-22 23:20:41
【问题描述】:
如何通过逗号将字符串列拆分为具有应用架构的新数据框?
例如,这是一个包含两列(id 和 value)的 pyspark DataFrame
df = sc.parallelize([(1, "200,201,hello"), (2, "23,24,hi")]).toDF(["id", "value"])
我想获取 value 列并将其拆分为一个新的 DataFrame 并应用以下架构:
from pyspark.sql.types import IntegerType, StringType, StructField, StructType
message_schema = StructType(
[
StructField("id", IntegerType()),
StructField("value", IntegerType()),
StructField("message", StringType()),
]
)
可行的方法是:
df_split = (
df.select(split(df.value, ",\s*"))
.rdd.flatMap(lambda x: x)
.toDF()
)
df_split.show()
但我仍然需要根据架构转换和重命名列:
df_split.select(
[
col(_name).cast(_schema.dataType).alias(_schema.name)
for _name, _schema in zip(df_split.columns, message_schema)
]
).show()
预期结果:
+---+-----+-------+
| id|value|message|
+---+-----+-------+
|200| 201| hello|
| 23| 24| hi|
+---+-----+-------+
【问题讨论】:
标签: python-3.x apache-spark pyspark apache-spark-sql