【问题标题】:Creating parquet files in spark with row-group size that is less than 100在 spark 中创建行组大小小于 100 的镶木地板文件
【发布时间】:2018-01-09 22:51:14
【问题描述】:

我有一个包含少量字段的 spark 数据框。一些字段是巨大的二进制 blob。整行的大小约为 50 MB。

我将数据框保存为镶木地板格式。我正在使用parquet.block.size 参数控制行组的大小。

Spark 会生成一个 parquet 文件,但是我总是会在一个行组中获得至少 100 行。这对我来说是个问题,因为块大小可能会变成千兆字节,这不适用于我的应用程序。

parquet.block.size 只要大小足以容纳超过 100 行,就可以正常工作。

我将InternalParquetRecordWriter.java 修改为MINIMUM_RECORD_COUNT_FOR_CHECK = 2,这解决了这个问题,但是我找不到支持调整这个硬编码常量的配置值。

是否有不同/更好的方法来获得小于 100 的行组大小?

这是我的代码的 sn-p:

from pyspark import Row
from pyspark.sql import SparkSession
import numpy as np

from pyspark.sql.types import StructType, StructField, BinaryType


def fake_row(x):
    result = bytearray(np.random.randint(0, 127, (3 * 1024 * 1024 / 2), dtype=np.uint8).tobytes())
    return Row(result, result)

spark_session = SparkSession \
    .builder \
    .appName("bbox2d_dataset_extraction") \
    .config("spark.driver.memory", "12g") \
    .config("spark.executor.memory", "4g")

spark_session.master('local[5]')

spark = spark_session.getOrCreate()
sc = spark.sparkContext
sc._jsc.hadoopConfiguration().setInt("parquet.block.size", 8 * 1024 * 1024)

index = sc.parallelize(range(50), 5)
huge_rows = index.map(fake_row)
schema = StructType([StructField('f1', BinaryType(), False), StructField('f2', BinaryType(), False)])

bbox2d_dataframe = spark.createDataFrame(huge_rows, schema).coalesce(1)
bbox2d_dataframe. \
    write.option("compression", "none"). \
    mode('overwrite'). \
    parquet('/tmp/huge/')

【问题讨论】:

标签: hadoop apache-spark parquet


【解决方案1】:

很遗憾,我还没有找到这样做的方法。我报告了this issue 以删除硬编码值并使其可配置。如果你有兴趣,我有一个补丁。

【讨论】:

【解决方案2】:

虽然PARQUET-409 尚未修复,但有几种变通方法可以使应用程序使用100 硬编码的每个行组的最小记录数。

第一个问题和解决方法: 您提到您的行可能高达 50Mb。 这给出了大约 5Gb 的行组大小。 同时,您的 spark 执行器只有 4Gb (spark.executor.memory)。 使其明显大于最大行组大小。
我建议spark.executor.memory 使用 12-20Gb 的大型 spark 执行器内存。玩这个,看看哪一个适用于您的数据集。 我们的大多数生产作业都使用此范围内的 spark executor 内存运行。 为使其适用于如此大的行组,您可能还需要将spark.executor.cores 调低至 1,以确保每个执行程序进程一次只占用一个如此大的行组。 (以损失一些 Spark 效率为代价)也许尝试将 spark.executor.cores 设置为 2 - 这可能需要将 spark.executor.memory 增加到 20-31Gb 范围。 (尽量保持under 32Gb,因为 jvm 切换到非压缩 OOP,这可能会占用 50% 的内存)

第二个问题和解决方法:如此大的 5Gb 行块很可能分布在许多 HDFS 块中,因为默认 HDFS 块在 128-256Mb 范围内。 (我假设你使用 HDFS 来存储那些 parquet 文件,因为你有“hadoop”标签) Parquet best practice 是一个行组完全驻留在一个 HDFS 块中:

行组大小:更大的行组允许更大的列块 可以进行更大的顺序 IO。更大的团体也 在写入路径中需要更多缓冲(或两次写入)。我们 推荐大行组(512MB - 1GB)。由于整个行组 可能需要读取,我们希望它完全适合一个 HDFS 块。 因此,HDFS 块大小也应该设置得更大。一个 优化的读取设置为:1GB 行组,1GB HDFS 块大小,1 每个 HDFS 文件的 HDFS 块。

以下是如何更改 HDFS 块大小的示例(在您创建此类 parquet 文件之前设置):

sc._jsc.hadoopConfiguration().set("dfs.block.size", "5g")

或在 Spark Scala 中:

sc.hadoopConfiguration.set("dfs.block.size", "5g")

我希望有时会在 Parquet 级别修复此问题,但是这两种解决方法应该允许您使用 Parquet 到如此大的行组。

【讨论】:

    猜你喜欢
    • 2017-06-27
    • 2021-02-28
    • 1970-01-01
    • 1970-01-01
    • 2018-04-05
    • 2022-12-01
    • 2020-12-14
    • 2016-10-07
    • 2020-09-09
    相关资源
    最近更新 更多