在pyspark中加载比内存hdf5文件更大的文件

pyg*_*iel 5 python hdf5 apache-spark pyspark

我有一个以HDF5格式存储的大文件(例如20 Gb)。该文件基本上是一组随时间变化的3D坐标(分子模拟轨迹)。这基本上是形状的数组(8000 (frames), 50000 (particles), 3 (coordinates))

在普通的python中,我只需要使用for h5py或加载hdf5数据文件,pytables并为该数据文件建立索引,就好像它是一个numpy一样(该库会延迟加载所需的任何数据)。

但是,如果我尝试使用SparkContext.parallelize它在Spark中加载此文件,显然会阻塞内存:

sc.parallelize(data, 10)
Run Code Online (Sandbox Code Playgroud)

我该如何解决这个问题?大型数组有首选的数据格式吗?我可以将rdd写入磁盘而不经过内存吗?

vvl*_*rov 5

Spark(和Hadoop)不支持读取HDF5二进制文件的某些部分。(我怀疑这样做的原因是HDF5是用于存储文档的容器格式,它允许为文档指定树状层次结构)。

但是,如果您需要从本地磁盘读取文件,Spark可以使用它,尤其是在您知道HDF5文件的内部结构的情况下。

这是一个示例 -假设您将运行本地Spark作业,并且您预先知道HDF5数据集“ / mydata”由100个块组成。

h5file_path="/absolute/path/to/file"

def readchunk(v):
    empty = h5.File(h5file_path)
    return empty['/mydata'][v,:]

foo = sc.parallelize(range(0,100)).map(lambda v: readchunk(v))
foo.count()
Run Code Online (Sandbox Code Playgroud)

更进一步,您可以修改程序以使用以下命令检测块数 f5['/mydata'].shape[0]

下一步将是遍历多个数据集(您可以使用列出数据集f5.keys())。

另外还有另一篇文章“从HDF5数据集到Apache Spark RDD”描述了类似的方法。

相同的方法在分布式集群上也可以使用,但是效率不高。h5py要求文件位于本地文件系统上。因此,可以通过以下几种方法实现:将文件复制到所有工作程序,并将其保存在工作程序磁盘上的同一位置;或将文件放入HDFS并使用保险丝安装HDFS-这样工作人员就可以访问该文件。两种方法都存在一些效率低下的问题,但是对于临时任务来说应该足够好了。

这是优化的版本,在每个执行程序上仅打开一次h5:

h5file_path="/absolute/path/to/file"

_h5file = None    
def readchunk(v):
    # code below will be executed on executor - in another python process on remote server
    # original value for _h5file (None) is sent from driver
    # and on executor is updated to h5.File object when the `readchunk` is called for the first time
    global _h5file
    if _h5file is None:
         _h5file = h5.File(h5file_path)
    return _h5file['/mydata'][v,:]

foo = sc.parallelize(range(0,100)).map(lambda v: readchunk(v))
foo.count()
Run Code Online (Sandbox Code Playgroud)