【发布时间】:2018-11-21 07:53:56
【问题描述】:
我在 python 2.7 中编写了一个脚本,使用 pyspark 将 csv 转换为 parquet 和其他东西。 当我在一个小数据上运行我的脚本时,它运行良好,但是当我在一个更大的数据(250GB)上运行时,我迷上了以下错误——总分配超过了堆内存的 95.00%(960,285,889 字节)。 我怎么解决这个问题?它发生的原因是什么? tnx!
部分代码:
导入的库:
import pyspark as ps
from pyspark.sql.types import StructType, StructField, IntegerType,
DoubleType, StringType, TimestampType,LongType,FloatType
from collections import OrderedDict
from sys import argv
使用 pyspark:
schema_table_name="schema_"+str(get_table_name())
print (schema_table_name)
schema_file= OrderedDict()
schema_list=[]
ddl_to_schema(data)
for i in schema_file:
schema_list.append(StructField(i,schema_file[i]()))
schema=StructType(schema_list)
print schema
spark = ps.sql.SparkSession.builder.getOrCreate()
df = spark.read.option("delimiter",
",").format("csv").schema(schema).option("header", "false").load(argv[2])
df.write.parquet(argv[3])
# df.limit(1500).write.jdbc(url = url, table = get_table_name(), mode =
"append", properties = properties)
# df = spark.read.jdbc(url = url, table = get_table_name(), properties =
properties)
pq = spark.read.parquet(argv[3])
pq.show()
只是为了澄清 schema_table_name 是为了保存所有表名(在适合 csv 的 DDL 中)。
function ddl_to_schema 只需要一个常规的 ddl 并将其编辑为 parquet 可以使用的 ddl。
【问题讨论】:
-
给我们看一些代码...
-
将代码添加到问题中,而不是在 cmets 中
-
@Lorelorelore tnx!
-
看来您唯一的解决方案是不将整个文件读入内存。
-
@usr2564301 我的意思是也许这是一个标志,我应该增加在那里定义的数字......因为我看到了一个像“set.memory.driver”这样的命令,但我真的不知道tnx 回复!
标签: python csv pyspark heap-memory parquet