相关疑难解决方法(0)

来自 spark 数据帧的块 topandas

我有一个包含 1000 万条记录和 150 列的 spark 数据框。我正在尝试将其转换为熊猫 DF。

x = df.toPandas()
# do some things to x
Run Code Online (Sandbox Code Playgroud)

它失败了ordinal must be >= 1。我假设这是因为一次处理太大了。是否可以将其分块并将其转换为每个块的熊猫 DF?

全栈:

ValueError                                Traceback (most recent call last)
<command-2054265283599157> in <module>()
    158 from db.table where snapshot_year_month=201806""")
--> 159 ps = x.toPandas()
    160 # ps[["pol_nbr",
    161 # "pol_eff_dt",

/databricks/spark/python/pyspark/sql/dataframe.py in toPandas(self)
   2029                 raise RuntimeError("%s\n%s" % (_exception_message(e), msg))
   2030         else:
-> 2031             pdf = pd.DataFrame.from_records(self.collect(), columns=self.columns)
   2032 
   2033             dtype = {}

/databricks/spark/python/pyspark/sql/dataframe.py in collect(self)
    480         with SCCallSiteSync(self._sc) as …
Run Code Online (Sandbox Code Playgroud)

python pandas apache-spark

6
推荐指数
1
解决办法
3017
查看次数

使用toPandas()方法将Spark数据帧转换为Pandas数据帧时会发生什么

我有一个spark数据框,可以使用将其转换为pandas数据框。

toPandas()
Run Code Online (Sandbox Code Playgroud)

pyspark中可用的方法。

我对此有以下疑问?

  1. 这种转换是否违反了使用spark本身(分布式计算)的目的?
  2. 数据集将是巨大的,那么速度和内存问题呢?
  3. 如果有人也可以解释,那么这一行代码究竟会发生什么,那将真正有帮助。

谢谢

python pandas apache-spark pyspark pyspark-sql

3
推荐指数
1
解决办法
667
查看次数

将 Spark Structure Streaming DataFrame 转换为 Pandas DataFrame

我设置了一个 Spark Streaming 应用程序,它从 Kafka 主题进行消费,我需要使用一些接受 Pandas Dataframe 的 API,但是当我尝试转换它时,我得到了这个

: org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start();;
kafka
        at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.org$apache$spark$sql$catalyst$analysis$UnsupportedOperationChecker$$throwError(UnsupportedOperationChecker.scala:297)
        at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:36)
        at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$$anonfun$checkForBatch$1.apply(UnsupportedOperationChecker.scala:34)
        at org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:127)
        at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.checkForBatch(UnsupportedOperationChecker.scala:34)
        at org.apache.spark.sql.execution.QueryExecution.assertSupported(QueryExecution.scala:63)
        at org.apache.spark.sql.execution.QueryExecution.withCachedData$lzycompute(QueryExecution.scala:74)
        at org.apache.spark.sql.execution.QueryExecution.withCachedData(QueryExecution.scala:72)
        at org.apache.spark.sql.execution.QueryExecution.optimizedPlan$lzycompute(QueryExecution.scala:78)
        at org.apache.spark.sql.execution.QueryExecution.optimizedPlan(QueryExecution.scala:78)
        at org.apache.spark.sql.execution.QueryExecution.completeString(QueryExecution.scala:219)
        at org.apache.spark.sql.execution.QueryExecution.toString(QueryExecution.scala:202)
        at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:62)
        at org.apache.spark.sql.Dataset.withNewExecutionId(Dataset.scala:2832)
        at org.apache.spark.sql.Dataset.collectToPython(Dataset.scala:2809)
        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
        at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
        at py4j.Gateway.invoke(Gateway.java:282)
        at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
        at py4j.commands.CallCommand.execute(CallCommand.java:79)
        at py4j.GatewayConnection.run(GatewayConnection.java:238)
        at java.lang.Thread.run(Thread.java:745)
Run Code Online (Sandbox Code Playgroud)

这是我的Python代码

spark = SparkSession\
    .builder\ …
Run Code Online (Sandbox Code Playgroud)

python pandas apache-spark pyspark spark-structured-streaming

3
推荐指数
1
解决办法
1884
查看次数

Dataframe.toPandas 总是在驱动程序节点上还是在工作程序节点上?

想象一下,您正在通过 SparkContext 和 Hive 加载一个大型数据集。所以这个数据集然后分布在你的 Spark 集群中。例如,对数千个变量的观察(值 + 时间戳)。

现在您将使用一些 map/reduce 方法或聚合来组织/分析您的数据。例如按变量名称分组。

分组后,您可以获得每个变量的所有观察值(值)作为时间序列数据框。如果您现在使用 DataFrame.toPandas

def myFunction(data_frame):
   data_frame.toPandas()

df = sc.load....
df.groupBy('var_name').mapValues(_.toDF).map(myFunction)
Run Code Online (Sandbox Code Playgroud)
  1. 这是在每个工作节点上转换为 Pandas 数据帧(每个变量),还是
  2. Pandas 数据帧是否总是在驱动程序节点上,因此数据从工作节点传输到驱动程序?

python hadoop pandas apache-spark pyspark

2
推荐指数
1
解决办法
1851
查看次数

2
推荐指数
1
解决办法
1万
查看次数