Serialization¶
Arrow serialization layer for automatic schema generation and IPC validation.
ArrowSerializableDataclass¶
Use ArrowSerializableDataclass as a mixin for dataclasses that need Arrow serialization. The Arrow schema is generated automatically from field annotations:
from dataclasses import dataclass
from typing import Annotated
import pyarrow as pa
from vgi_rpc import ArrowSerializableDataclass, ArrowType
@dataclass(frozen=True)
class Measurement(ArrowSerializableDataclass):
timestamp: str
value: float
count: Annotated[int, ArrowType(pa.int32())] # explicit Arrow type override
# Auto-generated schema
print(Measurement.ARROW_SCHEMA)
# timestamp: string
# value: double
# count: int32
# Serialize / deserialize
m = Measurement(timestamp="2024-01-01T00:00:00Z", value=42.0, count=7)
data = m.serialize_to_bytes()
m2 = Measurement.deserialize_from_bytes(data)
assert m == m2
These dataclasses work directly as RPC parameters and return types. They're also the base class for StreamState — any field you add to a stream state is automatically serialized between calls.
Type mappings¶
| Python type | Arrow type |
|---|---|
str |
utf8 |
bytes |
binary |
int |
int64 |
float |
float64 |
bool |
bool_ |
list[T] |
list_<T> |
dict[K, V] |
map_<K, V> |
frozenset[T] |
list_<T> |
Enum |
dictionary(int32, utf8) |
Optional[T] |
nullable T |
nested ArrowSerializableDataclass |
struct |
Annotated[T, ArrowType(...)] |
explicit type override |
Annotated[T \| None, ArrowType(...)] |
explicit type override, still nullable |
Annotated[T, ArrowType(..., nullable=...)] |
explicit type AND nullability |
Nullability¶
Optional[T] is a nullable column whether or not it is wrapped in Annotated.
Before 0.43.0 the wrapper hid it — get_origin(Annotated[T | None, ...]) is
Annotated, never a union — so every annotated optional field was described as
non-null while its values were free to be null. Nothing in one SDK notices that,
because the same schema describes both the write and the read; a peer comparing
schemas field by field rejects the batch.
Set nullable only where the wire genuinely disagrees with the Python type:
# Optional in Python because it has a default, never null on the wire.
cursor: Annotated[bytes | None, ArrowType(pa.binary(), nullable=False)] = None
# Not Optional in Python, but a peer may omit the value.
tag: Annotated[str, ArrowType(pa.string(), nullable=True)]
An explicit value is a claim about the bytes, so make it only where the serializer backs it up. Deriving from the annotation is right almost always.
API Reference¶
ArrowSerializableDataclass¶
ArrowSerializableDataclass
¶
Mixin for dataclasses with automatic Arrow IPC serialization.
Provides automatic schema generation and serialization/deserialization for frozen dataclasses. The ARROW_SCHEMA is auto-generated from field type annotations.
Auto-detected types: - Basic types: str, bytes, int, float, bool - Generic types: list[T], dict[K, V], frozenset[T] - NewType: unwraps to underlying type (e.g., NewType("Id", bytes) -> binary) - Enum: serializes as dictionary-encoded string via .name - ArrowSerializableDataclass: serializes as struct
Not supported:
- tuple: Arrow has no native heterogeneous-tuple type. Use a nested
dataclass (ArrowSerializableDataclass) for fixed, named fields, or
list[T] for homogeneous sequences.
Optional fields (annotated with | None) are marked as nullable.
To override specific field types, use Annotated with ArrowType.
| ATTRIBUTE | DESCRIPTION |
|---|---|
ARROW_SCHEMA |
Auto-generated Arrow schema from field annotations.
TYPE:
|
serialize
¶
Serialize this instance to an Arrow IPC stream.
| PARAMETER | DESCRIPTION |
|---|---|
dest
|
The destination to write to (must support binary writes, e.g., stdout pipe, BufferedWriter).
TYPE:
|
Source code in vgi_rpc/utils.py
serialize_to_bytes
¶
Serialize this instance to Arrow IPC bytes.
| RETURNS | DESCRIPTION |
|---|---|
bytes
|
Arrow IPC stream bytes containing a single-row RecordBatch. |
deserialize_from_batch
classmethod
¶
deserialize_from_batch(
batch: RecordBatch,
custom_metadata: KeyValueMetadata | None = None,
*,
ipc_validation: IpcValidation = FULL
) -> Self
Deserialize an instance from an Arrow RecordBatch.
| PARAMETER | DESCRIPTION |
|---|---|
batch
|
Single-row RecordBatch containing the serialized data.
TYPE:
|
custom_metadata
|
Optional metadata from the batch (unused, reserved for subclass overrides).
TYPE:
|
ipc_validation
|
Validation level for nested IPC batches.
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
Self
|
Deserialized instance of this class. |
| RAISES | DESCRIPTION |
|---|---|
ValueError
|
If the batch is invalid (wrong row count or missing fields). |
TypeError
|
If a field value has an unexpected type during conversion. |
KeyError
|
If an Enum name cannot be resolved. |
Source code in vgi_rpc/utils.py
1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 | |
deserialize_from_bytes
classmethod
¶
deserialize_from_bytes(
data: bytes, ipc_validation: IpcValidation = FULL
) -> Self
Deserialize an instance from Arrow IPC bytes.
| PARAMETER | DESCRIPTION |
|---|---|
data
|
Arrow IPC stream bytes containing a single-row RecordBatch.
TYPE:
|
ipc_validation
|
Validation level for the deserialized batch.
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
Self
|
Deserialized instance of this class. |
| RAISES | DESCRIPTION |
|---|---|
ValueError
|
If the batch is invalid (wrong row count or missing fields). |
IPCError
|
If the IPC stream is malformed or truncated. |
TypeError
|
If a field value has an unexpected type during conversion. |
KeyError
|
If an Enum name cannot be resolved. |
Source code in vgi_rpc/utils.py
ArrowType¶
ArrowType
dataclass
¶
Annotation marker to specify explicit Arrow type for a field.
Use with Annotated to override the default inferred Arrow type:
@dataclass(frozen=True)
class MyData(ArrowSerializableDataclass):
# Override int64 → int32
count: Annotated[int, ArrowType(pa.int32())]
# Override for nested list of int32
matrix: Annotated[
list[list[int]], ArrowType(pa.list_(pa.list_(pa.int32())))
]
Nullability normally follows the Python annotation — X | None is a
nullable column, with or without an Annotated wrapper. Set nullable
only where the wire genuinely disagrees with the Python type, which happens
in two directions:
# Optional in Python because it has a default, never null on the wire —
# the serializer substitutes b"" for None.
cursor: Annotated[bytes | None, ArrowType(pa.binary(), nullable=False)] = None
# Not Optional in Python, but a peer may omit the column's value.
tag: Annotated[str, ArrowType(pa.string(), nullable=True)]
An explicit value is a claim about the BYTES, so it should be made only
where the serializer's behaviour backs it up. Leaving it None (derive
from the annotation) is right almost always.
| ATTRIBUTE | DESCRIPTION |
|---|---|
arrow_type |
The Arrow type to use for this field.
TYPE:
|
nullable |
Wire nullability override;
TYPE:
|
IpcValidation¶
IpcValidation
¶
Bases: Enum
Level of validation applied to incoming IPC record batches.
| ATTRIBUTE | DESCRIPTION |
|---|---|
NONE |
No validation — batches are used as-is.
|
STANDARD |
Call
|
FULL |
Call
|
from_env
classmethod
¶
from_env(
default: IpcValidation | None = None,
) -> IpcValidation
Resolve a validation level from VGI_RPC_IPC_VALIDATION.
Lets an operator trade validation for throughput at deploy time
without a code change. FULL walks every buffer of every
incoming batch, which is real money on large Arrow payloads
(measured at ~7% of HTTP server time), but it is also the
defence against malformed attacker-supplied IPC — so the default
stays FULL and lowering it is an explicit, deliberate act.
Accepts none, standard or full (case-insensitive).
An unrecognised value warns and falls back to default rather
than raising: the failure direction is more validation, never
silently less than the operator believes they configured.
| PARAMETER | DESCRIPTION |
|---|---|
default
|
Level to use when the variable is unset or invalid.
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
IpcValidation
|
The resolved validation level. |
Source code in vgi_rpc/utils.py
ValidatedReader¶
ValidatedReader
¶
ValidatedReader(
reader: RecordBatchStreamReader,
ipc_validation: IpcValidation,
)
Wrapper around ipc.RecordBatchStreamReader that validates every batch on read.
Proxies the subset of the reader API used by the RPC framework
(read_next_batch, read_next_batch_with_custom_metadata,
schema, and the context manager protocol). Downstream code
needs zero changes — just wrap ipc.open_stream(...) in
ValidatedReader(..., ipc_validation).
When ipc_validation is IpcValidation.NONE, each read still
delegates to the inner reader with minimal extra overhead.
Wrap reader so every batch is validated at ipc_validation level.
Source code in vgi_rpc/utils.py
ipc_validation
property
¶
ipc_validation: IpcValidation
The validation level applied to every batch read.
read_next_batch
¶
Read the next batch, validating it before returning.
read_next_batch_with_custom_metadata
¶
Read the next batch with custom metadata, validating before returning.
Source code in vgi_rpc/utils.py
__enter__
¶
__exit__
¶
__exit__(
exc_type: type[BaseException] | None,
exc_val: BaseException | None,
exc_tb: TracebackType | None,
) -> None
Exit the context manager.
Source code in vgi_rpc/utils.py
Validation¶
validate_batch
¶
validate_batch(
batch: RecordBatch, ipc_validation: IpcValidation
) -> None
Validate a RecordBatch at the specified level.
| PARAMETER | DESCRIPTION |
|---|---|
batch
|
The batch to validate.
TYPE:
|
ipc_validation
|
Validation level (NONE, STANDARD, or FULL).
TYPE:
|
| RAISES | DESCRIPTION |
|---|---|
IPCError
|
If validation fails. |
Source code in vgi_rpc/utils.py
Errors¶
IPCError
¶
Bases: Exception
Error during IPC message reading or writing.