4 ms·
Hi! I'm the original author of Awkward Array (Jim Pivarski), though there are now many contributors with about five regulars. Two of my colleagues just pointed
by jpivarski 5y ago
Hi! I'm the original author of Awkward Array (Jim Pivarski), though there are now many contributors with about five regulars. Two of my colleagues just pointed me here—I'm glad you're interested! I can answer any questions you have about it.
First, sorry about all the TODOs in the documentation: I laid out a table of contents structure as a reminder to myself of what ought to be written, but haven't had a chance to fill in all of the topics. From the front page (https://awkward-array.org/ https://awkward-array.org/), if you click through to the Python API reference (https://awkward-array.readthedocs.io/ https://awkward-array.readthedocs.io/), that site is 100% filled in. Like NumPy, the library consists of one basic data type, `ak.Array`, and a suite of functions that act on it, `ak.this` and `ak.that`. All of those functions are individually documented, and many have examples.
The basic idea starts with a data structure like Apache Arrow (https://arrow.apache.org/)—a https://arrow.apache.org/)—a tree of general, variable-length types, organized in memory as a collection of columnar arrays—but performs operations on the data without ever taking it out of its columnar form. (3.5 minute explanation here: https://youtu.be/2NxWpU7NArk?t=661 https://youtu.be/2NxWpU7NArk?t=661) Those columnar operations are compiled (in C++); there's a core of structure-manipulation functions suggestively named "cpu-kernels" that will also be implemented in CUDA (some already have, but that's in an experimental stage).
A key aspect of this is that structure can be manipulated just by changing values in some internal arrays and rearranging the single tree organizing those arrays. If, for instance, you want to replace a bunch of objects in variable-length lists with another structure, it never needs to instantiate those objects or lists as explicit types (e.g. `struct` or `std::vector`), and so the functions don't need to be compiled for specific data types. You can define any new data types at runtime and the same compiled functions apply. Therefore, JIT compilation is not necessary.
We do have Numba extensions so that you can iterate over runtime-defined data types in JIT-compiled Numba, but that's a second way to manipulate the same data. By analogy with NumPy, you can compute many things using NumPy's precompiled functions, as long as you express your workflow in NumPy's vectorized way. Numba additionally allows you to express your workflow in imperative loops without losing performance. It's the same way with Awkward Array: unpacking a million record structures or slicing a million variable-length lists in a single function call makes use of some precompiled functions (no JIT), but iterating over them at scale with imperative for loops requires JIT-compilation in Numba.
Just as we work with Numba to provide both of these programming styles—array-oriented and imperative—we'll also be working with JAX to add autodifferentiation (Anish Biswas will be starting on this in January; he's actually continuing work from last spring, but in a different direction). We're also working with Martin Durant and Doug Davis to replace our homegrown lazy arrays with industry-standard Dask, as a new collection type (https://github.com/ContinuumIO/dask-awkward/ https://github.com/ContinuumIO/dask-awkward/). A lot of my time, with Ianna Osborne and Ioana Ifrim at my university, is being spent refactoring the internals to make these kinds of integrations easier (https://indico.cern.ch/event/855454/contributions/4605044/ https://indico.cern.ch/event/855454/contributions/4605044/). We found that we had implemented too much in C++ and need more, but not all, of the code to be in Python to be able to interact with third-party libraries.
If you have any other questions, I'd be happy to answer them!
- riskneutral 5y agoWas there no way to do this in Apache Arrow, or with some modifications to Arrow?
- jpivarski 5y agoNaturally, we considered this! :) We needed a larger set of tree node types than Apache Arrow in order to perform some of these operations without descending all the way down a subtree. For example, the implementation of the slice described in the video I linked above requires list offsets to be described as two separate arrays, which we call `starts` and `stops`, rather than a single set of `offsets`. So we have a ListArray (most general, uses `starts` and `stops`) and a ListOffsetArray (for the special case with `offsets`) and some operations produce one, other operations produce the other (as an internal detail, hidden from high-level users). Arrow's ListType is equivalent to the ListOffsetArray. If we had been forced to only use ListOffsetArrays, then the slice described in the video would have to propagate down to all of a tree node's children, and we want to avoid that because a wide record can have a lot of children. So Awkward Array has a superset of Arrow's node types. However, one of the great things about columnar data is that transformation between formats can share memory and be performed in constant time. In the ak.to_arrow/ak.from_arrow functions (https://awkward-array.readthedocs.io/en/latest/_auto/ak.to_arrow.html https://awkward-array.readthedocs.io/en/latest/_auto/ak.to_a...), ListOffsetArrays are replaced with Arrow's list node type with shared memory (i.e. we give it to Arrow as a pointer). Our ListArrays are rewritten as ListOffsetArrays, propagating down the tree, before giving it to Arrow. If you're doing some array slicing and your goal is to end up with Arrow arrays, what we've effectively done is delayed the evaluation of the changes that have to happen in the list node's children until they're needed to fit Arrow's format. You might do several slices in your workflow, but the expensive propagation into the list node's children happens just once when the array is finally being sent to Arrow. As far as what matters for users, Awkward Arrays are 100% compatible with Arrow through the ak.to_arrow/ak.from_arrow functions, and usually shares memory with O(1) cost (where "n" is the length of the array). When it isn't shared memory with O(1) conversion time, it's because it's doing evaluations that you were saved from having to do earlier.
- 5y ago