Trying to call pyspark.sql.functions.from_json from within a udf fails with: AttributeError: 'NoneType' object has no attribute '_jvm' Are there any restrictions/limitations calling pyspark.sql.functions from within udf's?
Spark 2.3.0 on Linux
# spark version 2.3.0.cloudera2
# Using Python version 3.6.8
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType
from pyspark.sql.functions import from_json, udf
spark = SparkSession.builder.appName('UdfFromJson').getOrCreate()
df1 = spark.createDataFrame(['{"f1":"val1"}'], "string").toDF("data_json_str")
df2 = df1.select('data_json_str') \
.withColumn('data_json_parsed', from_json(col=df1.data_json_str, schema=StructType([StructField("f1",StringType(),True)]))) \
.select('data_json_str','data_json_parsed')
df2.show()
# This works:
# +-------------+----------------+
# |data_json_str|data_json_parsed|
# +-------------+----------------+
# |{"f1":"val1"}| [val1]|
# +-------------+----------------+
# Calling from_json within the udf does not work
def parse_json_py(col):
return from_json(col=col, schema=StructType([StructField("f1",StringType(),True)]))
parse_json = udf(f=parse_json_py, returnType=StructType())
df3 = df1.select('data_json_str') \
.withColumn('data_json_parsed', parse_json(df1.data_json_str)) \
.select('data_json_str','data_json_parsed')
df3.show()
# Fails with: AttributeError: 'NoneType' object has no attribute '_jvm'
Expected:
+-------------+----------------+
|data_json_str|data_json_parsed|
+-------------+----------------+
|{"f1":"val1"}| [val1]|
+-------------+----------------+
Actual:
AttributeError: 'NoneType' object has no attribute '_jvm'