Как я могу сгруппировать кадр данных Spark в строки размером не более 50 000 и не более 90 МБ?Python

Программы на Python
Anonymous
Как я могу сгруппировать кадр данных Spark в строки размером не более 50 000 и не более 90 МБ?

Сообщение Anonymous »

Как я могу сгруппировать фрейм данных 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

Вернуться в «Python»