Eda*_*ame 12 python hadoop pandas apache-spark pyspark
我在Jupyter笔记本上使用pyspark.以下是Spark设置的方式:
import findspark
findspark.init(spark_home='/home/edamame/spark/spark-2.0.0-bin-spark-2.0.0-bin-hadoop2.6-hive', python_path='python2.7')
import pyspark
from pyspark.sql import *
sc = pyspark.sql.SparkSession.builder.master("yarn-client").config("spark.executor.memory", "2g").config('spark.driver.memory', '1g').config('spark.driver.cores', '4').enableHiveSupport().getOrCreate()
sqlContext = SQLContext(sc)
Run Code Online (Sandbox Code Playgroud)
然后,当我这样做:
spark_df = sqlContext.createDataFrame(df_in)
Run Code Online (Sandbox Code Playgroud)
哪里df_in是熊猫数据帧.然后我得到以下错误:
---------------------------------------------------------------------------
AttributeError Traceback (most recent call last)
<ipython-input-9-1db231ce21c9> in <module>()
----> 1 spark_df = sqlContext.createDataFrame(df_in)
/home/edamame/spark/spark-2.0.0-bin-spark-2.0.0-bin-hadoop2.6-hive/python/pyspark/sql/context.pyc in createDataFrame(self, data, schema, samplingRatio)
297 Py4JJavaError: ...
298 """
--> 299 return self.sparkSession.createDataFrame(data, schema, samplingRatio)
300
301 @since(1.3)
/home/edamame/spark/spark-2.0.0-bin-spark-2.0.0-bin-hadoop2.6-hive/python/pyspark/sql/session.pyc in createDataFrame(self, data, schema, samplingRatio)
520 rdd, schema = self._createFromRDD(data.map(prepare), schema, samplingRatio)
521 else:
--> 522 rdd, schema = self._createFromLocal(map(prepare, data), schema)
523 jrdd = self._jvm.SerDeUtil.toJavaArray(rdd._to_java_object_rdd())
524 jdf = self._jsparkSession.applySchemaToPythonRDD(jrdd.rdd(), schema.json())
/home/edamame/spark/spark-2.0.0-bin-spark-2.0.0-bin-hadoop2.6-hive/python/pyspark/sql/session.pyc in _createFromLocal(self, data, schema)
400 # convert python objects to sql data
401 data = [schema.toInternal(row) for row in data]
--> 402 return self._sc.parallelize(data), schema
403
404 @since(2.0)
AttributeError: 'SparkSession' object has no attribute 'parallelize'
Run Code Online (Sandbox Code Playgroud)
有谁知道我做错了什么?谢谢!
zer*_*323 16
SparkSession不是替代,SparkContext而是等同于SQLContext.只需使用它与您以前使用的方式相同SQLContext:
spark.createDataFrame(...)
Run Code Online (Sandbox Code Playgroud)
如果您必须访问SparkContextuse sparkContext属性:
spark.sparkContext
Run Code Online (Sandbox Code Playgroud)
因此,如果您需要SQLContext向后兼容性,您可以:
SQLContext(sparkContext=spark.sparkContext, sparkSession=spark)
Run Code Online (Sandbox Code Playgroud)
每当我们尝试从向后兼容的对象(例如 RDD)或 Spark 会话创建的数据帧创建 DF 时,您需要使 SQL 上下文感知会话和上下文。
就像前:
如果我创建一个 RDD:
ss=SparkSession.builder.appName("vivek").master('local').config("k1","vi").getOrCreate()
rdd=ss.sparkContext.parallelize([('Alex',21),('Bob',44)])
Run Code Online (Sandbox Code Playgroud)
但是如果我们希望从这个 RDD 创建一个 df,我们需要
sq=SQLContext(sparkContext=ss.sparkContext, sparkSession=ss)
那么只有我们可以将 SQLContext 与 pandas 创建的 RDD/DF 一起使用。
schema = StructType([
StructField("name", StringType(), True),
StructField("age", IntegerType(), True)])
df=sq.createDataFrame(rdd,schema)
df.collect()
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
28697 次 |
| 最近记录: |