2
votes

I create a dataframe DF from external file, which has the following schema:

(id, field1, field2, field3) partition column: id

data example is

 000,  11_field1,  22_field2,  33_field3
 001, 111_field1, 222_field2, 333_field3

I want to create another dataframe from DF which schema is

 (id, fieleName, fieldValue)

data example is

000, field1,  11_field1
000, field2,  22_field2
000, field3,  33_field3
001, field1, 111_field1
001, field2, 222_field2
001, field3, 333_field3

Could anyone tell me how to get the new dataframe?

1
Thank you for your answer. It works for me. - Yifei Xu
Could you please up vote my answer - User12345

1 Answers

1
votes

You can achieve this in pyspark like below using explode option

First import necessary libraries and functions

from pyspark.sql import SQLContext, Row

Say your data frame is df.

If you do df.show()

you should get result like below

+---+----------+----------+----------+
| id|    field1|    field2|    field3|
+---+----------+----------+----------+
|  0| 11_field1| 22_field2| 33_field3|
|  1|111_field1|222_field2|333_field3|
+---+----------+----------+----------+

Then map all columns you want to explode as 2 columns. Here you want all columns except id to explode. So, do the below

cols= df.columns[1:]

then convert the data frame to rdd like below

rdd = data.rdd.map(lambda x: Row(id=x[0], val=dict(zip(cols, x[1:]))))

To check how the rdd has been mapped do below

rdd.take()

you will get result like below

[Row(id=0, val={'field2': u'22_field2', 'field3': u'33_field3', 'field1': u'11_field1'}), Row(id=1, val={'field2': u'222_field2', 'field3': u'333_field3', 'field1': u'111_field1'})]

Then convert the rdd back to a data frame say df2

df2 = sqlContext.createDataFrame(rdd)

Then do df2.show(). you should get result like below

+---+--------------------+
| id|                 val|
+---+--------------------+
|  0|Map(field3 -> 33_...|
|  1|Map(field3 -> 333...|
+---+--------------------+

then register the data frame df2 as a temp table

df2.registerTempTable('mytempTable')

Then run a query like below on the data frame:

df3 = sqlContext.sql( """select id,explode(val) AS (fieldname,fieldvalue) from mytempTable""")

then do df3.show(), you should get the result as below

+---+---------+----------+
| id|fieldname|fieldvalue|
+---+---------+----------+
|  0|   field3| 33_field3|
|  0|   field2| 22_field2|
|  0|   field1| 11_field1|
|  1|   field3|333_field3|
|  1|   field2|222_field2|
|  1|   field1|111_field1|
+---+---------+----------+