我有一个3节点的星系团。并尝试使用雪花火花连接器和jdbc驱动程序访问雪花。
jdbc:雪花-jdbc-3.12.4.jar火花-连接器:火花-雪花_2.11-2.7.0-火花_2.4.jar
这是我的代码:
sfOptions = {
"sfURL" : "{}.snowflakecomputing.com".format(ACCOUNT_NAME),
"sfUser" : "{}@fmr.com".format(USER_ID),
"sfAccount" : "{}".format(ACCOUNT_ID),
"sfRole" : "{}".format(DEFAULT_ROLE),
"sfAuthenticator" : "oauth",
"sfToken" : "{}".format(oauth_token),
"sfDatabase" : "{}".format(DATABASE),
"sfSchema" : "{}".format(SCHEMA),
"sfWarehouse" : "{}".format(WAREHOUSE)
}
SNOWFLAKE_SOURCE_NAME = "net.snowflake.spark.snowflake"
....
conf = (SparkConf()
.setMaster("spark://<master-url>")
.setAppName("Spark-Snowflake-Connector")
)
spark = (SparkSession.builder.config(conf=conf)
.enableHiveSupport()
.getOrCreate())
spark._jvm.net.snowflake.spark.snowflake.SnowflakeConnectorUtils.enablePushdownSession(spark._jvm.org.apache.spark.sql.SparkSession.builder().getOrCreate())
sdf = spark.read.format(SNOWFLAKE_SOURCE_NAME) \
.options(**sfOptions) \
.option("query", "select * from TIME_AGE") \
.load()
sdf.show()除了以下例外,我在sdf.show()上的调用失败了。有什么建议吗?
文件"/apps/shared/spark/python/lib/pyspark.zip/pyspark/sql/dataframe.py",第378行
20/04/26 09:54:55 INFO DAGScheduler: Job0失败: showString at NativeMethodAccessorImpl.java:0,采取了5.494100 s的回溯(最近一次调用):文件“/fedata/a 393831/雪花/火花-驱动器”,第114行在show "/apps/shared/spark/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py",第1257行,在call File "/apps/shared/spark/python/lib/pyspark.zip/pyspark/sql/utils.py",第63行,在deco "/apps/shared/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py",第328行中,在get_return_value py4j.protocol.Py4JJavaError中:调用o68.showString时出错。::org.apache.spark.SparkException:由于阶段失败而中止作业:阶段0.0中的任务0失败4次,最近的失败:阶段0.0中的任务0.3失败(TID 3,10.240.62.46,执行者0):net.snowflake.client.core.SFArrowResultSet.getObject(SFArrowResultSet.java:570) at net.snowflake.client.jdbc.SnowflakeResultSetV1.getObject(SnowflakeResultSetV1.java:336) at net.snowflake.spark.snowflake.io.ResultIterator$$anonfun$2.apply(SnowflakeResultSetRDD.scala:115) at net.snowflake.spark.snowflake.io.ResultIterator$$anonfun$2.apply(SnowflakeResultSetRDD.scala:114) at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)的java.lang.NullPointerException在scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234) at scala.collection.immutable.Range.foreach(Range.scala:160) at scala.collection.TraversableLike$class.map(TraversableLike.scala:234) at scala.collection.AbstractTraversable.map(Traversable.scala:104) at net.snowflake.spark.snowflake.io.ResultIterator.next(SnowflakeResultSetRDD.scala:114) at scala.collection.Iterator$$anon$11.next(Iterator.scala:410) at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:256) at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247) at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836) at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$24.apply(RDD.scala:836) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(( org.apache.spark.rdd.RDD.iterator(RDD.scala:288) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90) at org.apache.spark.scheduler.Task.run(Task.scala:121) at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
)
发布于 2020-04-27 23:03:19
看起来雪花JDBC 3.12.4 jar版本出现了问题,当使用Spark连接器火花-雪花_2.11-2.7.0-spark_2.4.jar时,您可以尝试使用3.12.3版本的雪花JDBC驱动程序吗?这与上面的星火连接器版本很好。
发布于 2020-04-27 07:56:11
对于相同的连接器和驱动程序配置,我也有相同的问题。我的查询只是计算SF示例表- snowflake_sample_data.tpch_sf1.lineitem上的行数。
"sfDatabase" -> "snowflake_sample_data",
"sfSchema" -> "tpch_sf1",
"query" -> "select count(*) from lineitem"因此,我刚刚试用了3.12.0版本的jdbc驱动程序,它可以工作。因此,似乎在驱动程序新版本中出现了倒退。
https://stackoverflow.com/questions/61443127
复制相似问题