Azure Databricks Autoloader运行报Column is not iterable错误求助
问题根源
你在applytransform函数中通过inputdf._jdf.sparkSession()获取的是Java侧的SparkSession实例,调用它的sql()方法返回的是原生Java DataFrame对象(也就是你看到的py4j.java_gateway.JavaObject类型),没有经过PySpark的封装,所以不支持PySpark的withColumn等Python侧API,传参给PySpark的lit()函数时就会触发类型匹配错误。
修复方案
1. 修正applytransform函数
不要直接调用Java侧的SparkSession,使用当前PySpark环境的SparkSession即可,修改后代码如下:
from pyspark.sql import SparkSession def applytransform(inputdf,targettable,targetdatabase,fileconfig): logger.info('Inside applytransform for Database/Table {}.{}',targetdatabase,targettable) inputdf.createOrReplaceTempView("src_to_transform") # 替换原有获取Java SparkSession的代码,直接取当前PySpark的SparkSession spark = SparkSession.getActiveSession() if 'TransformQuery' in fileconfig and fileconfig['TransformQuery'] is not None: vsqlscript = fileconfig['TransformQuery'] # 用PySpark的sql方法,返回的是封装好的PySpark DataFrame df = spark.sql(vsqlscript) logger.info("Applied Tranform") return df else: logger.info("Passed DF") return inputdf
2. 可选优化点
如果担心多会话场景下获取SparkSession异常,也可以直接从输入的PySpark DataFrame获取所属会话,代码更稳妥:
spark = inputdf.sparkSession
这个方法是PySpark DataFrame原生提供的,直接返回绑定当前DataFrame的PySpark SparkSession实例,不需要碰底层Java对象。
验证
修改后先打印modifieddf的类型,确认返回为<class 'pyspark.sql.dataframe.DataFrame'>后再执行后续的加审计列操作即可,报错会自动消失。
内容的提问来源于stack exchange,提问作者bunker
相关产品推荐
相关产品推荐

