IPC Format¶
The IPC format streams Vortex arrays between processes as a sequence of self-delimiting messages.
Arrays keep their encodings on the wire, so compressed data crosses the process boundary without
being decompressed. The format is implemented by the vortex-ipc crate, re-exported as
vortex::ipc.
Warning
The IPC format is unstable. It is not covered by the versioning guarantees of the
file format, and it has a single message version, V0, which readers require
exactly. If you want to persist data, please consider writing a Vortex file instead.
Streams¶
An IPC stream carries a single logical array as a sequence of chunks:
[DType message] [Array message] [Array message] ...
The first message is a
DTypeMessagecarrying the stream’s data type.Every following message is an
ArrayMessageholding one chunk, whose data type must equal the stream’s.The stream ends at end-of-input after a complete message. There is no end-of-stream marker, so a stream with no chunks is a lone
DTypeMessage.
Message framing¶
Every message has the same framing:
<4 bytes> u32 little-endian header length, H
<H bytes> Message FlatBuffer
<body_size bytes> message body
The Message FlatBuffer is the header, which contains the following fields:
version(MessageVersion, defaultV0): the message format version. Readers reject any value other thanV0.header(MessageHeaderunion): the message type, which also determines how to interpret the body.body_size(uint64): the exact length of the body, so a reader always knows where the next message starts.
Messages are written back-to-back, with no padding between or within them. This means that header or body bytes are not guaranteed to land on an aligned offset. Thus, after receiving a message, a reader copies each FlatBuffer and data buffer into appropriately aligned memory.
Message types¶
DTypeMessage¶
The body is a DType FlatBuffer (root_type DType in dtype.fbs). The
DTypeMessage table itself has no fields.
ArrayMessage¶
The body is a serialized array. The header carries the information needed to decode it:
row_count(uint32) — the length of the root array.encodings([string]) — the encoding IDs the array references, such asvortex.primitive.
The body has the same shape as any other serialized array:
[buffer 0] [buffer 1] ... [buffer N-1] [Array FlatBuffer] [u32 little-endian FlatBuffer length]
A reader decodes the body from the end:
Read the trailing
u32to find the length of theArrayFlatBuffer that precedes it. The bytes before the FlatBuffer hold the data buffers.Locate each data buffer from the
Array.bufferstable, in order. Bufferistartspaddingbytes after the end of bufferi - 1, where the first buffer’s predecessor ends at offset 0, and spanslengthbytes. Align it to 2alignment_exponentbytes.Decode the
ArrayNodetree fromArray.root. Each node’sencodingis an index into the message’sencodingslist, and itsbuffersare indices into theArray.bufferstable.Decode the root node with the stream’s data type and
row_count. Each encoding derives the data types and lengths of its children from its own metadata. Anystatspresent on a node are attached to the decoded array.
Each message is self-contained: its encodings list and buffer table apply only to that message.
The IPC writer never pads buffers, so padding is always 0. Readers should still honor it,
because Vortex files use the same array layout, which do pad. The writer always sets
compression to None. LZ4 is reserved in the schema but not implemented; current readers
don’t check the field.
The ArrayNode tree and its statistics are described in
Serialization.
BufferMessage¶
The body is a single raw byte buffer. alignment_exponent is the alignment the receiver must give
it, as a power of two. Readers reject exponents above 16 (64 KiB), since the value is untrusted.
The stream writers never emit a BufferMessage, and the stream readers reject one. It is available
only through the message-level API below.
Rust API¶
The crate has three layers. Most callers only need the array-stream layer.
Layer |
Write |
Read |
|---|---|---|
Array streams |
|
|
Messages |
|
|
Framing |
|
|
ArrayStreamIPC and ArrayIteratorIPC are implemented for every ArrayStream and
ArrayIterator. They write the DTypeMessage, then one ArrayMessage per item. The iterator and
stream from ArrayRef::to_array_iterator and ArrayRef::to_array_stream yield one item per chunk
of a chunked array.
AsyncIPCReader and SyncIPCReader read such a stream back as an ArrayStream or ArrayIterator.
They fail if the first message is not a DTypeMessage, if a later message is not an
ArrayMessage, or if a decoded array’s data type differs from the stream’s.
MessageDecoder holds no IO. It reports the total number of bytes it needs to make progress, so
callers can drive it over any transport.
use std::io::Cursor;
use vortex::VortexSessionDefault;
use vortex::array::IntoArray;
use vortex::array::iter::ArrayIteratorExt;
use vortex::buffer::buffer;
use vortex::error::VortexResult;
use vortex::ipc::iterator::ArrayIteratorIPC;
use vortex::ipc::iterator::SyncIPCReader;
use vortex::session::VortexSession;
fn round_trip() -> VortexResult<()> {
let session = VortexSession::default();
let array = buffer![1i32, 2, 3].into_array();
let bytes = array.to_array_iterator().write_ipc(Vec::new(), &session)?;
let reader = SyncIPCReader::try_new(Cursor::new(bytes), &session)?;
assert_eq!(reader.read_all()?.len(), 3);
Ok(())
}
Decoding resolves encoding IDs through the session, so the reader’s session must register every encoding the writer used.
Limitations¶
No zero-copy reads. Buffers are unpadded on the wire, so the reader copies them into aligned memory.
No state shared across messages. Each message repeats its own encoding list, and data such as a dictionary shared between chunks is sent again with every chunk.
Row count limit. An
ArrayMessageholds at mostu32::MAXrows. Larger arrays must be split into chunks.No end-of-stream marker. A stream cut exactly at a message boundary is indistinguishable from a complete one. Transports that can drop data must detect truncation themselves.
FlatBuffer definition¶
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors
include "vortex-array/array.fbs";
include "vortex-dtype/dtype.fbs";
enum MessageVersion: uint8 {
V0 = 0,
}
/// Indicates the message body contains a flatbuffer Array message, followed by array buffers.
table ArrayMessage {
/// The row count of the array.
row_count: uint32;
/// The encodings referenced by the array.
encodings: [string];
}
/// Indicates the body contains a regular byte buffer.
table BufferMessage {
alignment_exponent: uint8;
}
/// Indicates the body contains a flatbuffer DType message.
table DTypeMessage {}
union MessageHeader {
ArrayMessage,
BufferMessage,
DTypeMessage,
}
table Message {
version: MessageVersion = V0;
header: MessageHeader;
body_size: uint64;
}
root_type Message;