【问题标题】:Pyspark - Index from monotonically_increasing_id changes after list aggregationPyspark - 列表聚合后来自 monotonically_increasing_id 的索引更改
【发布时间】: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在聚合前有601295546156012955461660129554617三个ID,但聚合后数字变为601295547426012955474360129554744

为什么?我无法想象这应该发生。 monotonically_increasing_id() 的结果不就是一个简单的 long,在创建后保持它的值吗?

编辑:正如预期的那样,解决方法是在创建 id 之前coalesce(1) DataFrame。

【问题讨论】:

    标签: apache-spark pyspark


    【解决方案1】:

    dgdf2 是两个不同的数据帧,每个都有自己的 DAG。当调用其中一个数据帧上的操作时,这些 DAG 将彼此独立地执行。因此,每次调用 show() 时,都会评估相应数据帧的 DAG,并在评估期间调用 f.monotonically_increasing_id()

    为防止f.monotonically_increasing_id() 被调用两次,您可以在withColumn 转换后添加cache

    dg = dg.sort("B").withColumn("ID0", f.monotonically_increasing_id()).cache()
    

    使用缓存,f.monotonically_increasing_id() 的第一次评估结果被缓存并在评估第二个数据帧时重复使用。

    【讨论】:

    • 感谢您的解释,解决了很多困惑。
    猜你喜欢
    • 2013-03-30
    • 2021-06-10
    • 2018-05-07
    • 1970-01-01
    • 1970-01-01
    • 2013-10-04
    • 1970-01-01
    • 2015-08-11
    相关资源
    最近更新 更多