我读到具有相同分区器的 RDD 将位于同一位置。这对我很重要,因为我想加入几个未分区的大型 Hive 表。我的理论是,如果我可以将它们分区(通过字段调用 date_day)并位于同一位置,那么我将避免 shuffling 。
这是我试图为每张桌子做的事情:
def date_day_partitioner(key):
return (key.date_day - datetime.date(2017,05,01)).days
df = sqlContext.sql("select * from hive.table")
rdd = df.rdd
rdd2 = rdd.partitionBy(100, date_day_partitioner)
df2 = sqlContext.createDataFrame(rdd2, df_log_entry.schema)
print df2.count()
Run Code Online (Sandbox Code Playgroud)
不幸的是,我什至无法测试我关于协同定位和避免改组的理论,因为当我尝试 partitionBy 时出现以下错误:ValueError: too many values to unpack
Traceback (most recent call last):
File "/tmp/zeppelin_pyspark-118755547579363441.py", line 346, in <module>
raise Exception(traceback.format_exc())
Exception: Traceback (most recent call last):
File "/tmp/zeppelin_pyspark-118755547579363441.py", line 339, in <module>
exec(code)
File "<stdin>", line 15, in <module>
File "/usr/lib/spark/python/pyspark/sql/dataframe.py", line 380, in count …Run Code Online (Sandbox Code Playgroud)