Forum Discussion
Reusable function for data transformation - user data functions
- 1 year ago
HI v-ssriganesh ,
As I explained in my previous response, based on Microsoft's response, the User data Functions cannot be used for transformations in the dataframe. Therefore, we need to utilize PySpark's native functions to transform data in the User data functions. So, I need to tweak my solution a bit and not use User data functions for Fabric but instead use pyspark.udf to do the transformation.
I think I know the way ahead now. Thanks for your support and help. We can close the ticket now.
Hi v-ssriganesh ,
Thanks for your response. I created a User data function as suggested, however, It did not work. Here is the code and scenario:
User data function
Code in the notebook:
Instantiate the function
data_functions = notebookutils.udf.getFunctions('data_functions')
Test the function
data_functions.convert_julian_to_date('123241')
Output: '2023-08-29 00:00:00'
Call the function for a dataframe column, and it fails. Error: TypeError: Column is not iterable
df_silver = df_silver.withColumn('request_date', data_functions.convert_julian_to_date(df_silver['request_date']))
---------------------------------------------------------------------------
TypeError Traceback (most recent call last)
Cell In[46], line 1
----> 1 df_silver = df_silver.withColumn('request_date', data_functions.convert_julian_to_date(df_silver['request_date']))
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/notebookutils/mssparkutils/handlers/udfHandler.py:95, in UDF.__create_dynamic_function.<locals>.dynamic_function(*args, **kwargs)
93 workspace_id = self.__metadata.get("folderObjectId", "")
94 capacity_id = self.__metadata.get("capacityObjectId", "")
---> 95 result = self.__udf_handler.run(artifact_id, name, parameters, workspace_id, capacity_id)
96 if json.loads(result).get("status", "").lower() != "succeeded":
97 raise Exception(f"Function {name} failed with error: {result}")
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/notebookutils/mssparkutils/handlers/udfHandler.py:27, in UdfHandler.run(self, artifact_id, function_name, parameters, workspace_id, capacity_id)
24 if not workspace_id:
25 workspace_id = self.getCurrentWorkspaceId()
---> 27 return self.jvm.notebookutils.udf.run(artifact_id, function_name, parameters, workspace_id, capacity_id)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1314, in JavaMember.__call__(self, *args)
1313 def __call__(self, *args):
-> 1314 args_command, temp_args = self._build_args(*args)
1316 command = proto.CALL_COMMAND_NAME +\
1317 self.command_header +\
1318 args_command +\
1319 proto.END_COMMAND_PART
1321 answer = self.gateway_client.send_command(command)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1277, in JavaMember._build_args(self, *args)
1275 def _build_args(self, *args):
1276 if self.converters is not None and len(self.converters) > 0:
-> 1277 (new_args, temp_args) = self._get_args(args)
1278 else:
1279 new_args = args
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1264, in JavaMember._get_args(self, args)
1262 for converter in self.gateway_client.converters:
1263 if converter.can_convert(arg):
-> 1264 temp_arg = converter.convert(arg, self.gateway_client)
1265 temp_args.append(temp_arg)
1266 new_args.append(temp_arg)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_collections.py:523, in MapConverter.convert(self, object, gateway_client)
521 java_map = HashMap()
522 for key in object.keys():
--> 523 java_map[key] = object[key]
524 return java_map
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_collections.py:82, in JavaMap.__setitem__(self, key, value)
81 def __setitem__(self, key, value):
---> 82 self.put(key, value)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1314, in JavaMember.__call__(self, *args)
1313 def __call__(self, *args):
-> 1314 args_command, temp_args = self._build_args(*args)
1316 command = proto.CALL_COMMAND_NAME +\
1317 self.command_header +\
1318 args_command +\
1319 proto.END_COMMAND_PART
1321 answer = self.gateway_client.send_command(command)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1277, in JavaMember._build_args(self, *args)
1275 def _build_args(self, *args):
1276 if self.converters is not None and len(self.converters) > 0:
-> 1277 (new_args, temp_args) = self._get_args(args)
1278 else:
1279 new_args = args
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_gateway.py:1264, in JavaMember._get_args(self, args)
1262 for converter in self.gateway_client.converters:
1263 if converter.can_convert(arg):
-> 1264 temp_arg = converter.convert(arg, self.gateway_client)
1265 temp_args.append(temp_arg)
1266 new_args.append(temp_arg)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/py4j/java_collections.py:510, in ListConverter.convert(self, object, gateway_client)
508 ArrayList = JavaClass("java.util.ArrayList", gateway_client)
509 java_list = ArrayList()
--> 510 for element in object:
511 java_list.add(element)
512 return java_list
File /opt/spark/python/lib/pyspark.zip/pyspark/sql/column.py:710, in Column.__iter__(self)
709 def __iter__(self) -> None:
--> 710 raise TypeError("Column is not iterable")
TypeError: Column is not iterable
Hi tinbaj,
Thank you for sharing the details and code. The error TypeError: Column is not iterable occurs because the User Data Function (UDF) is being applied to a Spark DataFrame column directly, which isn't compatible with the function's expectation of a single string input. To fix this, you need to register the UDF with Spark to handle DataFrame columns.
Here’s how to resolve it:
In your notebook, after instantiating the UDF, register it with Spark:
- Use spark.udf.register to make the UDF available for DataFrame operations.
- Then, apply it using withColumn with the registered UDF.
Update your notebook code as follows:
- Instantiate the UDF: data_functions = notebookutils.udf.getFunctions('data_functions')
- Register the UDF: spark.udf.register("convert_julian_to_date", data_functions.convert_julian_to_date)
- Apply to the DataFrame: df_silver = df_silver.withColumn('request_date', spark.sql.functions.expr("convert_julian_to_date(request_date)"))
This ensures the UDF processes each row’s request_date column value correctly. Also, verify that the request_date column in df_silver contains valid Julian date strings (e.g: '123241'). If the column has mixed or invalid data types, you may need to preprocess it to ensure all values are strings.
If this helps, please mark it as “Accept as solution” and feel free to give a “Kudos” to help others in the community as well.
Thank you.
- tinbaj1 year agoFrequent Visitor
Hi v-ssriganesh ,
Thanks for your response. I am getting this error when I am running the command to register the UDF: Register the UDF: spark.udf.register("convert_julian_to_date", data_functions.convert_julian_to_date)
--> 612 self.sparkSession._jsparkSession.udf().registerPython(name, register_udf._judf) 613 return return_udf File /opt/spark/python/lib/pyspark.zip/pyspark/sql/udf.py:321, in UserDefinedFunction._judf(self) 314 @property 315 def _judf(self) -> JavaObject: 316 # It is possible that concurrent access, to newly created UDF, 317 # will initialize multiple UserDefinedPythonFunctions. 318 # This is unlikely, doesn't affect correctness, 319 # and should have a minimal performance impact. 320 if self._judf_placeholder is None: --> 321 self._judf_placeholder = self._create_judf(self.func) 322 return self._judf_placeholder File /opt/spark/python/lib/pyspark.zip/pyspark/sql/udf.py:330, in UserDefinedFunction._create_judf(self, func) 327 spark = SparkSession._getActiveSessionOrCreate() 328 sc = spark.sparkContext --> 330 wrapped_func = _wrap_function(sc, func, self.returnType) 331 jdt = spark._jsparkSession.parseDataType(self.returnType.json()) 332 assert sc._jvm is not None File /opt/spark/python/lib/pyspark.zip/pyspark/sql/udf.py:59, in _wrap_function(sc, func, returnType) 57 else: 58 command = (func, returnType) ---> 59 pickled_command, broadcast_vars, env, includes = _prepare_for_python_RDD(sc, command) 60 assert sc._jvm is not None 61 return sc._jvm.SimplePythonFunction( 62 bytearray(pickled_command), 63 env, (...) 68 sc._javaAccumulator, 69 ) File /opt/spark/python/lib/pyspark.zip/pyspark/rdd.py:5251, in _prepare_for_python_RDD(sc, command) 5248 def _prepare_for_python_RDD(sc: "SparkContext", command: Any) -> Tuple[bytes, Any, Any, Any]: 5249 # the serialized command will be compressed by broadcast 5250 ser = CloudPickleSerializer() -> 5251 pickled_command = ser.dumps(command) 5252 assert sc._jvm is not None 5253 if len(pickled_command) > sc._jvm.PythonUtils.getBroadcastThreshold(sc._jsc): # Default 1M 5254 # The broadcast will have same life cycle as created PythonRDD File /opt/spark/python/lib/pyspark.zip/pyspark/serializers.py:469, in CloudPickleSerializer.dumps(self, obj) 467 msg = "Could not serialize object: %s: %s" % (e.__class__.__name__, emsg) 468 print_exec(sys.stderr) --> 469 raise pickle.PicklingError(msg) PicklingError: Could not serialize object: PySparkRuntimeError: [CONTEXT_ONLY_VALID_ON_DRIVER] It appears that you are attempting to reference SparkContext from a broadcast variable, action, or transformation. SparkContext can only be used on the driver, not in code that it run on workers. For more information, see SPARK-5063.
- v-ssriganesh1 year ago
Community Support
Hi tinbaj,
Thank you for providing the error details. The PicklingError: [CONTEXT_ONLY_VALID_ON_DRIVER] occurs because the User Data Function (UDF) is being serialized in a way that references the SparkContext, which isn't allowed in Spark's distributed environment. This is likely due to how the UDF is defined or accessed in your notebook.To resolve this, try the following steps:
Instead of directly registering the UDF with spark.udf.register, use the Fabric UDF directly in the DataFrame operation, as Fabric’s UDFs are designed to work seamlessly with Spark. Update your notebook code as follows:
- Instantiate the UDF: data_functions = notebookutils.udf.getFunctions('data_functions')
- Apply the UDF to the DataFrame: df_silver = df_silver.withColumn('request_date', data_functions.convert_julian_to_date(df_silver.request_date))
- Ensure your UDF (convert_julian_to_date) in the User Data Functions item doesn’t reference SparkContext or other non-serializable objects. Your provided UDF code looks fine, but confirm it only uses standard Python libraries (e.g., datetime, timedelta) and avoids Spark-specific calls.
- Before applying to the DataFrame, test the UDF with a single value to confirm it works: print(data_functions.convert_julian_to_date('123241')). This should return '2023-08-29 00:00:00'.
If the error persists, please share:
- The schema of df_silver (df_silver.printSchema()).
- Any modifications made to the UDF code.
- Whether you’re running this in a Fabric notebook with a Spark session active.
Please try these steps and let me know the outcome. If it resolves the issue, consider marking it as “Accept as solution” and giving a “Kudos” to help others in the community.
Thank you.- tinbaj1 year agoFrequent Visitor
Hi v-ssriganesh ,
Thanks for your response, but the suggested code did not fix the problem. I can confirm that I am using standard python libraries in UDF and does not use spark context.
The implementation as per the suggestion and the error message is as below:
data_functions = notebookutils.udf.getFunctions('data_functions')print(data_functions.convert_julian_to_date('123241'))Return Value: 2023-08-29df_silver = df_silver.withColumn("request_date", data_functions.convert_julian_to_date(df_silver.request_date))Error Message: PySparkTypeError: [NOT_ITERABLE] Column is not iterable.Now if we go through the documentation for UDF's (Link: https://learn.microsoft.com/en-us/fabric/data-engineering/user-data-functions/python-programming-model), column data type is not one of the acceptable data type in UDF's. could this be a reason for this error?Thanks