This page introduces the concepts underlying K2.
The core primitive in K2 is what we call a stream — an ordered, durable, log of events. The stream sits between producers writing events and consumers reading them. Unlike a traditional queue, where consumption removes items, production and consumption in a log are completely decoupled. Writes append events to the log, while reads merely advance a pointer (or offset) within the log.
This has some useful properties:
- Writes and reads are completely independent, so we never run out of space or otherwise block writes due to slow reads
- We can support multiple independent readers consuming the entire stream (pub-sub style) as readers do not affect each other or the log
- We can support historical replay, as data is only removed based on a configurable time-to-live (TTL)
Each stream has a name and a retention period. When you create a stream, K2 assigns it a unique ID, which producers and consumers use to communicate with it. Producers write to a stream through an HTTP endpoint (with or without authentication), a Workers binding, or both.
Records are the data written to a stream. Each record has the following fields:
content: a binary message, containing arbitrary data; base64-encoded in the HTTP APIsheaders: an optional map of string keys to string values that can be used to describe the data in the content
Headers are useful for describing the content, without needing to deserialize it. Common use cases for headers include:
- Storing the encoding, so the reader knows how to deserialize the content
- Representing data used to route or filter events, improving efficiency by avoiding deserialization when not necessary
- Annotating content in a pass-through pipeline without needing to modify the underlying data
Records can be up to 1 MB, counting across both content and headers.
Reads from K2 are performed via subscriptions. Each subscription will receive all messages in the stream, and multiple consumers can share a single subscription.
This enables K2 to support two delivery strategies: reads can be shared amongst a set of consumers (such that each consumer gets a subset of the messages), or delivered to all consumers independently (such that each consumer gets all messages). These strategies can also be mixed, with multiple groups of consumers which each get a subset of the messages.
Subscriptions are created with an initial position in the log: either earliest, which receives
all retained (not deleted according to the TTL) data in the stream, or latest which receives
all events from the time the subscription is created.
Once a subscription is created, clients can consume from it by POSTing to the subscription's
/consume endpoint. This returns a list of messages which this client is expected to process.
These messages are leased to that client for a particular amount of time — the lease period,
which is 5 minutes. The client is expected, before the lease expires, to either ack the messages,
telling the subscription that they are successfully consumed, or nack them, indicating a processing
failure. Events owned by a nack'd or timed-out lease will be redelivered on a subsequent call to
consume.