【问题标题】:Create Spark DataFrame from Pandas DataFrame从 Pandas DataFrame 创建 Spark DataFrame
【发布时间】:2019-07-08 22:35:55
【问题描述】:

我正在尝试从一个简单的 Pandas DataFrame 构建一个 Spark DataFrame。这是我遵循的步骤。

import pandas as pd
pandas_df = pd.DataFrame({"Letters":["X", "Y", "Z"]})
spark_df = sqlContext.createDataFrame(pandas_df)
spark_df.printSchema()

到目前为止一切正常。输出是:


|-- 字母:字符串(可为空=真)

当我尝试打印 DataFrame 时出现问题:

spark_df.show()

这是结果:

调用 o158.collectToPython 时出错。 : org.apache.spark.SparkException:作业因阶段失败而中止: 阶段 5.0 中的任务 0 失败 1 次,最近一次失败:丢失任务 0.0 在 5.0 阶段(TID 5、本地主机、执行程序驱动程序): org.apache.spark.SparkException:
来自 python 工作者的错误:
错误 执行 Jupyter 命令 'pyspark.daemon': [Errno 2] No such file 或 目录 PYTHONPATH 是:
/home/roldanx/soft/spark-2.4.0-bin-hadoop2.7/python/lib/pyspark.zip:/home/roldanx/soft/spark-2.4.0-bin-hadoop2.7/python/lib/ py4j-0.10.7-src.zip:/home/roldanx/soft/spark-2.4.0-bin-hadoop2.7/jars/spark-core_2.11-2.4.0.jar:/home/roldanx/soft/ spark-2.4.0-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip:/home/roldanx/soft/spark-2.4.0-bin-hadoop2.7/python/: org.apache.spark.SparkException:pyspark.daemon 中没有端口号 标准输出

这些是我的 Spark 规格:

SparkSession - 蜂巢

SparkContext

火花用户界面

版本: v2.4.0

大师: 本地[*]

应用名称: PySparkShell

这是我的 venv:

导出 PYSPARK_PYTHON=jupyter

导出 PYSPARK_DRIVER_PYTHON_OPTS='lab'

事实:

正如错误所述,它与从 Jupyter 运行 pyspark 有关。使用 'PYSPARK_PYTHON=python2.7' 和 'PYSPARK_PYTHON=python3.6' 运行它可以正常工作

【问题讨论】:

    标签: python pandas pyspark apache-spark-sql


    【解决方案1】:

    导入并初始化 findspark,创建一个 spark 会话,然后使用该对象将 pandas 数据帧转换为 spark 数据帧。然后将新的 spark 数据框添加到目录中。使用 python 3.6.6 在 Jupiter 5.7.2 和 Spyder 3.3.2 中测试和运行。

    import findspark
    findspark.init()
    
    import pyspark
    from pyspark.sql import SparkSession
    import pandas as pd
    
    # Create a spark session
    spark = SparkSession.builder.getOrCreate()
    
    # Create pandas data frame and convert it to a spark data frame 
    pandas_df = pd.DataFrame({"Letters":["X", "Y", "Z"]})
    spark_df = spark.createDataFrame(pandas_df)
    
    # Add the spark data frame to the catalog
    spark_df.createOrReplaceTempView('spark_df')
    
    spark_df.show()
    +-------+
    |Letters|
    +-------+
    |      X|
    |      Y|
    |      Z|
    +-------+
    
    spark.catalog.listTables()
    Out[18]: [Table(name='spark_df', database=None, description=None, tableType='TEMPORARY', isTemporary=True)]
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-10-20
      • 2021-11-26
      • 1970-01-01
      • 1970-01-01
      • 2019-02-10
      • 1970-01-01
      • 1970-01-01
      • 2017-03-17
      相关资源
      最近更新 更多