如何在pyspark.sql中创建一个表作为选择 [英] How to create a table as select in pyspark.sql

查看:41
本文介绍了如何在pyspark.sql中创建一个表作为选择的处理方法,对大家解决问题具有一定的参考价值,需要的朋友们下面随着小编来一起学习吧!

问题描述

是否可以使用 select 语句在 spark 上创建表?

Is it possible to create a table on spark using a select statement?

我执行以下操作

import findspark
findspark.init()
import pyspark
from pyspark.sql import SQLContext

sc = pyspark.SparkContext()
sqlCtx = SQLContext(sc)

spark_df = sqlCtx.read.format('com.databricks.spark.csv').options(header='true', inferschema='true').load("./data/documents_topics.csv")
spark_df.registerTempTable("my_table")

sqlCtx.sql("CREATE TABLE my_table_2 AS SELECT * from my_table")

但我收到错误

/Users/user/anaconda/bin/python/Users/user/workspace/Outbrain-Click-Prediction/test.py 使用 Spark 的默认 log4j 配置文件:org/apache/spark/log4j-defaults.properties将默认日志级别设置为警告".调整日志级别使用sc.setLogLevel(newLevel).17/01/21 17:19:43 警告 NativeCodeLoader:无法为您的平台加载本机 Hadoop 库...使用适用的内置java类回溯(最近调用最后):文件"/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/utils.py",第 63 行,装饰中返回 f(*a, **kw) 文件/Users/user/spark-2.0.2-bin-hadoop2.7/python/lib/py4j-0.10.3-src.zip/py4j/protocol.py",第 319 行,在 get_return_value py4j.protocol.Py4JJavaError: An error调用 o19.sql 时发生.:org.apache.spark.sql.AnalysisException:未解析的运算符'CreateHiveTableAsSelectLogicalPlan CatalogTable( 表:my_table_2创建时间:美国东部时间 2017 年 1 月 21 日星期六 17:19:53 最后访问时间:12 月 31 日星期三18:59:59 EST 1969 类型:托管存储(输入格式:org.apache.hadoop.mapred.TextInputFormat, OutputFormat:org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false;;'CreateHiveTableAsSelectLogicalPlan CatalogTable( 表:my_table_2创建时间:美国东部时间 2017 年 1 月 21 日星期六 17:19:53 最后访问时间:12 月 31 日星期三18:59:59 EST 1969 类型:托管存储(输入格式:org.apache.hadoop.mapred.TextInputFormat, OutputFormat:org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false :+- 项目 [document_id#0, topic_id#1,confidence_level#2] : +- SubqueryAlias my_table : +-关系[document_id#0,topic_id#1,confidence_level#2] csv

/Users/user/anaconda/bin/python /Users/user/workspace/Outbrain-Click-Prediction/test.py Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties Setting default log level to "WARN". To adjust logging level use sc.setLogLevel(newLevel). 17/01/21 17:19:43 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable Traceback (most recent call last): File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/utils.py", line 63, in deco return f(*a, **kw) File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/lib/py4j-0.10.3-src.zip/py4j/protocol.py", line 319, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o19.sql. : org.apache.spark.sql.AnalysisException: unresolved operator 'CreateHiveTableAsSelectLogicalPlan CatalogTable( Table: my_table_2 Created: Sat Jan 21 17:19:53 EST 2017 Last Access: Wed Dec 31 18:59:59 EST 1969 Type: MANAGED Storage(InputFormat: org.apache.hadoop.mapred.TextInputFormat, OutputFormat: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false;; 'CreateHiveTableAsSelectLogicalPlan CatalogTable( Table: my_table_2 Created: Sat Jan 21 17:19:53 EST 2017 Last Access: Wed Dec 31 18:59:59 EST 1969 Type: MANAGED Storage(InputFormat: org.apache.hadoop.mapred.TextInputFormat, OutputFormat: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false : +- Project [document_id#0, topic_id#1, confidence_level#2] : +- SubqueryAlias my_table : +- Relation[document_id#0,topic_id#1,confidence_level#2] csv

在org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.failAnalysis(CheckAnalysis.scala:40)在org.apache.spark.sql.catalyst.analysis.Analyzer.failAnalysis(Analyzer.scala:58)在org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:374)在org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:67)在org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:126)在org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.checkAnalysis(CheckAnalysis.scala:67)在org.apache.spark.sql.catalyst.analysis.Analyzer.checkAnalysis(Analyzer.scala:58)在org.apache.spark.sql.execution.QueryExecution.assertAnalyzed(QueryExecution.scala:49)在 org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:64) 在org.apache.spark.sql.SparkSession.sql(SparkSession.scala:582) 在sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)在sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)在 java.lang.reflect.Method.invoke(Method.java:498) 在py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:237) 在py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) 在py4j.Gateway.invoke(Gateway.java:280) 在py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在py4j.GatewayConnection.run(GatewayConnection.java:214) 在java.lang.Thread.run(Thread.java:745)

at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.failAnalysis(CheckAnalysis.scala:40) at org.apache.spark.sql.catalyst.analysis.Analyzer.failAnalysis(Analyzer.scala:58) at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:374) at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$1.apply(CheckAnalysis.scala:67) at org.apache.spark.sql.catalyst.trees.TreeNode.foreachUp(TreeNode.scala:126) at org.apache.spark.sql.catalyst.analysis.CheckAnalysis$class.checkAnalysis(CheckAnalysis.scala:67) at org.apache.spark.sql.catalyst.analysis.Analyzer.checkAnalysis(Analyzer.scala:58) at org.apache.spark.sql.execution.QueryExecution.assertAnalyzed(QueryExecution.scala:49) at org.apache.spark.sql.Dataset$.ofRows(Dataset.scala:64) at org.apache.spark.sql.SparkSession.sql(SparkSession.scala:582) 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:237) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:280) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:214) at java.lang.Thread.run(Thread.java:745)

在处理上述异常的过程中,又发生了一个异常:

During handling of the above exception, another exception occurred:

回溯(最近一次调用最后一次):文件/Users/user/workspace/Outbrain-Click-Prediction/test.py",第 16 行,在sqlCtx.sql("CREATE TABLE my_table_2 AS SELECT * from my_table") 文件"/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/context.py",第 360 行,在 sql 中返回 self.sparkSession.sql(sqlQuery) 文件/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/session.py",第 543 行,在 sql 中返回 DataFrame(self._jsparkSession.sql(sqlQuery), self._wrapped) 文件"/Users/user/spark-2.0.2-bin-hadoop2.7/python/lib/py4j-0.10.3-src.zip/py4j/java_gateway.py",第 1133 行,在 call 文件中"/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/utils.py",第 69 行,装饰中引发 AnalysisException(s.split(': ', 1)[1], stackTrace) pyspark.sql.utils.AnalysisException:未解析的运算符'CreateHiveTableAsSelectLogicalPlan CatalogTable(\n\tTable:my_table_2\n\t创建时间:美国东部时间 2017 年 1 月 21 日星期六 17:19:53\n\t最后访问时间:1969 年 12 月 31 日星期三 18:59:59 EST\n\tType: MANAGED\n\tStorage(InputFormat:org.apache.hadoop.mapred.TextInputFormat, OutputFormat:org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)),false;;\n'CreateHiveTableAsSelectLogicalPlan CatalogTable(\n\tTable:my_table_2\n\t创建时间:美国东部时间 2017 年 1 月 21 日星期六 17:19:53\n\t最后访问时间:1969 年 12 月 31 日星期三 18:59:59 EST\n\tType: MANAGED\n\tStorage(InputFormat:org.apache.hadoop.mapred.TextInputFormat, OutputFormat:org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)),假\n:+- 项目 [document_id#0,topic_id#1,confidence_level#2]\n:+- 子查询别名 my_table\n:+-关系[document_id#0,topic_id#1,confidence_level#2] csv\n"

Traceback (most recent call last): File "/Users/user/workspace/Outbrain-Click-Prediction/test.py", line 16, in sqlCtx.sql("CREATE TABLE my_table_2 AS SELECT * from my_table") File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/context.py", line 360, in sql return self.sparkSession.sql(sqlQuery) File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/session.py", line 543, in sql return DataFrame(self._jsparkSession.sql(sqlQuery), self._wrapped) File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/lib/py4j-0.10.3-src.zip/py4j/java_gateway.py", line 1133, in call File "/Users/user/spark-2.0.2-bin-hadoop2.7/python/pyspark/sql/utils.py", line 69, in deco raise AnalysisException(s.split(': ', 1)[1], stackTrace) pyspark.sql.utils.AnalysisException: "unresolved operator 'CreateHiveTableAsSelectLogicalPlan CatalogTable(\n\tTable: my_table_2\n\tCreated: Sat Jan 21 17:19:53 EST 2017\n\tLast Access: Wed Dec 31 18:59:59 EST 1969\n\tType: MANAGED\n\tStorage(InputFormat: org.apache.hadoop.mapred.TextInputFormat, OutputFormat: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false;;\n'CreateHiveTableAsSelectLogicalPlan CatalogTable(\n\tTable: my_table_2\n\tCreated: Sat Jan 21 17:19:53 EST 2017\n\tLast Access: Wed Dec 31 18:59:59 EST 1969\n\tType: MANAGED\n\tStorage(InputFormat: org.apache.hadoop.mapred.TextInputFormat, OutputFormat: org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat)), false\n: +- Project [document_id#0, topic_id#1, confidence_level#2]\n: +- SubqueryAlias my_table\n: +- Relation[document_id#0,topic_id#1,confidence_level#2] csv\n"

推荐答案

我已经通过使用 HiveContext 而不是 SQLContext 更正了这个问题,如下所示:

I've corrected this issue by using HiveContext instead of SQLContext as below:

import findspark
findspark.init()
import pyspark
from pyspark.sql import HiveContext

sqlCtx= HiveContext(sc)

spark_df = sqlCtx.read.format('com.databricks.spark.csv').options(header='true', inferschema='true').load("./data/documents_topics.csv")
spark_df.registerTempTable("my_table")

sqlCtx.sql("CREATE TABLE my_table_2 AS SELECT * from my_table")

这篇关于如何在pyspark.sql中创建一个表作为选择的文章就介绍到这了,希望我们推荐的答案对大家有所帮助,也希望大家多多支持IT屋!

查看全文
登录 关闭
扫码关注1秒登录
发送“验证码”获取 | 15天全站免登陆