Как я могу сгруппировать фрейм данных Spark в строки размером не более 50 000 и не более 90 МБ.
Я пробовал следующее, но иногда я получаю разделы размером более 90 МБ.
from pyspark.sql.window import Window
from pyspark.sql.functions import expr, length, sum as sql_sum, row_number, col, count
PARTITION_MB = 90
ROW_LIMIT = 50000
try:
# Select required columns
sdf = spark.table("table_name")
# Calculate total memory usage per row in MB
sdf = sdf.withColumn("json_string", expr("to_json(struct(*))")) \
.withColumn("memory_usage_per_row_MB", length("json_string") / 1024 / 1024) \
.drop("json_string")
# Generate row_number for ordering purposes
window_spec = Window.orderBy(expr("monotonically_increasing_id()"))
# Add row_number column to keep track of row order
sdf = sdf.withColumn("row_number", row_number().over(window_spec))
# Calculate cumulative memory usage
sdf = sdf.withColumn("cumulative_memory_MB", sql_sum("memory_usage_per_row_MB").over(window_spec)) \
.withColumn("cumulative_row_count", row_number().over(window_spec))
# Assign partition id based on memory and row limits
sdf = sdf.withColumn(
"partition_id",
expr(f"""
greatest(
floor(cumulative_memory_MB / {PARTITION_MB}),
floor((row_number - 1) / {ROW_LIMIT})
)
""")
)
# Validate partitions
partition_counts = sdf.groupBy("partition_id").agg(
sql_sum("memory_usage_per_row_MB").alias("partition_memory_MB"),
count("*").alias("row_count")
)
# Count the number of distinct partitions
num_partitions = sdf.select("partition_id").distinct().count()
# Assert that all partitions meet both memory and row constraints
assert partition_counts.filter(col("partition_memory_MB")
Подробнее здесь: https://stackoverflow.com/questions/790 ... -more-than