Documentation
¶
Overview ¶
Package wshf owns the WSHF columnar shuffle wire format: the magics, the envelope codecs, and the one bounds-checked decoder every consumer uses.
The format replaces Parquet for inter-stage shuffle data. It avoids per-row goparquet.Value allocation (nRows × nCols objects), alphabetical column reordering, and Parquet page/RLE encoding overhead.
Magic "WSHF" (4 bytes)
NumChunks uint32 (4 bytes)
NumCols uint16 (2 bytes)
Schema: for each column:
NameLen uint16
Name []byte
TypeID uint8
Scale, Precision uint8 ×2 — DECIMAL only
Chunks: for each chunk:
NumRows uint32 (4 bytes)
For each column:
NullBitmapWords uint32 (number of uint64 words)
NullBitmap []uint64
DataLen uint32 (byte length of column data)
Data []byte (type-dependent raw data)
The WRITER lives in internal/worker (it needs the engine's batch gather, view resolution and the WIDX extent-index footer). This package is the read side, and it is the only read side: the coordinator's inline-result path and the worker's file/stream/pread paths all decode through it, so a payload cannot be interpreted two ways (#422).
Every read goes through Cursor, which returns an error rather than panicking on short input. That matters because the bytes are untrusted: the coordinator decodes a NATS payload from a worker in the decode goroutine of readInlineResults, where a panic is not a failed query but a dead coordinator.
Index ¶
- Constants
- Variables
- func DecodeBatches(data []byte) ([]*batch.RecordBatch, error)
- func DecodeChunk(schema []parquet.Column, numRows int, chunkBytes []byte, chunkIdx uint32) (*batch.RecordBatch, error)
- func Decompress(data []byte) ([]byte, error)
- func DecompressStream(src io.Reader, dst io.Writer, codec Codec) error
- func FixedTypeLen(typ parquet.TypeID, numRows int) (int, error)
- func IsShuffleFormat(data []byte) bool
- func ParseHeader(c *Cursor) (schema []parquet.Column, numChunks uint32, err error)
- func ReadColumn(c *Cursor, vec *batch.Vector, numRows int, typ parquet.TypeID) error
- func ValidateChunkBytes(schema []parquet.Column, numRows int, buf []byte) error
- type ChunkReader
- type Codec
- type Cursor
- func (c *Cursor) Len32(what string, max int) (int, error)
- func (c *Cursor) Peek(n int, what string) ([]byte, error)
- func (c *Cursor) Pos() int
- func (c *Cursor) Remaining() int
- func (c *Cursor) Size() int
- func (c *Cursor) Skip(n int, what string) error
- func (c *Cursor) Take(n int, what string) ([]byte, error)
- func (c *Cursor) U8(what string) (uint8, error)
- func (c *Cursor) U16(what string) (uint16, error)
- func (c *Cursor) U32(what string) (uint32, error)
Constants ¶
const ( MaxCols = 1 << 12 MaxNameLen = 1 << 12 MaxRows = 1 << 26 MaxBytesLen = 1 << 31 )
Plausibility ceilings. A length field is not a promise: these bound what a header may claim before the decoder allocates or skips by it. They are far above anything the writer emits (2048-row batches, engine schemas) and far below anything that would exhaust memory on a corrupt field.
const ( // LenBytes: [dataLen u32][data][numRows × u32 end offsets]. LenBytes = -1 // LenContainer: [payloadLen u32][payload]. The payload is // self-describing (batch.EncodeContainerColumn) and the walk skips it // whole — ARRAY/ROW/MAP/VECTOR have no per-row width at all. LenContainer = -2 )
Sentinels returned by FixedTypeLen for the two classes whose byte count is not a function of the row count.
const HeaderLen = 10
HeaderLen is the smallest possible header: magic + NumChunks + NumCols.
Variables ¶
var ( MagicWSHF = [4]byte{'W', 'S', 'H', 'F'} MagicWSHC = [4]byte{'W', 'S', 'H', 'C'} MagicWSHZ = [4]byte{'W', 'S', 'H', 'Z'} )
Wire magics. These four-byte constants ARE the wire contract (ADR-0010): MagicWSHF is the raw payload, MagicWSHC an s2 stream of it, MagicWSHZ a zstd stream of it (docs/design/exchange-zstd-wire.md).
Functions ¶
func DecodeBatches ¶
func DecodeBatches(data []byte) ([]*batch.RecordBatch, error)
DecodeBatches decodes every chunk in a raw WSHF payload. Callers holding the whole payload in memory already (inline results, gather replies) use this; file-backed readers use ChunkReader.
func DecodeChunk ¶
func DecodeChunk(schema []parquet.Column, numRows int, chunkBytes []byte, chunkIdx uint32) (*batch.RecordBatch, error)
DecodeChunk materializes one staged chunk's column segments (the bytes AFTER the row-count word) into a fresh RecordBatch. Shared by the serial stream path, the decode-ahead workers and the index-mode pread workers so they cannot diverge on payload interpretation; chunkIdx is for error text.
func Decompress ¶
Decompress unwraps a WSHC (s2) or WSHZ (zstd) envelope back to raw WSHF. Plain WSHF — or anything that is not a shuffle payload at all, e.g. a parquet result file — is returned unchanged, so callers can sniff once and branch after.
Both envelopes, not just WSHC: WSHZ is what an S3 stage upload carries under WADJET_EXCHANGE_ZSTD=1, and a reader that knows only WSHC hands the compressed bytes on to a parquet decoder and fails with a parquet error on a perfectly good shuffle file.
The worker's own DecompressShuffleData is the pooled, streaming variant of this for its hot file paths; this is the whole-payload form for callers that already hold the bytes.
func DecompressStream ¶
DecompressStream copies the compressed body that follows a WSHC/WSHZ magic from src to dst. codec names the envelope (the caller sniffed the magic); the WSHF magic itself is inside the compressed body, so dst receives a complete WSHF payload.
func FixedTypeLen ¶
FixedTypeLen returns the exact payload byte length for fixed-width shuffle types, or one of the sentinels above for the variable-length classes. Shared by the decoder, the streaming stage walk and the index-mode extent validation so the three cannot diverge.
func IsShuffleFormat ¶
IsShuffleFormat reports whether data starts with any shuffle magic.
func ParseHeader ¶
ParseHeader consumes the WSHF magic, chunk count and schema from c, leaving the cursor at the first chunk's row-count word.
func ReadColumn ¶
ReadColumn decodes one column segment into vec, advancing c past it. Every count and length is checked against the bytes that remain before it is used — this is the bounds-checked replacement for the two hand-copied unchecked walks (#422).
func ValidateChunkBytes ¶
ValidateChunkBytes walks one chunk's column segments in buf (the bytes AFTER the row-count word) and requires the walk to consume buf in full. Index-mode decode workers run it over their pread extent before decoding: the decoder is bounds-checked on its own, but "these bytes are exactly one chunk" is a stronger claim than "this decode did not run off the end", and an extent that is off by a column is a wrong answer, not a crash.
Types ¶
type ChunkReader ¶
type ChunkReader struct {
// contains filtered or unexported fields
}
ChunkReader iterates over the chunks in a WSHF byte slice one at a time, allocating a single RecordBatch per Next call. Callers hold only one batch in memory at a time instead of materializing the whole payload.
func NewChunkReader ¶
func NewChunkReader(data []byte) (*ChunkReader, error)
NewChunkReader parses the WSHF header and returns a reader positioned at the first chunk. The caller retains ownership of data — it must remain valid for the lifetime of the reader (batches copy their bytes out, so the data may be released once the last batch is in hand).
func (*ChunkReader) Next ¶
func (r *ChunkReader) Next() (*batch.RecordBatch, error)
Next returns the next RecordBatch, or (nil, nil) when all chunks have been consumed. Allocates exactly one RecordBatch per non-empty chunk.
func (*ChunkReader) NumChunks ¶
func (r *ChunkReader) NumChunks() uint32
NumChunks is the chunk count the header promised.
func (*ChunkReader) Pos ¶
func (r *ChunkReader) Pos() int
Pos returns the reader's byte offset into the WSHF slice — everything below it has been fully decoded (batches copy column data out), so the drop-behind walk can discard those pages. Strictly monotonic.
func (*ChunkReader) Schema ¶
func (r *ChunkReader) Schema() []parquet.Column
Schema is the decoded column schema.
type Codec ¶
type Codec uint8
Codec identifies the envelope around a WSHF payload.
func CodecForMagic ¶
CodecForMagic maps a 4-byte magic to its codec. ok=false means the payload is not a shuffle format at all (e.g. parquet).
type Cursor ¶
type Cursor struct {
// contains filtered or unexported fields
}
Cursor is a bounds-checked walk over a WSHF byte slice. Every read returns an error instead of indexing past the end, which is the whole point: the decoder's counts and lengths come OUT of the bytes it is walking, so a truncated or corrupt payload otherwise turns a length field into a slice bound (#422 — a short inline result panicked the coordinator's decode goroutine, which nothing above recovers).
The zero value is not usable; construct with NewCursor. Cursors are values, not pointers to the data they walk — Take returns a subslice that ALIASES the underlying bytes, so callers must copy anything that outlives the payload (the column readers all copy).
func NewCursorAt ¶
NewCursorAt returns a cursor positioned at pos, for the callers that resume a walk they interrupted (the chunk reader's per-chunk position).
func (*Cursor) Len32 ¶
Len32 reads a little-endian uint32 length field and returns it as an int bounded by max — a length is a claim about bytes that have not been checked yet, so it is range-checked before anything allocates or slices by it. On a 32-bit platform a uint32 near 2^32 would also wrap negative as an int; the max check rejects it either way.
func (*Cursor) Peek ¶
Peek returns the next n bytes without advancing. See Take for the bounds-check shape.
func (*Cursor) Take ¶
Take advances by n and returns those bytes, or an error if fewer than n remain. A negative n is a corrupt length field, not a rewind: the uint(n) conversion wraps a negative n to a huge value, so the same single comparison rejects it and an insufficient remainder alike (c.pos never exceeds len(c.data), so len(c.data)-c.pos is never negative and the conversion on that side is always exact).