Skip to article frontmatterSkip to article content
Site not loading correctly?

This may be due to an incorrect BASE_URL configuration. See the MyST Documentation for reference.

Things on this page are fragmentary and immature notes/thoughts of the author. Please read with your own judgement!

Tips and Traps

  1. If you use PySpark instead of Spark/Scala, pandas udf is a great alternative to all those (complicated) collections functions discussed here. Leveraging pandas udf, each partition of a Spark DataFrame can be converted to a pandas DataFrame without copying the underlying data, you can then do transforms on pandas DataFrames which will be converted back to partitons of a Spark DataFrame.

  2. When converting a pandas DataFrame to a Spark DataFrame,

    • a column of list is converted to a column of ArrayType

    • a column of tuple is converted to a column of StructType

    • a column of dict is converted to a column of MapType

Q: how about dict?

Comemnts

There are multiple ways (vanilla string, JSON string, StructType and ArrayType) to represent complex data types in Spark DataFrames. Notice that a Tuple is converted to a StructType in Spark DataFrames and an Array is converted to a ArrayType in Spark DataFrames. Starting from Spark 2.4, you can use ArrayType which is more convenient if the elements have the same type.

Vanilla String

  • string, substring, regexp_extract, locate, left, concat_ws

JSON String

  • json_tuple

  • get_json_object

  • from_json

StructType

ArrayType

  • array

  • element_at

  • array_min, array_max, array_join, array_interesect, array_except, array_distinct, array_contains, array, array_position, array_remove, array_repeat, array_sort, array_union, array_overlap, array_zip

Python Types to DataType in PySpark

A column of list is converted to a Column of ArrayType.

+------+----+
|  col1|col2|
+------+----+
|[1, 2]| how|
|[2, 3]| are|
|[3, 4]| you|
+------+----+

StructType(List(StructField(col1,ArrayType(LongType,true),true),StructField(col2,StringType,true)))

A column of tuple is converted to a Column of StructType.

+------+----+
|  col1|col2|
+------+----+
|{1, 2}| how|
|{2, 3}| are|
|{3, 4}| you|
+------+----+

StructType(List(StructField(col1,StructType(List(StructField(_1,LongType,true),StructField(_2,LongType,true))),true),StructField(col2,StringType,true)))

A column of dict is converted to a column of MapType.

+----------------+----+
|            col1|col2|
+----------------+----+
|{x -> 1, y -> 2}| how|
|{x -> 2, y -> 3}| are|
|{x -> 3, y -> 4}| you|
+----------------+----+

StructType(List(StructField(col1,MapType(StringType,LongType,true),true),StructField(col2,StringType,true)))
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2| how|
|   2|   3| are|
|   3|   4| you|
+----+----+----+

+----------+
|       map|
+----------+
|{how -> 1}|
|{are -> 2}|
|{you -> 3}|
+----------+

+----------+
|       map|
+----------+
|{how -> 1}|
|{are -> 2}|
|{you -> 3}|
+----------+

explode

+---------------+
|          words|
+---------------+
|[how, are, you]|
+---------------+

+-----+
|words|
+-----+
|  how|
|  are|
|  you|
+-----+

split

collect

+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   2| how|
|   2|   3| are|
|   3|   4| you|
+----+----+----+

+------+
|struct|
+------+
|{1, 2}|
|{2, 3}|
|{3, 4}|
+------+

+------+
|struct|
+------+
|{1, 2}|
|{2, 3}|
|{3, 4}|
+------+

StructType(List(StructField(struct,StructType(List(StructField(col1,LongType,true),StructField(col2,LongType,true))),false)))

Work with StructType

Notice that a Tuple is converted to StructType in Spark DataFrames.

+------+----+
|  col1|col2|
+------+----+
|{1, 2}| how|
|{2, 3}| are|
|{3, 4}| you|
+------+----+

Split all elements of a StructType into different columns.

+---+---+
| _1| _2|
+---+---+
|  1|  2|
|  2|  3|
|  3|  4|
+---+---+

Extract elements from StructTypes by position and rename the columns.

+---+---+
| v1| v2|
+---+---+
|  1|  2|
|  2|  3|
|  3|  4|
+---+---+

Work with ArrayType

Notice that an Array is converted to an ArrayType in Spark DataFrames. Note: ArrayType requires Spark 2.4.0+.

+------+----+
|  col1|col2|
+------+----+
|[1, 2]| how|
|[2, 3]| are|
|[3, 4]| you|
+------+----+

null
+---+---+
| v1| v2|
+---+---+
|  1|  2|
|  2|  3|
|  3|  4|
+---+---+

ArrayType

+------+----+----+
|  col1|col2|col3|
+------+----+----+
|[1, 2]| how|   1|
|[2, 3]| are|   2|
|[3, 4]| you|   3|
+------+----+----+

+----+
|word|
+----+
|   1|
|   2|
|   3|
+----+

+----+
|word|
+----+
|   1|
|   2|
|   3|
+----+

+------+
|    f1|
+------+
|[1, 1]|
|[2, 1]|
|[3, 1]|
+------+

StructType(List(StructField(f1,ArrayType(IntegerType,true),true)))
+---+---+
| v1| v2|
+---+---+
|  1|  1|
|  2|  1|
|  3|  1|
+---+---+

StructType

+------+----+----+
|  col1|col2|col3|
+------+----+----+
|[1, 2]| how|   1|
|[2, 3]| are|   2|
|[3, 4]| you|   3|
+------+----+----+

StructType(List(StructField(col1,StructType(List(StructField(_1,LongType,true),StructField(_2,LongType,true))),true),StructField(col2,StringType,true),StructField(col3,LongType,true)))
+---+---+
| _1| _2|
+---+---+
|  1|  2|
|  2|  3|
|  3|  4|
+---+---+

+------+
|    f1|
+------+
|[1, 1]|
|[2, 1]|
|[3, 1]|
+------+

StructType(List(StructField(f1,StructType(List(StructField(_1,IntegerType,true),StructField(_2,IntegerType,true))),true)))
+---+---+
| _1| _2|
+---+---+
|  1|  1|
|  2|  1|
|  3|  1|
+---+---+

+---+---+
| v1| v2|
+---+---+
|  1|  1|
|  2|  1|
|  3|  1|
+---+---+