Loading NumPy arrays from disk: mmap() vs. Zarr/HDF5
pythonspeed.com
pythonspeed.com
I can't figure out if one of those three methods can do what I want. I want to include a multi-valued lookup table so given a string id I can find the 0 or more corresponding records. (I can handle single-valued lookup tables, if multi-valued isn't available.)
I looked at Parquet, but it doesn't seem to handle id->index table lookups. Or 256 byte fields. Do any of Zarr, HDF5 and TileDB offer a better fit?
Background: it's hard to compare molecules directly. Instead, there are hashing methods to turn a molecular structure into a fixed-length bitstring of length nearly always between 166 and 4096 bits in length.
Bitstrings are easy to compare using Jaccard similarity, which gives a decent-enough and very fast way to compare molecular similarity. With help from another HN user (who I paid), we were able to get the performance to within a factor of 2 of the theoretical memory bandwidth.
Going faster requires compressed RAM->uncompressed cache transfer, better indexing, and distributed computing. I don't know how to do it using one of the existing Big Data formats, and developing my own format is hard.
My FPB file format is a FourCC format with several chunks. One chunk stores all N fingerprints, cache size aligned. These are ordered by popcount, so in principle some should be easily compressible - I haven't tried, and can't figure out how BLOSC is supposed to be used.
Another chunk stores the corresponding fingerprint identifiers, UTF-8/non-fixed-length strings, in the same order (fingerprint[i] corresponds to id[i]).
A third chunk stores a multi-valued hash table so I can look up fingerprint(s) by name, instead of by index. The hash table is a variant of the djb's "cdb" hash table. This is a 32-bit hash, which means for large data sets I get many collisions. (I designed it for <100M entries. Some people now want to handle >1B entries.)
I'm hoping to use an existing format rather than roll my own. Again.
https://blog.rstudio.com/2016/03/29/feather/
http://www.fstpackage.org/ https://github.com/fstpackage/fstlib
I had the same issue with fstlib - how do I handle id lookups?
Any pointers?
For fst, I only played with the R interface, which would be called like this to retrieve row 12345 from the "fingerprint" column:
read_fst("library.fst", columns="fingerprint", from=12345, to=12345)
However fst didn't offer a raw/binary column datatype last I checked, which is frustrating. It has chr (string) but R can't have embedded NUL bytes in strings, so that was a dead-end for efficient storage of binary fingerprints for R. I didn't check if the underlying fstlib structures accept NULs in string columns.[1] https://feedback.tiledb.com/tiledb-core/p/support-axes-label...
[2] https://docs.tiledb.com/developer/dask/quickstart
[3] https://docs.tiledb.com/developer/spark/quickstart
Disclosure: I am a member of the TileDB team.
The big advantage of zarray is when reading/writing to a cloud bucket, where having a single file containing chunks/indicies is impractical and you'd rather work with direct references to objects in a bucket.
Otherwise, HDF5 offers every single advantage that zarray has and is much more mature, stable, better documented, and has better support. If you're working in a situation where it makes sense to have things on disk or some sort of NFS share, use HDF5. If you're working with objects in a cloud bucket, you'll incur additional overhead with HDF5, as you'll have to read its table of indices, then make range requests to each chunk. Zarr is optimized for the cloud use case.
I recommended Zarr on the basis that threading seems to work much better and compression options are better.
pytables might make it possible though, haven't looked.
It's more that I keep hearing that statement ("Zarr is chunked and HDF5 isn't") in the wild a lot.
As for threading, yeah, HDF5 can be a bit cumbersome there.
I'm not sure I'd agree on compression, though. HDF5 supports fully arbitrary compression, after all, it's just the client reading the data also needs to have the compression filter you're using installed. Zarr is often used with BLOSC, and that was originally developed specifically for HDF5, FWIW.
My impression from the talk linked in the article was that HDF5 doesn't do BLOSC by default, and that in general it's much easier to plug new things in, add caching, etc..
Absolutely not. HDF5 is an awful format with terrible implementations. For example, try writing a python program with multiple threads where each thread writes to a different HDF5 file. This should just work -- there's no concurrent access. And yet it doesn't because HDF5 implementations are piles of ancient C code that use lots of global state. There's no technical reason for this; one could easily store all the state needed in a per-file object. But back in the day, software eng standards were lower (especially for scientists) and HDF5 changes at a glacial place.
I've been bitten by this particular bug, but you really have to wonder: given how poorly it speaks to the software engineering behind HDF5 implementations, what else is broken in the code or specifications?
If you're working in a situation where it makes sense to have things on disk or some sort of NFS share, use HDF5. If you're working with objects in a cloud bucket, you'll incur additional overhead with HDF5, as you'll have to read its table of indices, then make range requests to each chunk. Zarr is optimized for the cloud use case.
When last I looked, there were no open source HDF5 implementations that were smart enough to do range requests to cloud hosted hdf files. Has this changed?
And the library does quite a lot of work when you call into it; chunk lookup, decompression, and type conversions all happen behind that lock. You can use the "direct chunk access" functions (H5Dread_chunk?) to bypass a lot of that work and do it yourself, so you get back to using multiple threads again, and that can be a big win, but having to do it sucks, and I don't think h5py exposes this functionality at all.
Many python programs I write end up using 8+ cores on a single machine using either multiprocessing or C functions with released GIL.
[0] https://github.com/jjhelmus/pyfive
[1] https://h5s3.github.io/h5s3/python.html
[2] https://www.hdfgroup.org/solutions/hdf-kita/
[3] http://docs.h5py.org/en/stable/high/file.html#python-file-like-objectsEfficient access to scientific datasets hosted on S3/GCP is a full blown crisis in the scientific computing community. People aren't switching to zarr for the fun of it, but because zarr is here, today, and isn't a joke, and is actually open.
Pity, as it has on paper a lot of great concepts and features. Maybe it'll be mature enough someday, though my money is on something better from the ground up coming along.
Honestly, most of the portability advantage is moot nowadays. Chunk s3-like storage, smb, and ability to copy files from ext to ntfs (at least on nix) means that sharing your data across platforms isn't the struggle it used to be. Windows is rapidly becoming/already is a second class citizen in science-data heavy workflows.
I ended up going with a NAS and just file system primitives for my computer vision image workflow, works great.
https://stackoverflow.com/questions/35837243/hdf5-possible-d...
A glib high level overview of my last job for 6 years was "write out HDF5 files". In that time, I don't recall seeing a true data corruption problem with HDF5.
Now, I ran into many other problems with HDF5, typically surrounding the newer features that came along in 1.10, and its threading limitations. The older folks at that job would mention historical issues with data corruption (often from reading files as they're being written to), but I never saw it myself.
Docs: https://docs.tiledb.com/developer/
Website: https://tiledb.com/
Earlier HN post: https://news.ycombinator.com/item?id=15547749
Disclosure: I am a member of the TileDB team.
Taking a look at the compression filters, it seems that they are mostly suited for byte-oriented and integer data? The bit width reduction filter, for example, seems to act only on integers, and it doesn't look like any floating-point-specific compressors (like fpzip) are implemented -- just general-purpose ones like bzip2.
The webpage does not tell me exactly what TileDB is, so I'm having a hard time getting understanding what's underneath all the marketing mumbo-jumbo.
But if it is indeed a better SQLite/Parquet, DAMN, I'm so going to spread the gospel of this to all my students.
For a lot of what I do, I want a hierarchical containment system- the equivalent of folders with files. And the files themselves are leaves in the hierarchy, containing multidimensional array data. WOrks great when the arrays are composed of fairly straightforward payloads, like float[x][y][z] but also works if your array values are structs. Much of the value in zarr and tiledb comes from specifically how they arrange the arrays, for convenient read access to slices of the arrays. Access is going to look like: ages = root["user"]["age"][100:100:2]
Parquet is mostly a column file format, but with nesting. I'd use it to store large amounts of structured data with a relatively straightforward schema, although the schema itself can be fairly nested so some records have very complex structure. Access would often be in a loop over all records: for record in records: if record.user.has_age(): print("User age:", record.user.age)
SQLIte is a library/CLI that implements a relational database. It has a SQL interface and stores data using classic relational DB approaches, including secondary indices, etc, and permitting joins directly within the engine: SELECT age FROM user WHERE user.country == 'Bulgaria'
The good news is that we offer efficient integrations with MariaDB, PrestoDB and Spark, so you can directly process SQL queries on TileDB data via those engines (which work even for dense data). With MariaDB, we even have a embedded version which allows running SQL queries directly from Python[2]. This combines the ease of use of sqlite with MariaDB's speed and TileDB's fast access to AWS S3 (and Azure Blob Store in the next version).
[1] https://docs.tiledb.com/main/use-cases/dataframes
[2] https://docs.tiledb.com/developer/api-usage/embedded-sql
Disclosure: I am a member of the TileDB team.
Atomic access wise, that’s fine as long as you don’t map it through a network file system if any sort; also, $DIETY help you recover a consistent state if the power goes out or you kernel panic
Mmap is awesome, I use it all the time. But writing into mmaps robustly is hard. Best example to learn from I am familiar with is symas LMDB source code.
I'll take a look at the symas LMDB source code, thanks for pointing that out.
If nothing else, it's very useful to have a "don't use mmap" option in your code.
And for that Dask supports Zarr, and HDF5, and pre-existing arrays (where mmap() comes in), and a few other formats: https://docs.dask.org/en/latest/array-creation.html#
If the answer is "much less than the whole file", accessing the data over the network makes sense. Otherwise, you should almost certainly bring it to a local disk, where mmap shines.
You sure about this? You can mount S3 locally, and the filesystem should be transparent to the kernel.
1. FUSE filesystems. This means sending really slow queries to S3, and so you're much better off using compression to get better performance.
2. Syncing to filesystem with AWS DataSync. This would work, yes, but it's a S3 specific feature and not all blob storage systems have it.
As author of https://github.com/kahing/goofys/ I respectfully disagree :-)