在HDP 2.6.5 Zeppelin中执行Spark DataFrame代码遇Py4JJavaError及Row未定义
解决PySpark中Row未定义的NameError问题
问题场景
在Hortonworks Sandbox HDP 2.6.5环境下,通过Zeppelin执行Spark2 PySpark的单词计数实验,将RDD转换为DataFrame时触发以下错误:
NameError: global name 'Row' is not defined
问题原因
代码中直接使用了Row构造DataFrame行对象,但未从pyspark.sql模块显式导入Row类,导致Spark Worker节点无法识别该全局名称。
解决方案
方案1:显式导入Row类
在代码开头添加Row的导入语句,修改后的完整代码如下:
%spark2.pyspark # 导入Row类 from pyspark.sql import Row # First, let's transform our RDD to a DataFrame. # We will use a Row to define column names. wordsCounts = (filteredWordCounts.map(lambda (w, c): Row(word=w, count=c)) .toDF()) # Print schema wordsCounts.printSchema() # Output: As you can see, the count and word types have been inferred without having to explicitly define long and string types respectively.
方案2:无需Row的简化写法
如果不想导入Row,可以直接使用元组结合toDF()指定列名,代码更简洁:
%spark2.pyspark # 直接用元组转换并指定列名 wordsCounts = filteredWordCounts.toDF(["word", "count"]) # Print schema wordsCounts.printSchema() # Output: As you can see, the count and word types have been inferred without having to explicitly define long and string types respectively.
验证说明
运行修改后的代码,Spark将成功识别行结构,生成带有正确列名的DataFrame,并打印出自动推断的Schema。
内容的提问来源于stack exchange,提问作者Viacheslav
相关产品推荐
相关产品推荐

