Lakestream Storage Spec
The Storage Spec defines format version 3: WAL objects, the offset index, compacted objects, and the append, compaction and trim protocols.
This is a specification for the Lakestream storage format, which stores the logs of a stream as entries in immutable objects in object storage, located by an offset index in a linearizable metadata store.
The key words "MUST", "MUST NOT", "REQUIRED", "SHALL", "SHALL NOT", "SHOULD", "SHOULD NOT", "RECOMMENDED", "NOT RECOMMENDED", "MAY", and "OPTIONAL" in this document are to be interpreted as described in BCP 14 [RFC 2119] [RFC 8174] when, and only when, they appear in all capitals, as shown here.
Each requirement begins with the role it binds, in brackets. Conformance defines the roles.
Format Versioning
Version 3 of the Lakestream storage format is complete. Ursa 1.0 writes it by default.
Versions 1 and 2 are older encodings that a reader can meet in deployments that predate version 3. Appendix C describes them.
The format version is a property of a deployment, not of an individual object or record.
[All] All WAL objects and index records of a deployment MUST use one format version.
Version 3 artifacts carry the number 3 in two places: the index_version field that opens the index region of every WAL object, and the version field of every EntryIndex.
[Reader] A reader SHOULD reject a WAL object whose index_version is not 3, and an EntryIndex whose version is not 3.
Ursa 1.0 selects the format version from its configuration and does not check either number when it reads; see Implementation status.
Every change to the format goes through a Lakestream Improvement Proposal (LIP) and produces a new format version.
| v1 | v2 | v3 | Feature |
|---|---|---|---|
| required | required | Legacy WAL object layout: an int64 index length, then 32-byte index tuples | |
| required | Version 3 WAL object layout: an int32 index length, index_version, and one section per log | ||
| required | required | Index values as UTF-8 strings | |
| required | Index values as protobuf EntryIndex messages | ||
| one per entry | one per entry | one per log section | Index records per WAL object |
| required | required | File type (RAW or PARQUET) in the index value | |
| required | Index value fields index_type, file_size and entry_count | ||
| required | Entry boundaries (entry_offsets) in index records that hold more than one entry | ||
| required | extraData with a file index (CompactedObjectFileIndex) in every PARQUET record |
Version 3
In version 3, a WAL object's index region lists each log's entries in a log section, and one index record indexes each log section. Index values are protobuf EntryIndex messages that carry entry boundaries, and each PARQUET record locates its compacted data through a file index. Specification defines version 3.
Versions 1 and 2
Versions 1 and 2 use a legacy WAL object layout and index every entry with a record of its own, whose value is a UTF-8 string. Version 2 adds the file type to that string. This specification does not define them normatively; Appendix C describes them for readers of older deployments.
Goals
- Protocol independence -- Entry payloads are opaque bytes at every layer. A binding, such as the Kafka binding in Appendix A, defines what a payload holds and which identifiers compacted objects carry; the format carries no protocol rules.
- Immutable data objects -- WAL objects and compacted objects are written once, under names that are never reused, and are never modified. Index records in the metadata store are replaced by compaction and removed by trimming.
- Leaderless appends -- Any writer holding a write lease on a log can append to it without an elected leader. Offsets come from a sequenced put in the metadata store, so concurrent appends receive disjoint, contiguous offset ranges in the store's commit order.
- Continuous readability -- Compaction installs the record for a compacted range before it deletes the records that range replaces, and objects are reclaimed only after a grace period, so every acknowledged offset at or above the trim marker stays readable.
- Metadata and data separation -- Object storage holds entry bytes. A linearizable metadata store holds offsets, the offset index, entry timestamps, the trim marker, write fences and write lease records.
Overview
A log's lifecycle has five parts:
- Append. A writer groups entries for one or more logs into a WAL object and writes it to object storage. For each log in the object, it then publishes one index record with a sequenced put, which assigns the entries' offsets and the log's cumulative size. The writer acknowledges the entries once their index record has committed.
- Read. A reader applies the log's trim marker, finds the index record that covers an offset with one lookup, and reads the entries from a WAL object or from compacted objects, as the record's
file_typedirects. - Compaction. A compactor rewrites the entries of a range of RAW records into compacted objects: Parquet files, each with an
.indexcompanion file. It then installs one PARQUET record for the range, and only after that deletes the RAW records the range replaces. - Trim. A writer advances the trim marker. Compacted objects that lie wholly below the marker are deleted, and then the index records that name them are removed. WAL objects that no index record names are deleted after a grace period.
- Lifecycle. A writer or compactor holds a write lease while it changes a log, and announces the lease with an ephemeral record in the metadata store. Deleting a log installs a permanent write fence, waits until the log has no lease records, and then removes the log's index records.
Storage operations
The format uses three object-storage operations:
- Write once -- Create an object under a name that has not been used before. The format never writes an object twice, appends to it, or changes it in place.
- Read -- Read an object, whole or by byte range. Parquet files are read by byte range.
- Delete -- Remove an object that no index record at or above the trim marker needs, within the rules of Trim, deletion and reclamation.
The format renames nothing and needs no operation that spans several objects. These operations match what object stores such as S3 provide.
The metadata store is a linearizable key-value store divided into partitions. Each log's records live in one partition: its offset index, its trim marker, its write fence and its write lease records. Offsets come from a sequenced put, which creates a new index record under a key that it derives, in one atomic step, from the log's greatest key. Metadata store model lists every primitive the format uses.
Not defined by this specification
The following are implementation-defined. Appendix B and Appendix E describe Ursa's choices.
- The serialization of catalog records.
- How compactors claim and schedule compaction work, and the compaction cursor.
- Cursor and acknowledgement state.
- The mechanism that reclaims objects, beyond the safety rules in Trim, deletion and reclamation.
- How WAL objects are named, beyond uniqueness.
- How log IDs are allocated, beyond uniqueness.
- Materialization into external tables, which the Materialization Spec covers.
Specification
Terms
- Deployment -- The systems that share one metadata store, one WAL root and one compacted root.
- Stream -- A named collection of one or more logs, identified by a namespace and a stream name.
- Log -- An append-only sequence of records in a stream, identified by a log ID.
- Log ID -- A non-negative 64-bit integer that identifies a log. This specification writes it
L. - Record -- The unit that offsets count. The format never looks inside a record.
- Entry -- The payload of one append: an opaque byte string that holds one or more records.
- Offset -- The position of a record in its log. The first record of a log has offset 0.
- Index record -- A metadata-store record that locates the entries holding a contiguous range of a log's offsets. Its key carries the end offset of the range; its value is an
EntryIndex. - End offset -- The exclusive end of an index record's range: the offset that follows its last record.
- Cumulative size -- The total payload size, in bytes, of all the entries a log has received up to a given offset.
- RAW record, PARQUET record -- An index record whose entries are in a WAL object (
file_typeRAW), or in compacted objects (file_typePARQUET). The PARQUET record that a compaction installs hasindex_typeCOMPACT, and the overview diagram labels it COMPACT. - WAL object -- An immutable object that holds entries of one or more logs, and an index region that lists each log's entries.
- Log section -- The part of a WAL object's index region that lists one log's entries.
- Payload region -- The part of a WAL object after its index region, which holds the entry payloads.
- Compacted object -- An immutable Parquet file that holds entries of a compacted range, one row per entry, together with its companion file.
- Companion file -- The
.indexfile beside a compacted Parquet file. It holds each row's metadata and a map from entry offsets to rows. - File index -- The
CompactedObjectFileIndexvalue of a PARQUET record: the names of the record's compacted files, each with the last offset it holds. - Compaction -- Rewriting the entries of a range of RAW records into compacted objects, and replacing those records with one PARQUET record.
- Trim marker -- The exclusive lower bound of a log's readable offsets, written
B. Offsets belowBare trimmed. - Soft trim -- Advancing the trim marker.
- Hard trim -- Removing the index records of a log that end at or below a given offset.
- Write lease -- A claim on a log that a writer or compactor holds while it changes the log's records, announced by an ephemeral metadata-store record.
- Write fence -- A permanent metadata-store record that stops new write leases on a log. A catalog installs it when it deletes the log.
- Metadata store -- The linearizable key-value store that holds index records, trim markers, write fences and write lease records.
- Partition -- A unit of the metadata store within which operations are atomic and linearizable. All records of log
Lare in its partition,P(L). - Sequenced put -- The metadata-store primitive that creates an index record under a key it derives from the log's greatest key.
- WAL root, compacted root -- The object-storage locations, set by deployment configuration, that hold WAL objects and compacted objects.
- Binding -- A protocol-specific definition of what entry payloads hold and which identifiers compacted objects carry. Appendix A defines the Kafka binding.
Conformance
The format has four roles:
| Role | Scope |
|---|---|
| Reader | Finds and reads a log's entries: version checks in Format Versioning, compacted file names in Object storage, roots and naming, parsing the structures in WAL objects, Offset index keys, Offset index values, Entry offsets and Compacted objects, Reading, and Kafka entry payloads in Appendix A. |
| Writer | Appends entries and advances the trim marker: WAL object names in Object storage, roots and naming, Logs, entries and offsets, WAL objects, RAW records in Offset index values, Entry offsets, Appending, write leases and session loss in Write leases and fences, soft trim in Trim, deletion and reclamation, and Kafka entry payloads in Appendix A. |
| Compactor | Rewrites ranges into compacted objects, updates the offset index, and removes trimmed data: compacted file names in Object storage, roots and naming, PARQUET records in Offset index values, write leases and session loss in Write leases and fences, Compacted objects, Compaction, and hard trim and object deletion in Trim, deletion and reclamation. |
| Catalog | Creates and deletes logs and streams: namespaces and stream names in Object storage, roots and naming, fencing and deletion in Write leases and fences, and Catalog semantics. |
Requirements tagged [All] apply to every role. A requirement tagged with two roles applies to each of them.
An implementation conforms to format version 3 in a role when it satisfies every requirement of that role. An implementation states the roles it implements. Implementation status lists each implementation's roles.
Tables that define the fields of a structure have a v3 column, which says what format version 3 expects of the implementation that writes the structure.
[All] An implementation that writes a structure defined by a field table MUST write every field marked required, and MAY write a field marked optional.
Where the text sets a condition on a field marked optional, the condition applies.
Metadata store model
The metadata store is a key-value store. Keys are UTF-8 strings and values are byte strings. Every record belongs to one partition, which a string names. A store can keep several partitions together, so a lookup in one partition can return a record of another.
[All] An implementation MUST keep every record of log L that this specification defines in the partition P(L), whose name is the decimal representation of L without leading zeros.
| Record | Key | Value |
|---|---|---|
| Index record | {log_id:020d}-{end_offset:020d}-{cumulative_size:020d} | An EntryIndex; see Offset index values |
| Trim marker | /mark-deleted-offsets/{L} | B, as ASCII decimal digits |
| Write fence | /stream-write-fences/{L} | The single byte 0x01 |
| Write lease record | /stream-write-leases/{L}/{token} | The lease's token, in UTF-8 |
In these keys, {L} is the log ID in decimal without leading zeros, and {log_id:020d} is the log ID as exactly 20 decimal digits.
[All] An implementation MUST NOT store any record other than an index record of L under a key that begins with {log_id:020d}- for L.
[All] An implementation MUST use a metadata store that provides:
- Linearizable partitions -- Every primitive below acts atomically on one partition, and the operations on each partition are linearizable.
- Key order -- Within a partition, keys that contain no
/are ordered bytewise, and no key that contains/sorts between two keys that contain no/. - Versions -- Every record has a version, which changes each time the record is written.
- Creation timestamps -- Every record has a creation timestamp, in milliseconds since the Unix epoch, which the store sets when it creates the record. Replacing the value of an existing record keeps its creation timestamp.
- Sessions -- A client holds a session with the store. An ephemeral record exists until the session that created it ends.
The format uses the primitives below, except Ceiling and Scan, which are listed for completeness: no procedure in this specification uses them, and a metadata store need not provide them. P is a partition, k, a and b are keys, and v is a value.
| Primitive | Semantics | Used by |
|---|---|---|
Get(P, k) | The record under k, with its value, version and creation timestamp, or nothing. | Reading, leases, trim, compaction |
Floor(P, k) | The record with the greatest key at or below k. | Reading |
Ceiling(P, k) | The record with the smallest key at or above k. | None of the procedures here |
Higher(P, k) | The record with the smallest key above k. | Reading |
Scan(P, a, b) | The records whose keys are at or above a and below b, in key order. | None of the procedures here |
List(P, prefix) | The keys that begin with prefix. | Log deletion |
SequencedPut(P, p, d1, d2, v) | Atomically: find the greatest key in P that begins with p followed by -; take its two 20-digit fields as n1 and n2, or 0 and 0 if there is none; create the key p-{n1+d1:020d}-{n2+d2:020d} with value v; return the new key and its creation timestamp. d1 is greater than 0 and d2 is at least 0. The primitive takes no version condition. | Appending |
Put(P, k, v) | Create k with v, or replace the value under k. | Compaction |
CreateIfAbsent(P, k, v) | Create k only if no record exists under it; the record can be ephemeral. | Trim marker, leases, fences |
CompareAndSet(P, k, version, v) | Replace the value under k only if the record's version equals version. | Trim marker |
Delete(P, k) | Remove the record under k. | Leases |
DeleteRange(P, a, b) | Remove every record whose key is at or above a and below b. | Compaction, hard trim, log deletion |
Index keys contain no /, and the other records of a log begin with /. With the key-order property, every key that sorts between two index keys of a log is an index key of the same log. A lookup can still return a key beyond a log's first or last index key, which is why Reading checks the prefix of the key it gets back.
The sequenced put derives new offsets from the greatest existing index key of the log. Removing a log's greatest index record would therefore make the next append reuse offsets; the rules in Trim, deletion and reclamation keep that record in place.
Object storage, roots and naming
The WAL root and the compacted root are locations in object storage, such as a bucket and a key prefix, that deployment configuration sets.
[All] All implementations in a deployment MUST use the same WAL root and the same compacted root.
WAL object names
A RAW record's location is the name of its WAL object, relative to the WAL root. With a key prefix r, the object's key is r/ followed by the name.
[Writer] A writer MUST give every WAL object a name that no other object in the WAL root has had, and MUST NOT reuse a name, even after its object has been deleted.
[Writer] A writer MUST NOT change or overwrite a WAL object after writing it.
Beyond uniqueness, WAL object names are implementation-defined. Appendix E describes Ursa's names and what its garbage collector assumes about them.
Compacted object names
The compacted objects of a stream live in one directory:
{compacted root}/{namespace}/{stream name}/The parts are joined with /. The namespace and stream name identify the stream that a log belongs to, so every log of a stream uses the same directory.
[Catalog] In format version 3, a stream namespace MUST NOT contain /.
[Catalog] In format version 3, a stream name MUST NOT contain /.
With both rules, the directory of every stream lies exactly two levels below the compacted root.
[Compactor] A compactor MUST write the compacted files of a log, and their companion files, in the directory of the log's stream.
[Compactor] A compactor MUST give each compacted file a name that ends in .parquet, contains .parquet nowhere else, and has not been used before in the directory. The file's companion has the same name with .parquet replaced by .index, in the same directory.
[Compactor] A compactor MUST NOT change or overwrite a compacted file or a companion file after writing it.
[Reader] A reader MUST resolve the file names in a file index against the directory of the log's stream, and MUST locate a file's companion by replacing the .parquet suffix of the file's name with .index.
A reader therefore needs the namespace and name of the stream that a log belongs to; it learns them from the catalog. The Parquet footer key rowMetadata also names the companion file, but readers do not use it.
Logs, entries and offsets
A log is an append-only sequence of records, numbered by offset from 0. An append adds one entry, which holds one or more records and takes the next offsets in order.
[Writer] Every entry MUST hold at least one record.
An index record covers the range [start, end) of whole entries, where end is the end offset in its key and start is end - message_count. The rules of Appending and Compaction keep the index records of a log that has not been deleted disjoint and adjacent: each record starts where the previous one ends. There are two exceptions, both from compaction. Between steps C3 and C4, the new PARQUET record overlaps the records it replaces. And when, in a compaction of the range [s, e), a hard trim that removes records ending above e runs between the Get and the Put of step C3, the Put re-creates the record for [s, e) and step C4 keeps it, which leaves a gap between that record and the first record that the hard trim kept. The record and the gap both lie below B; see Offset index keys. In a log that has been deleted, the late writes that Fences and log deletion describes can leave records that overlap or have gaps between them: a sequenced put that finds no index record of the log starts again at offset 0.
The cumulative size in a record's key is the total payload size of every entry the log has received, from offset 0 to the record's end offset. For adjacent records Q and R, the cumulative size of R minus its entry_size_bytes is the cumulative size of Q. Trimming does not change cumulative sizes.
Within a record, entry j holds the offsets [o_j, o_j + record_count_j), where o_0 is start and each later entry starts where the one before it ends. The entry's cumulative size is the record's cumulative size, minus the record's entry_size_bytes, plus the payload sizes of entries 0 through j.
The timestamp of an entry is the creation timestamp of the RAW record that the writer's sequenced put created for it; every entry of that record shares it. Compaction keeps each entry's timestamp in the entry's row (see Compacted objects).
WAL objects
A WAL object holds the entries of one or more logs. Its index region lists each log's entries in a log section, and its payload region holds the entry payloads. The object records no offsets and no timestamps; those come from the index records.
[Writer] A writer MUST write every WAL object in the layout below.
[Reader] A reader MUST parse WAL objects by the layout below.
WAL object
index_length int32, big-endian
index region index_length bytes
index_version int32, big-endian: 3
log section repeated, one per log in the object
log_id int64, big-endian
entry_count varint32
entry repeated entry_count times
payload_offset varint32
payload_size varint32
record_count varint32
payload region from byte 4 + index_length to the end of the object| v3 | Field | Type | Description |
|---|---|---|---|
| required | index_length | int32, big-endian | The length in bytes of the index region, which includes index_version and every log section. At least 4. |
| required | index_version | int32, big-endian | The format version: 3. |
| required | log sections | repeated | One section per log, filling the rest of the index region. |
| required | payload region | bytes | Every byte after the index region. |
Each log section has these fields:
| v3 | Field | Type | Description |
|---|---|---|---|
| required | log_id | int64, big-endian | The log ID. |
| required | entry_count | varint32 | The number of entries in the section. At least 1. |
| required | payload_offset | varint32 | Per entry: the start of the entry's payload, relative to the start of the payload region. |
| required | payload_size | varint32 | Per entry: the length of the entry's payload, in bytes. |
| required | record_count | varint32 | Per entry: the number of records in the entry. At least 1. |
A varint32 is an unsigned LEB128 number: seven bits per byte, least-significant group first, with the high bit set on every byte except the last.
[Writer] A writer MUST encode each varint32 in at most 5 bytes, with a value between 0 and 2³¹ − 1.
[Writer] A log ID MUST appear in at most one log section of an object.
[Writer] The entries of a log section MUST be listed in the order of their offsets.
[Writer] The payload of every entry, payload_size bytes from payload_offset, MUST lie within the payload region.
[Reader] A reader MUST find a log's section by its log_id, and MUST NOT depend on the order of the sections.
[Reader] A reader MUST locate each payload by its payload_offset, and MUST NOT assume that a section's payloads are adjacent, or in section order, in the payload region.
Example. Log 7 with two entries: 2 records in the 5-byte payload 0102030405, then 1 record in the 3-byte payload 060708. The object is 31 bytes:
00000013 00000003 00000000 00000007 02000502 05030101 02030405 060708| Bytes | Field | Value |
|---|---|---|
00000013 | index_length | 19 |
00000003 | index_version | 3 |
0000000000000007 | log_id | 7 |
02 | entry_count | 2 |
00 05 02 | entry 0 | payload offset 0, payload size 5, 2 records |
05 03 01 | entry 1 | payload offset 5, payload size 3, 1 record |
0102030405060708 | payload region | 8 bytes |
Offset index keys
The key of an index record is 62 ASCII bytes: three decimal numbers, each zero-padded to exactly 20 digits, separated by -. The key in the diagram illustrates the format only; see the examples below.
{log_id:020d}-{end_offset:020d}-{cumulative_size:020d}
00000000000000000007-00000000000000000003-00000000000000000000
|----- log_id -----| |--- end_offset ---| |--cumulative_size-|| Field | Description |
|---|---|
log_id | The log ID, L. |
end_offset | The end offset of the record's range. |
cumulative_size | The log's cumulative size at end_offset. |
The record's start offset is its end_offset minus the message_count in its value. Because the fields have a fixed width, the bytewise order of a log's keys is the order of its records.
This specification writes K(L, x, y) for the key with log ID L, end offset x and cumulative size y, and MAX for 9223372036854775807 (2⁶³ − 1). K(L, x, MAX) sorts after every key of L whose end offset is x or less, and before every key whose end offset is greater than x.
[All] An implementation MUST NOT create an index key except by the sequenced put of step A3 in Appending, or by the Put of step C3 in Compaction.
Step C3 replaces the value under an existing key. Its Put creates the key again only if a hard trim, or the deletion of the log, removes the key between the compactor's Get and its Put. A hard trim can remove the key only when the whole range of the compaction lies at or below B, so the record it creates is one that readers never read. A deletion can remove it only after the session behind the compactor's write lease has ended, and the record it creates is then left in a deleted log; Fences and log deletion describes that window and the rule that bounds it.
Writes that choose their own offsets, creating index keys without a sequenced put, are outside format version 3.
Examples. These keys illustrate the format, with the cumulative size set to 0: log 7 with end offset 3, then log 7 with end offset 5.
00000000000000000007-00000000000000000003-00000000000000000000
00000000000000000007-00000000000000000005-00000000000000000000They are not keys of the running example that the other examples on this page use. In that example the RAW record is the log's first record, covering offsets 0 to 2 with 8 payload bytes, so its key is K(7, 3, 8).
Offset index values
The value of an index record is an EntryIndex message in the protobuf (proto2) binary encoding:
syntax = "proto2";
package storage.proto;
enum CompressionType {
NONE = 0;
FastPFOR = 1;
ZSTD = 2;
}
enum FileType {
RAW = 0;
PARQUET = 1;
}
enum IndexType {
NORMAL = 0;
COMPACT = 1;
}
enum Version {
v3 = 3;
}
message ExtraData {
optional string key = 1;
optional string value = 2;
}
message EntryIndex {
optional Version version = 1;
optional string location = 2;
optional FileType file_type = 3;
optional IndexType index_type = 4;
optional int32 message_count = 5;
optional int64 entry_size_bytes = 6;
optional int64 file_size = 7;
optional int64 offset_in_file = 8;
optional int32 entry_count = 9;
optional EntryOffsets entry_offsets = 10;
repeated ExtraData extraData = 11;
}
message EntryOffsets {
optional CompressionType compression_type = 1;
optional int32 uncompressed_size = 2;
optional bytes compressed_payload = 3;
}| v3 | Field | Type | Description |
|---|---|---|---|
| required | version (1) | Version | The format version: v3 (3). |
| required | location (2) | string | Where the entries are; see the next table. |
| required | file_type (3) | FileType | RAW (0): the entries are in a WAL object. PARQUET (1): they are in compacted objects. |
| required | index_type (4) | IndexType | NORMAL (0) or COMPACT (1). |
| required | message_count (5) | int32 | The number of records in the record's range. At least 1. |
| required | entry_size_bytes (6) | int64 | The total payload size, in bytes, of the entries in the range. |
| required | file_size (7) | int64 | See the next table. |
| required | offset_in_file (8) | int64 | See the next table. |
| required | entry_count (9) | int32 | See the next table. |
| optional | entry_offsets (10) | EntryOffsets | Entry boundaries; see Entry offsets. |
| optional | extraData (11) | repeated ExtraData | String key-value pairs. |
The values depend on the file type:
| Field | RAW record | PARQUET record for the range [s, e) |
|---|---|---|
location | The name of the WAL object, relative to the WAL root | The name of the first file in the record's file index |
file_type | RAW | PARQUET |
index_type | NORMAL or COMPACT; readers ignore it | COMPACT |
message_count | The sum of the section's record_count values | e - s |
entry_size_bytes | The sum of the section's payload_size values | c_e - c_s, the cumulative sizes at e and at s |
file_size | The length of the object's whole payload region, which every log section of the object shares; not the length of this log's payloads, and not the length of the object | Equal to entry_size_bytes; not the size of a Parquet file |
offset_in_file | −1 | 0 |
entry_count | The section's entry_count | 1 |
entry_offsets | Present when entry_count is greater than 1 | Absent |
extraData | Absent | Holds the file index under the key CompactedObjectFileIndex |
[Writer] A writer MUST write every RAW record with the values that this table gives for RAW records.
[Compactor] A compactor MUST write every PARQUET record with the values that this table gives for PARQUET records.
[Reader] A reader MUST decide how to read a record's entries from its file_type, and MUST NOT use index_type for that decision.
[All] Every ExtraData pair MUST set both key and value.
[All] The keys of the extraData pairs of one record MUST be distinct.
[Reader] A reader MUST NOT read the entries of a PARQUET record whose extraData has no CompactedObjectFileIndex key.
[Compactor] A compactor SHOULD keep entry_size_bytes below 2³¹; Appendix E gives the reason.
Example: RAW record. The record for the WAL object in the example above, with the illustrative name 7-00000001.wal, is 61 bytes:
0803120e 372d3030 30303030 30312e77 616c1800 20012803 30083808 40ffffff ffffffff ffff0148 02521208 0110081a 0c000000 02000000 00000081 82| Field | Value |
|---|---|
version | 3 |
location | 7-00000001.wal |
file_type | RAW |
index_type | COMPACT, which Ursa writes for a WAL object with one log section (Appendix E) |
message_count | 3 |
entry_size_bytes | 8 |
file_size | 8, the payload region; the object itself is 31 bytes |
offset_in_file | −1 |
entry_count | 2 |
entry_offsets | FastPFOR, uncompressed_size 8, compressed_payload 00000002 00000000 00008182 |
extraData | Absent |
Example: PARQUET record. A record for the range [0, 7), with c_s = 0 and c_e = 28, whose entries are in two compacted files with the illustrative names 7-0.parquet and 7-3.parquet, is 123 bytes:
0803120b 372d302e 70617271 75657418 01200128 07301c38 1c400048 015a5c0a 18436f6d 70616374 65644f62 6a656374 46696c65 496e6465 78124041 41414141 41414141 41494141 41414c4e 7930774c 6e426863 6e46315a 58514141 41414141 41414142 67414141 4173334c 544d7563 47467963 58566c64 413d3d| Field | Value |
|---|---|
version | 3 |
location | 7-0.parquet |
file_type | PARQUET |
index_type | COMPACT |
message_count | 7 |
entry_size_bytes | 28 |
file_size | 28 |
offset_in_file | 0 |
entry_count | 1 |
entry_offsets | Absent |
extraData | CompactedObjectFileIndex = AAAAAAAAAAIAAAALNy0wLnBhcnF1ZXQAAAAAAAAABgAAAAs3LTMucGFycXVldA== |
File index decodes this file index.
Entry offsets
A RAW record that holds n entries, with n greater than 1, carries their boundaries in entry_offsets. The boundaries form the array A of n integers, where A[i] is the sum of the record counts of entries 0 through i, counted from the record's start offset. A[n-1] equals message_count. Entry i holds the offsets from start + A[i-1] up to, but not including, start + A[i], where A[-1] is 0.
| v3 | Field | Type | Description |
|---|---|---|---|
| required | compression_type (1) | CompressionType | FastPFOR (1). |
| required | uncompressed_size (2) | int32 | 4 × entry_count. Informational. |
| required | compressed_payload (3) | bytes | The encoded array A, as described below. |
The encoding is defined by reference to the class IntegratedIntCompressor (me.lemire.integercompression.differential.IntegratedIntCompressor) in JavaFastPFOR 0.2.1 (me.lemire.integercompression:JavaFastPFOR:0.2.1), created with its default constructor. compressed_payload is the array of 32-bit words that IntegratedIntCompressor.compress(A) returns, each word written as a big-endian int32. A reader decodes it by reading the bytes as big-endian int32 words and passing them to IntegratedIntCompressor.uncompress. Appendix D describes the layout of these words; that description is informative.
[Writer] A RAW record whose entry_count is greater than 1 MUST carry entry_offsets, computed from its log section.
[Writer] compression_type MUST be FastPFOR (1), and compressed_payload MUST hold exactly the words that IntegratedIntCompressor.compress(A) returns, each as a big-endian int32. The values NONE and ZSTD are not used in format version 3.
[Reader] A reader MAY derive the entry boundaries from the record's WAL log section instead of from entry_offsets.
A record that holds one entry carries no entry_offsets, and neither does a PARQUET record.
Example. The RAW record above holds two entries of 2 and 1 records, so A is [2, 3] and compressed_payload is these 12 bytes:
00000002 00000000 00008182Appending
An append adds entries to the end of one or more logs. A writer groups entries into a WAL object, writes the object, and then indexes each log section of the object.
[Writer] A writer MUST append entries by the following steps.
- A1: lease. The writer holds a write lease on every log whose entries the object holds; see Write leases and fences.
- A2: object. The writer writes the WAL object under a new name in the WAL root; see WAL objects and Object storage, roots and naming.
- A3: index. After the object has been written, the writer publishes one index record for each log section of the object:
SequencedPut(P(L), p, d1, d2, v), wherepis{log_id:020d},d1is the sum of the section's record counts,d2is the sum of its payload sizes, andvis the section's RAWEntryIndex. - A4: acknowledge. Once a log's index record has committed, the entries of that log's section are appended.
The key that the sequenced put returns gives the record's end offset e and cumulative size c. The section's entries take the offsets from e - message_count up to e, in section order, and share the record's creation timestamp as their timestamp.
[Writer] A writer MUST NOT publish an index record for a WAL object before the object has been written.
[Writer] If writing the object fails, the writer MUST NOT publish an index record for it, and MUST NOT report any of its entries as appended.
[Writer] A writer MUST NOT report an entry as appended before the index record of its log section has committed.
Index records for different logs of one object are published independently; the object does not make them atomic. One log's entries can be appended while the index put of another log in the same object fails.
Offsets follow the order in which the store commits sequenced puts, not the order in which writers issue them. The format orders two appends to one log only when the second one's sequenced put is issued after the first one's has committed.
A sequenced put whose outcome is unknown, for example after a timeout, can have committed. The format has no deduplication: an entry that a writer reports as failed can still be readable, and appending it again stores it twice.
Index puts carry no condition, and a sequenced put cannot carry one. Fencing works through write leases, not through conditional index writes.
Write leases and fences
A write lease tells a catalog that a writer or compactor is changing a log. A write fence, which a catalog installs when it deletes a log, stops new leases. The fence, and the lease record that announces each lease, live in the log's partition; Metadata store model gives their keys and values.
Leases
An implementation holds a write lease from the time it acquires the lease, by the two rules below, until it releases the lease by deleting its lease record. If the session that created the lease record ends, the store removes the record, but that does not release the lease in this sense. The last two rules of Fences and log deletion address that case for a writer's index records, and for a compactor's compaction index updates and hard trims.
[Writer, Compactor] To acquire a write lease on log L, an implementation MUST, in this order, create its lease record with CreateIfAbsent(P(L), ...) as an ephemeral record, under a token that no other lease of L uses, and then read the write fence of L with Get, and MUST NOT use the lease unless that Get succeeded and found no fence.
[Writer, Compactor] If the fence exists, or the Get fails, the implementation MUST delete its lease record and MUST NOT change the log.
[Writer, Compactor] An implementation MUST hold a write lease on L while it publishes an index record of L, writes the trim marker of L, performs steps C3 and C4 of a compaction of L, or hard-trims L.
[Writer, Compactor] An implementation MUST NOT delete a lease record before every change it started under that lease has completed.
Fences and log deletion
[Catalog] To fence log L, a catalog MUST create the write fence of L with CreateIfAbsent. A fence that already exists counts as success.
[All] An implementation MUST NOT delete or change a write fence. A fenced log gets no new leases, and its log ID is never reused; the end of this section describes the writes that leases acquired before the fence can still make.
[Catalog] To delete log L, a catalog MUST fence L; then list the lease records of L, the keys that begin with /stream-write-leases/{L}/, until none remains; and only then delete every index record of L with DeleteRange(P(L), K(L, 0, 0), K(L, MAX, MAX)).
[Catalog] If lease records remain when the catalog stops waiting for them, it MUST leave the log's records and objects unchanged, and MUST leave the fence in place.
The lease and the fence are in the same partition, whose operations are linearizable. So, when the writer's Get of the fence succeeds, either it reads the fence after the fence was created, and gives up its lease, or the catalog lists the lease records after the writer's lease record was created, and waits for it.
A fence stops new leases; it does not revoke leases that already exist. Because index puts carry no condition, a writer that acquired a lease before the fence was installed can keep publishing index records while it holds the lease. Deleting a log waits for these leases to drain. The catalog counts a lease as drained once its lease record is gone, which also happens when the session that created the record ends while a writer or compactor still holds the lease. That writer or compactor can still change the log after the drain, until it learns that the session has ended. In that window, in a log that a catalog has just fenced and deleted, a writer can still publish index records, and a Put in a compactor's step C3 can re-create an index key. The next two rules bound that window for a writer's index records, and for a compactor's compaction index updates and hard trims.
[Writer] A writer SHOULD stop publishing index records for L as soon as it learns that the session that created its lease record on L has ended.
[Compactor] A compactor SHOULD stop the compaction index update (C3–C5) and any hard trim once it learns that the session behind its write lease has ended.
Reading
A reader needs no write lease. To read log L from offset o:
- R1: trim marker. [Reader] A reader MUST read
Bfrom the trim marker ofL, takingBas 0 when the marker is absent, and MUST start reading ato′, the greater ofoandB. - R2: covering record. [Reader] A reader MUST read the entries at
o′from the index record ofLwith the smallest end offset greater thano′.Higher(P(L), K(L, o′, MAX))returns that record; a result whose key does not begin with{log_id:020d}-forLmeans that no record holdso′, becauseo′is at or past the end of the log. - R3: dispatch. [Reader] A reader MUST read a RAW record's entries from its WAL object, and a PARQUET record's entries from compacted objects.
- For a RAW record, the reader reads the object named by
locationin the WAL root, finds the log section ofL, and assigns offsets to the section's entries from the record's start offset (see Logs, entries and offsets). Every entry takes the record's creation timestamp. - For a PARQUET record, the reader follows Reading a PARQUET record.
- For a RAW record, the reader reads the object named by
- R4: entries. Reads return whole entries. The first entry returned is the entry that holds
o′, and it can start beforeo′. [Reader] A reader MUST NOT return an entry that ends at or belowB. A following read starts where the last entry returned ends: its offset plus its record count.
The log's first index record above the trim marker is Higher(P(L), K(L, B, MAX)), or Higher(P(L), K(L, 0, 0)) when the marker is absent, with the same prefix check as in step R2. Its start offset can be below B.
[Reader] A reader SHOULD report the greater of that record's start offset and B as the log's first offset.
The log's last index record is Floor(P(L), K(L, MAX, MAX)), with a result whose key does not begin with {log_id:020d}- for L meaning that the log has no records. The next append to commit starts at its end offset.
[Reader] When the object that a record names does not exist, a reader SHOULD look up the covering record again before it reports an error. A compaction or a hard trim can have replaced or removed the record it used.
Compacted objects
A compacted object is an immutable Parquet file that holds part of a compacted range, one row per entry, together with its companion .index file. The structure in this section is protocol-neutral. A binding (see Appendix A) defines two identifiers: the value of the footer key entrySerDeType, and the name of the Avro record that describes a row.
Object storage, roots and naming gives the directory and the names of compacted files and companions.
Parquet files
[Compactor] Each compacted file MUST be a Parquet file with exactly one row per entry, holding consecutive entries of one log in ascending offset order.
[Compactor] A compactor MUST write each compacted file with an Avro record schema that has one field, payload, of type ["null", "bytes"], under the record name that the binding defines. The Parquet schema of the file then has one nullable BINARY column, payload.
[Compactor] The payload of each row MUST hold the entry's payload unchanged.
The file's footer carries this key-value metadata:
| v3 | Key | Value |
|---|---|---|
| required | schemaType | avro |
| required | entrySerDeType | The binding's identifier, for example KAFKA_BATCHED_RAW_PARQUET |
| required | rowMetadata | The name of the companion file. Readers do not use it. |
The Parquet writer can add keys of its own.
[Reader] A reader MUST reject a compacted file whose entrySerDeType is not the identifier of a binding that the reader implements.
Companion files
A companion file is one gzip stream. Row i of the companion describes row i of the Parquet file, with rows numbered from 0 across the blocks in order.
companion file (gzip stream)
block repeated
length int32, big-endian: N
rows N bytes of UTF-8 JSON: an array of objects, one per row
sentinel int32, big-endian: 0x7FFFFFFF
offset map
length int32, big-endian: N
map N bytes of UTF-8 JSON: an object from entry start offset to row| Part | Content |
|---|---|
| block | A JSON array that holds one JSON object per row, in row order. Each object maps string keys to string values. The number of rows in a block is implementation-defined. Every block except the last holds at least one row; the last block can be empty. |
| sentinel | The int32 0x7FFFFFFF, which ends the blocks. |
| offset map | A JSON object that maps the start offset of each entry in the file, as a decimal string, to the number of its row. |
[Compactor] A compactor MUST write every companion file in this layout.
[Compactor] Every block except the last MUST hold at least one row.
[Reader] A reader MUST accept an empty last block.
Each row object has these keys:
| v3 | Key | Value |
|---|---|---|
| required | metadata | A serialized LakehouseEntryMetadata message, in standard Base64 (RFC 4648, with padding). |
| required | offset | The entry's start offset, in decimal without leading zeros. |
The keys of the offset map are the offset values of the rows.
LakehouseEntryMetadata is a protobuf (proto3) message:
syntax = "proto3";
package materialization.proto;
message EntryHeader {
optional int64 offset = 1;
optional int32 numberOfMessages = 2;
optional int64 writtenTimestamp = 3;
optional int32 entrySize = 4;
optional int64 cumulativeSize = 5;
}
message LakehouseEntryMetadata {
reserved 1, 3;
optional EntryHeader entryHeader = 2;
optional int64 schemaVersion = 4;
}| v3 | Field | Type | Description |
|---|---|---|---|
| required | entryHeader (2) | EntryHeader | The entry's header, as below. |
| optional | schemaVersion (4) | int64 | A schema version, for bindings that record one. The Kafka binding does not. |
| required | offset (1) | int64 | In EntryHeader: the entry's start offset. |
| required | numberOfMessages (2) | int32 | In EntryHeader: the entry's record count. |
| required | writtenTimestamp (3) | int64 | In EntryHeader: the entry's timestamp. |
| required | entrySize (4) | int32 | In EntryHeader: the entry's payload size, in bytes. |
| required | cumulativeSize (5) | int64 | In EntryHeader: the log's cumulative size at the end of the entry. |
[Compactor] The entryHeader of each row MUST equal the header of the entry as read from its RAW record: its offset, record count, timestamp, payload size and cumulative size (see Logs, entries and offsets).
[Reader] A reader MUST take the timestamp of an entry read from a compacted object from the entry's row, not from the PARQUET record.
[Reader] To find the row of the entry that holds offset o, a reader MUST look up the decimal string of o in the offset map and, when o is not the start offset of an entry, MUST use the row of the greatest start offset below o, comparing the keys as numbers.
File index
A PARQUET record lists its compacted files in its file index, the extraData pair whose key is CompactedObjectFileIndex. The value is the standard Base64 encoding (RFC 4648, with padding) of these fields, repeated once per file, in offset order:
| v3 | Field | Type | Description |
|---|---|---|---|
| required | last_offset | int64, big-endian | The last offset that the file holds, inclusive. |
| required | name_length | int32, big-endian | The length of name, in bytes. |
| required | name | bytes | The file's name in UTF-8, relative to the directory of the log's stream. |
[Compactor] The files of a file index MUST appear once each, in ascending last_offset order, and the last_offset of the last file MUST be e - 1 for the range [s, e).
The file that holds offset o is the first file whose last_offset is at or above o. A file's base offset is the last_offset of the file before it plus 1, or the record's start offset s for the first file.
Example. The PARQUET record above names two files: 7-0.parquet holds offsets 0 to 2, and 7-3.parquet holds offsets 3 to 6. Before Base64 the file index is 46 bytes:
00000000 00000002 0000000b 372d302e 70617271 75657400 00000000 00000600 00000b37 2d332e70 61727175 6574Its Base64 form, the value in extraData:
AAAAAAAAAAIAAAALNy0wLnBhcnF1ZXQAAAAAAAAABgAAAAs3LTMucGFycXVldA==Reading a PARQUET record
To read offset o′ from a PARQUET record of log L for the range [s, e):
- Decode the file index, and select the file that holds
o′. - Open the file in the directory of the log's stream, and check its
entrySerDeType. - Read the offset map of its companion file, and find the row of the entry that holds
o′. - Read rows from that row onward. Each row gives one entry: its header from
metadata, and its payload frompayload.
[Reader] A reader MUST read a PARQUET record's entries by these steps.
Compaction
A compaction rewrites the entries of a range of RAW records into compacted objects, and replaces those records with one PARQUET record. The index update is not atomic. It runs in steps, and it installs the new record before it deletes the old ones, so every offset of the range is covered by a record at every point.
A compaction of log L covers the range [s, e). The last record it replaces, R_k, has the key K(L, e, c_e); R_1 to R_(k-1) are the other records it replaces, and c_s is the log's cumulative size at s.
- C1: range. [Compactor] The range MUST start at offset 0 or at the end offset of an earlier index record of
L, and MUST end at the end offset of a RAW record; every index record ofLwhose end offset lies in(s, e]MUST be a RAW record, and the compaction replaces all of them. - C2: data. [Compactor] Every compacted file of the range and its companion MUST be completely written to object storage before step C3.
- C3: install. [Compactor] The compactor MUST read the key
K(L, e, c_e)withGet; if the key exists, it MUST replace the key's value with the PARQUET record'sEntryIndexbyPut, keeping the key, and if the key no longer exists, it MUST skip this step. - C4: retire. [Compactor] After step C3, or after skipping it, the compactor MUST delete the other records it replaces with
DeleteRange(P(L), K(L, s + 1, 0), K(L, e - 1, MAX)), which removes every record ofLwhose end offset lies in[s + 1, e - 1]. - C5: progress. [Compactor] A compactor MAY record its compaction progress; see Appendix B.
Step C3 keeps the key of R_k, so the PARQUET record keeps the creation timestamp of R_k, and the log's greatest key is unchanged when R_k is the log's last record. The key of R_k is absent only when a hard trim has removed it, which means the whole range lies at or below B, or when the log has been deleted. A hard trim or a log deletion that runs between the Get and the Put of step C3 lets the Put create the key again; Offset index keys describes both cases.
[Compactor] A compactor MUST acquire a write lease on L before each attempt at step C3 and hold it until that attempt's step C4 has completed or the attempt has failed.
[Compactor] Compaction index updates of one log MUST NOT run concurrently with each other. How compactors coordinate this is implementation-defined; Appendix E describes Ursa's approach.
[Compactor] A compaction MUST NOT replace or delete a PARQUET record that an earlier compaction installed.
[Compactor] The compacted files of a range MUST hold one row for every entry of the range that ends above B, in ascending offset order.
[Compactor] A compactor MAY leave out of the compacted files the entries of the range that end at or below B.
[Compactor] If a compaction fails after step C3 has taken effect, the compactor MUST repeat steps C3 and C4, with the same range and the same value, until step C4 has completed, or until it can no longer acquire a write lease on L because L has been fenced for deletion. Both steps are idempotent, and the conditions of step C1 apply only when a compaction starts. When L has been fenced, deleting it removes the records that the compaction left before the deletion. A Put of step C3 that lands after the deletion can still leave a record in the deleted log; see Logs, entries and offsets, Offset index keys and Fences and log deletion.
Between steps C3 and C4 the records R_1 to R_(k-1) overlap the new PARQUET record. A lookup in that interval returns either one of those RAW records or the PARQUET record, and for every offset at or above B both hold the same entries at the same offsets.
Relation to the leaderless log protocol (informative). The compaction index update of the leaderless log protocol has the same order:
| Protocol action | Step |
|---|---|
CompactStart | C1 and C2 |
CompactWriteCompactedIndex | C3 |
CompactDeleteOldEntries | C4 |
CompactUpdateCursor | C5, implementation-defined |
The protocol's model checks that this order leaves no reader without a covering record (its property NoReaderError). The protocol numbers offsets from 1 and keys a record by its inclusive end offset, which it finds with a ceiling lookup. This specification numbers offsets from 0 and keys a record by its exclusive end offset, which it finds with a strictly-higher lookup. Under that translation, a record's key and the keys that step C4 deletes are the same in both.
Trim, deletion and reclamation
Trim marker
Each log has at most one trim marker, the record /mark-deleted-offsets/{L} in P(L), whose value is B in ASCII decimal digits. When the marker is absent, B is 0. Readers skip offsets below B, and the data that only those offsets use can be deleted.
[All] An implementation MUST NOT decrease B.
Soft trim
[Writer] A soft trim through offset o, inclusive, MUST compute B′ as the smaller of o + 1 and the start offset of the log's last index record, and, when B′ is greater than the current B, MUST write B′ with CreateIfAbsent if no marker exists, or otherwise with CompareAndSet on the version it read.
A soft trim of a log that has no index records changes nothing. When the conditional write fails because another trim changed the marker, the writer can read the marker again and retry. Capping B′ at the start of the log's last record keeps that record, and so the log's greatest key, in place: the next sequenced put continues from it.
Hard trim
A hard trim to offset X removes every index record of L that ends at or below X: DeleteRange(P(L), K(L, 0, 0), K(L, X, MAX)).
[Compactor] A hard trim MUST use an X at or below B, and MUST NOT remove the log's last index record.
Deleting compacted objects
[Compactor] A compactor MAY delete a compacted file and its companion once every index record whose file index names the file ends at or below B.
[Compactor] A compactor SHOULD delete such files before it hard-trims the records that name them, so that a failed deletion can be retried from those records.
[All] An implementation MUST NOT delete a compacted file that a compaction in progress can still install, unless the range of that compaction ends at or below B.
Readers never read a record whose range ends at or below B, so deleting its files leaves the log readable even when a retry of step C3 installs the record again afterwards. Offset index keys describes the same case for index keys.
Reclaiming compacted files that no index record names, such as the output of a compaction that failed before step C3, or the files of a deleted log, is implementation-defined.
Reclaiming WAL objects
[All] An implementation MUST NOT delete a WAL object while an index record of any log names it.
[All] After the last index record that names a WAL object has been removed, an implementation MUST keep the object for a grace period before it deletes it.
The length of the grace period, and the mechanism that finds and deletes objects, are implementation-defined. A reader that resolved a RAW record before a compaction or a trim removed it can still read the record's object during the grace period. Appendix E describes Ursa's mechanism.
Deleting a log
Deleting a log, as Write leases and fences describes, removes its index records and leaves its write fence. It deletes no objects: WAL objects and compacted files that only the log's records named are then unreferenced, and the rules above govern their reclamation.
Catalog semantics
A catalog creates and deletes streams and their logs. The serialization of its records is implementation-defined; Appendix B describes Ursa's. Format version 3 defines the behavior that other roles depend on.
[Catalog] A catalog MUST give every log a log ID that no other log of the deployment has had, and MUST NOT reuse a log ID, even after its log has been deleted.
How log IDs are allocated is implementation-defined. Interoperation in format version 3 means appending to, reading and compacting logs that already exist: an implementation without the Catalog role works on logs that a catalog has created.
[Catalog] Opening an existing log MUST NOT allocate a log ID, register a stream partition, grow a stream, or change catalog metadata.
[Catalog] Creating a stream MUST fail when its identity, the pair of namespace and stream name, has been dropped.
[Catalog] Dropping a stream MUST record a durable tombstone for its identity, even when no live stream has that identity.
[Catalog] Adding stream partitions MUST publish the stream's new layout atomically, and only after the log of every new partition exists.
[Catalog] Deleting a log MUST be idempotent: deleting a log that is already absent succeeds.
[Catalog] In format version 3, a stream namespace MUST NOT contain /.
[Catalog] In format version 3, a stream name MUST NOT contain /.
A reader of compacted objects learns from the catalog which stream a log belongs to; see Object storage, roots and naming.
Appendix A: Kafka binding
Ursa for Apache Kafka (UFK), a Kafka distribution built on the Lakestream API and specification, stores Kafka topics with this binding. UFK 4.3.1 implements it.
Stream identity
Each incarnation of a Kafka topic is one stream:
- The namespace is
default. - The stream name is
{topic}-topic-id-{topic ID}: the topic name, the text-topic-id-, and the Kafka topic ID in its string form. A topic that is deleted and then created again under the same name gets a new stream. - Partition
pof the topic is the log at positionpof the stream's layout.
The stream's properties include lakestream.source.logical.name and lakestream.kafka.topic.name, which both hold the topic name. The compacted directory of a topic is therefore {compacted root}/default/{topic}-topic-id-{topic ID}/.
Entry payloads
[Writer] A Kafka entry payload MUST be exactly the bytes of one or more complete Kafka record batches, in the Kafka record-batch format, with no envelope and no prefix, and the entry's record count MUST equal the total number of records in those batches.
[Writer] Each record batch of magic 2 or later MUST have base offset 0, and MUST declare a record count greater than 0 that equals its offset range, its last offset minus its base offset plus 1.
[Writer] Each record batch of magic 0 or 1 MUST hold at least one record.
[Writer] A Kafka entry payload MUST NOT contain a control batch or a transactional batch.
[Reader] A reader MUST assign offsets to the batches of an entry in order, starting at the entry's offset: each batch takes as many offsets as it has records, after the offsets of the batches before it.
Before it returns batches to a Kafka client, UFK's reader rewrites their offsets to these values.
Compacted object identifiers
| Identifier | Value |
|---|---|
entrySerDeType | KAFKA_BATCHED_RAW_PARQUET |
| Avro record name | KafkaMessage |
The payload of each row holds the entry's record batches.
KAFKA_PARQUET
Ursa can also write compacted objects whose entrySerDeType is KAFKA_PARQUET. That layout stores one row per Kafka record, decoded through a schema registry, and does not preserve the entry payload byte for byte. Format version 3 does not cover it.
Entry payloads with a length prefix, written before this binding stored raw record batches, are out of scope; see Appendix C.
Appendix B: Oxia mapping
This appendix is informative. It maps the primitives of Metadata store model to Oxia, which Ursa 1.0 uses as the metadata store, and lists the records that Ursa keeps in Oxia beyond those this specification defines. Those other records, and the compaction cursor, are Ursa's choices; format version 3 does not define them.
Primitives
| Model | Oxia |
|---|---|
Partition P(L) | The Oxia partition key. Every operation on a record of log L passes the option PartitionKey with the decimal log ID, which places all of the log's records in one Oxia shard. |
| Linearizable partitions | The operations of one shard are linearizable. |
| Key order | With its default hierarchical key sorting, Oxia orders keys first by the number of / they contain, and keys with the same number bytewise, with each / compared as the byte 0xFF. Keys without / therefore compare bytewise and sort before every key that contains /. A key that ends in // counts one / fewer, and Oxia's internal keys, under __oxia/, sort after all others; neither detail affects these two properties. Both properties also hold under Oxia's earlier comparer, which compared keys segment by segment at each /. |
| Version | The record's Version.versionId. |
| Creation timestamp | The record's Version.createdTimestamp, in milliseconds. A put that replaces a value keeps it. |
| Sessions | The Oxia client session. Oxia removes a session's ephemeral records when the session ends. |
Get(P, k) | get(k) |
Floor, Ceiling, Higher | get(k) with ComparisonFloor, ComparisonCeiling or ComparisonHigher |
Scan(P, a, b) | rangeScan(a, b) |
List(P, prefix) | list(prefix, end), where end is prefix followed by the character U+FFFF |
SequencedPut(P, p, d1, d2, v) | put(p, v) with SequenceKeysDeltas([d1, d2]). Oxia rejects a sequenced put without a partition key, with an expected version, or with a first delta of 0, and returns the key it created. |
Put(P, k, v) | put(k, v) without a condition |
CreateIfAbsent(P, k, v) | put(k, v) with IfRecordDoesNotExist, plus AsEphemeralRecord for an ephemeral record |
CompareAndSet(P, k, version, v) | put(k, v) with IfVersionIdEquals(version) |
Delete(P, k) | delete(k) |
DeleteRange(P, a, b) | deleteRange(a, b) |
Beyond the procedures of this specification, Ursa uses Ceiling to find a log's first record, trimmed or not, and Scan to read runs of index records, in batched reads and in its compacted-data cleaner.
Ursa's other records
| Key | Purpose |
|---|---|
/stream-id-generator | Log-ID allocation. Ursa writes this record with an empty value, and uses the version ID that Oxia returns as the new log ID. Log IDs therefore increase but are not contiguous. |
/stream-id-generator/{key} | Keyed log IDs, written with the partition key /stream-id-generator. The value is a decimal log ID, or a JSON object with the log ID, an owner and a state. |
/stream-id/{L} | Registered logs. The value is JSON that holds the log's key; the key of log i of a stream is {namespace}/{stream name}-partition-{i}. |
/streams/..., /admin/streams/... | Catalog records: streams and their partitions, namespaces (_namespaces), table catalogs (_tablecatalogs) and tombstones (_tombstones). |
/first-uncompacted-offset/{L} | The compaction cursor. |
/compact-stream-tasks/..., /compact-stream-tasks-dlq/..., /task-locks/... | Compaction tasks and the claims on them. |
/compact/leader | An ephemeral record, holding a host name, that elects the compaction commit leader. |
mark-delete-{L}-{cursorId}, individual-acks-... | Cursor state. |
ursa-wal-delete-marker, ursa-wal-cleanup-lock | State of the WAL garbage collector: the last WAL name it reached, and an ephemeral lock. |
Compaction cursor
Ursa's compactor writes /first-uncompacted-offset/{L} after step C4, under the same write lease, as ASCII decimal digits and without a partition key. For a compacted range [s, e) it writes e + 1. Ursa's WAL garbage collector reads it; see Appendix E.
Appendix C: Format versions 1 and 2
This appendix is informative. It describes the encodings that Ursa wrote before format version 3, so that a reader of an older deployment can recognize them.
Legacy WAL objects
WAL object (versions 1 and 2)
index_length int64, big-endian
index region index_length bytes of 32-byte tuples
log_id int64, big-endian
entry_id int64, big-endian
payload_offset int64, big-endian
payload_length int64, big-endian
payload region the rest of the objectThe tuples appear in no particular order. entry_id numbers the entries of the whole object from 0, across all of its logs, in the order they were added. payload_offset is relative to the start of the payload region.
Legacy index records
Each entry has an index record of its own, created with the same sequenced put as in version 3, with the entry's record count and payload size as the deltas, so the keys have the version 3 format. The values are UTF-8 strings:
version 1: {message_count:020d}-{entry_size_bytes:020d}-{location}-{entry_id}
version 2: {message_count:020d}-{entry_size_bytes:020d}-{location}-{entry_id}-{file_type}file_type is RAW or PARQUET. A parser reads the two numbers from characters 0 to 19 and 21 to 40, and takes everything from character 42 as the position. It splits the position's trailing fields off at the last -, because a location can contain -: in version 1 that is entry_id; in version 2 it is file_type, and then entry_id at the next - to the left. The values carry no entry count, since each record holds one entry, no file size, and no index type; Ursa treats a PARQUET record as COMPACT and a RAW record as NORMAL.
Ursa 1.0 writes version 2 values when it is configured with a format version below 3.
Telling the versions apart
Ursa 1.0 does not detect the version from the bytes. A version 3 index value that Ursa writes begins with the bytes 08 03, while a version 1 or 2 value begins with an ASCII digit. A legacy WAL object begins with an 8-byte length whose four high-order bytes are 0 for any index region shorter than 4 GiB, while a version 3 object begins with a 4-byte length of at least 4.
Legacy compacted data
Two older forms of compacted data are out of scope for format version 3:
- PARQUET records whose
extraDatahas noCompactedObjectFileIndexkey. Ursa's standalone Kafka reader rejects them. - Kafka entry payloads with a length prefix, written before the Kafka binding stored raw record batches. Their compacted objects carry the same
entrySerDeTypeas current ones, so a reader cannot tell them apart by metadata.
Appendix D: Integrated FastPFOR encoding
This appendix is informative. The class IntegratedIntCompressor in JavaFastPFOR 0.2.1 defines the encoding of entry_offsets; this appendix describes the words it returns for an array A of n increasing non-negative integers, so that other implementations can check their output. Entry offsets says how the words are stored.
The encoder works on deltas: the first delta is A[0], and each later delta is A[i] - A[i-1]. It writes 32-bit words in this order:
- Length. Word 0 is
n. - Binary packing. The first
32 × floor(n / 32)deltas are bit-packed in blocks of 32.- For each group of four blocks (128 deltas), one header word holds the four bit widths
b0,b1,b2andb3, wherebjis the bit length of the largest delta in blockj:b0in bits 24 to 31,b1in bits 16 to 23,b2in bits 8 to 15, andb3in bits 0 to 7. The packed words of the four blocks follow,b0words for block 0, thenb1,b2andb3words. - For each remaining block of 32 deltas, one header word holds its bit width
b, andbpacked words follow. - In a block of bit width
b, deltajof the block occupies bitsj × btoj × b + b - 1of the block's packed words, counted from the least-significant bit of its first word.
- For each group of four blocks (128 deltas), one header word holds the four bit widths
- Marker. When
nis less than 32, there is no binary-packed part, and a single 0 word follows word 0 instead. - Variable-byte tail. The remaining
n mod 32deltas are written as bytes: each delta in 7-bit groups, least-significant group first, with the high bit (0x80) set on the last byte of the delta and clear on the others. This is the reverse of the LEB128 convention. The bytes are padded with zero bytes to a multiple of 4 and read as little-endian 32-bit words.
compressed_payload then stores every word as a big-endian int32.
A | Words |
|---|---|
[2, 3] | 00000002 00000000 00008182 |
[1, 2, 3, …, 40], the cumulative counts of 40 one-record entries | 00000028 00000001 ffffffff 81818181 81818181 |
For [2, 3]: n is 2; with no binary-packed part, a 0 word follows; the deltas 2 and 1 become the bytes 82 and 81, padded to 82 81 00 00 and read little-endian as 00008182.
For the 40 entries: n is 40 (0x28); the first 32 deltas are all 1, so one block has bit width 1 (the header word 1) and a single packed word with all 32 bits set (ffffffff); the last 8 deltas are also 1, each the byte 81, which fill the two words 81818181.
Appendix E: Implementation Notes
This appendix is informative. It describes how Ursa 1.0 implements this specification where the specification leaves a choice open.
Format versions
The StorageConfig setting indexSerializeFormatVersion, 3 by default, selects both the WAL object layout and the encoding of index values, for writing and for reading. Ursa does not detect either from the bytes, and does not check the version numbers it reads. Its protobuf parser skips fields it does not recognize.
Writing WAL objects
- Ursa collects appends for up to 250 ms, then writes them in WAL objects of up to 4 MiB, each holding the entries of at most 4 logs. An entry larger than 4 MiB gets a WAL object of its own.
- The payload region is the concatenation of the entry payloads in the order they were added, with no padding, checksum or trailer. Log sections appear in hash-map order.
- Ursa writes the
index_typeof a RAW record asCOMPACTwhen the WAL object has exactly one log section, and asNORMALotherwise. - WAL object names have the form
yyyy/MM/dd/HH/mm/ss/{UUID}, with the date and time from the local clock and time zone of the process that writes the object, and a random UUID. - The S3 backend sends a CRC32C checksum with each PUT and validates checksums on each GET.
- After a WAL object's PUT succeeds, Ursa issues the index puts of all its logs together, and handles objects in the order they were flushed without waiting for one object's index puts to commit before it issues the next object's.
- Ursa's Oxia client pipelines puts in order per shard: it sends the puts for a shard over one write stream, in the order they are issued, and Oxia commits the requests of a stream in the order it receives them. One Ursa writer's appends to a log therefore take offsets in the order the writer issues their index puts. This specification does not require that order: Appending orders two appends to a log only when the second one's sequenced put is issued after the first one's has committed.
- Before each index put, Ursa also checks a fenced state that it keeps in the memory of the process; the durable protection is the lease and fence protocol.
Reads
- Ursa reads a WAL object whole and keeps recently read objects in a cache.
- It keeps index records in a cache for up to 600 seconds (
entryIndexCacheTTLInSecs), so a read can use an index record for up to that long after a compaction or a trim has removed it from the store. A deployment's WAL grace period therefore needs to exceed 600 seconds for Ursa's readers. This specification does not set the length of the grace period; see Reclaiming WAL objects. - For records with several entries, Ursa's
EntryIndexAPI estimates each entry's payload size by dividingentry_size_bytesevenly; the exact sizes are in the WAL log section. - Ursa builds entry headers with a 32-bit payload size, which is why a compactor keeps
entry_size_bytesbelow 2³¹.
Leases and deletion
A catalog waits up to 30 seconds for leases to drain when it deletes a log, and polls the lease records every 25 ms.
Compaction scheduling
- Ursa publishes compaction work for each log as tasks over ascending, adjacent ranges, starting at offset 0. A range closes when it holds 256 MiB of payload, or when its oldest entry has been in the log for 180 seconds.
- Before it schedules compaction for a stream, Ursa looks up the stream's source topic in a schema registry. A stream whose schema lookup fails, or whose schema type Ursa does not support, is skipped for a quarantine period.
- One commit runner per deployment, chosen by leader election on
/compact/leader, performs steps C3 and C4. Its commit interval is 180 seconds by default. - Compacted files are named
ursa-{UUID}.parquet, with companionsursa-{UUID}.index. Ursa writes them with Parquet's Avro binding, starts a new file when the record schema changes, and writes 10,000 rows per companion block by default (ursaIndexFileWriterMaxBufferedRecords). The compacted root comes from the compaction bucket and prefix settings.
Trimming
When the conditional write of a soft trim fails because another trim changed the marker, Ursa reads the marker again and retries, up to three times. If the write still fails, or fails for another reason, Ursa logs a warning and returns the computed B′ as the soft trim's result anyway, so the result does not show that the marker advanced. Implementation status lists this against the API contract.
Reclamation
- Compacted objects. Every 12 hours, Ursa's compacted-data cleaner finds, for each log with a trim marker, the PARQUET records that end at or below
B. It deletes the files that their file indexes name, with their companions, and then hard-trims the log to the end of the last such record, under a write lease. - WAL objects. Every 12 hours, under the ephemeral lock
ursa-wal-cleanup-lock, Ursa's WAL garbage collector reads, for each registered log, the index record that covers the log's compaction cursor (Appendix B), and parses the date and time at the start of the name of the WAL object that the record names. The earliest of these is the new horizon; the WAL name that the previous run reached, kept inursa-wal-delete-marker, gives the old one. The collector then installs object-storage lifecycle rules that expire the objects whose names begin with any of a set of prefixes: the old horizon's date and time, and each later time in one-hour steps from it, provided that the time is before the new horizon, each written in the fullyyyy/MM/dd/HH/mm/ssform. The collector therefore assumes that every WAL object in the deployment has a name with Ursa's date prefix.
Specification
Lakestream is specified by two documents: the Storage Spec, which defines format version 3, and the Materialization Spec.
Lakestream Materialization Spec
The Materialization Spec defines the policy model of the stream materialization framework: policies, resolution, table naming and table catalogs.