【发布时间】:2019-08-22 16:34:51
【问题描述】:
使用 sqlContext.sql 查询在 EMR 上运行 pyspark 代码。 其中一个查询导致引发 driver.maxResultSize 相关错误。 尝试对查询产生的数据框使用解释以了解原因。在那里,我看到 spark 出于某种原因使用带有嵌套连接的广播(没有明确的说明)。 我想了解:
1) 为什么 spark 使用广播和嵌套连接来执行这个查询?
2) 为什么广播要经过驱动程序?
3) 如何重写我的代码,使 spark 不使用广播(因为广播,或者它通过驱动程序,似乎是问题的根源)?
导致问题的查询:
df1.createOrReplaceTempView("temp_df_sql_view1")
df2.createOrReplaceTempView("temp_df_sql_view2")
# Get values from df1 that exist only in df1
df = sqlContext.sql("""SELECT * FROM temp_df_sql_view1 WHERE id NOT IN (SELECT id FROM temp_df_sql_view2)""")
df.explain()
我得到的错误信息是:Total size of serialized results of 79 tasks (2.1 GB) is bigger than spark.driver.maxResultSize (2.0 GB) 尽管 driver.maxResultSize 曾经是 1g,但为了修复错误而被放大。但是,结果的总大小似乎也随之扩大了。
在意识到这可能是一个广播问题后,我禁用了自动广播:
conf = SparkConf()
# This should've disabled auto-broadcast
conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
sc = SparkContext.getOrCreate(conf=conf)
sqlContext = SQLContext(sc)
但是在 df 上使用 explain() 仍然显示相同的以下计划(包括广播):
BroadcastNestedLoopJoin BuildLeft, LeftAnti, ((id#22 = id#19) || isnull((id#22 = id#19)))
:- BroadcastExchange IdentityBroadcastMode
: +- *(1) FileScan parquet [id#22,data#23] Batched: false, Format: Parquet, Location: InMemoryFileIndex[s3://bucket], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<...
+- Generate explode(id#2), false, [id#19]
+- *(2) Scan ElasticsearchRelation(Map(...,org.apache.spark.sql.SQLContext@6caa1e7e,None) [id#2] PushedFilters: [], ReadSchema: struct<id:array<string>>
None
【问题讨论】:
-
如何在
conf.set("spark.sql.autoBroadcastJoinThreshold", -1)中创建conf 变量? -
@moriarty007 定义 conf = SparkConf() 然后用它来创建 sqlContext
标签: apache-spark pyspark amazon-emr broadcast