【发布时间】: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