【发布时间】:2021-05-12 00:16:47
【问题描述】:
我正在使用 Pyspark 3.1.1 中的 monotonically_increasing_id() 函数创建索引。
我知道该函数的具体特征,但它们没有解释我的问题。
创建索引后,我对创建的索引应用collect_list() 函数进行简单聚合。
如果我比较结果,索引在某些情况下会发生变化,即在输入数据不太小的情况下,特别是在远程的上端。
完整示例代码:
import random
import string
from pyspark.sql import SparkSession
from pyspark.sql import functions as f
from pyspark.sql.types import StructType, StructField, StringType
spark = SparkSession.builder\
.appName("test")\
.master("local")\
.config('spark.sql.shuffle.partitions', '8')\
.getOrCreate()
# Create random input data of around length 100000:
input_data = []
ii = 0
while ii <= 100000:
L = random.randint(1, 3)
B = ''.join(random.choices(string.ascii_uppercase, k=5))
for i in range(L):
C = random.randint(1,100)
input_data.append((B,))
ii += 1
# Create Spark DataFrame:
input_rdd = sc.parallelize(tuple(input_data))
schema = StructType([StructField("B", StringType())])
dg = spark.createDataFrame(input_rdd, schema=schema)
# Create id and aggregate:
dg = dg.sort("B").withColumn("ID0", f.monotonically_increasing_id())
dg2 = dg.groupBy("B").agg(f.collect_list("ID0"))
Output:
dg.sort('B', ascending=False).show(10, truncate=False)
dg2.sort('B', ascending=False).show(5, truncate=False)
这当然会在每次运行时创建不同的数据,但如果长度足够大(问题在 10000 处出现轻微,但在 1000 处不出现),它应该每次都出现。这是一个示例结果:
+-----+-----------+
|B |ID0 |
+-----+-----------+
|ZZZVB|60129554616|
|ZZZVB|60129554617|
|ZZZVB|60129554615|
|ZZZUH|60129554614|
|ZZZRW|60129554612|
|ZZZRW|60129554613|
|ZZZNH|60129554611|
|ZZZNH|60129554609|
|ZZZNH|60129554610|
|ZZZJH|60129554606|
+-----+-----------+
only showing top 10 rows
+-----+---------------------------------------+
|B |collect_list(ID0) |
+-----+---------------------------------------+
|ZZZVB|[60129554742, 60129554743, 60129554744]|
|ZZZUH|[60129554741] |
|ZZZRW|[60129554739, 60129554740] |
|ZZZNH|[60129554736, 60129554737, 60129554738]|
|ZZZJH|[60129554733, 60129554734, 60129554735]|
+-----+---------------------------------------+
only showing top 5 rows
条目ZZZVB在聚合前有60129554615、60129554616和60129554617三个ID,但聚合后数字变为60129554742、60129554743、60129554744。
为什么?我无法想象这应该发生。 monotonically_increasing_id() 的结果不就是一个简单的 long,在创建后保持它的值吗?
编辑:正如预期的那样,解决方法是在创建 id 之前coalesce(1) DataFrame。
【问题讨论】:
标签: apache-spark pyspark