HTTP and binary API

Every operation is an HTTP call that sends and receives JSON. Produce and consume are also available over a binary protocol, for services that need the throughput.

Conventions

  • Requests and responses are JSON. An error is {"error": "…"} with a 4xx or 5xx status.
  • A 500 says internal error (id …); the detail is in the broker's log under that id, never in the response.
  • Request bodies are limited to 10 MiB (--max-request-body-bytes) and each record to 8 MiB (--max-record-bytes).
  • A panic in a handler becomes a 500; the broker keeps running.

A first session

shell
# Create a topic with 2 partitions
curl -s -X POST localhost:9000/topics \
  -H 'content-type: application/json' \
  -d '{"name":"orders","partitions":2}'

# Produce two records; the key picks the partition
curl -s -X POST localhost:9000/topics/orders/produce \
  -H 'content-type: application/json' \
  -d '{"records":[{"key":"user-1","value":"checkout"},{"key":"user-2","value":"signup"}]}'
{"results":[{"partition":0,"offset":0},{"partition":1,"offset":0}]}

# Read partition 0 from the start
curl -s 'localhost:9000/topics/orders/consume?partition=0&offset=0'

# Record progress for the group "analytics", then look at it
curl -s -X POST localhost:9000/groups/analytics/commit \
  -H 'content-type: application/json' \
  -d '{"topic":"orders","partition":0,"offset":1}'
curl -s localhost:9000/groups/analytics/offsets
{"offsets":{"orders":{"0":1}}}

Authentication

Authentication is off unless the broker runs with --auth required. The first time it starts that way with no keys, it creates an admin key and writes its secret to <data-dir>/bootstrap.key (readable only by the broker's user). Use it to create real keys, then delete the file.

Send a key in either header:

request headers
Authorization: Bearer esk_AbCdEf…
X-Es-Key: esk_AbCdEf…
StatusMeaning
401No key, or a key the broker does not know.
403The key exists but has no grant for this topic, group or producer id.
429Over the key's byte-rate quota; Retry-After says how many seconds to wait.

Grants

A key carries rules of the form {"action", "topic_prefix"}:

read
Describe, consume, and use consumer groups on matching topics.
write
Everything read allows, plus produce.
admin
On *: create and delete topics, change configs, manage keys and producers, trigger retention and compaction.

A prefix covers the topic of that name and the topics below it, starting at a separator (., - or _): orders matches orders, orders.eu and orders-archive, not ordersecret. A prefix that ends in a separator, such as orders., matches the names that start with it. * matches every topic; an empty prefix is refused.

Keys

  • POST/admin/keysCreate a key. The secret is returned once.
  • GET/admin/keysList keys, without secrets.
  • DELETE/admin/keys/:key_idRevoke a key. It stops working at once, on open binary connections too.
POST /admin/keysjson
{
  "name": "orders-writer",
  "acls": [{ "action": "write", "topic_prefix": "orders." }],
  "produce_bytes_per_sec": 1048576,
  "consume_bytes_per_sec": 10485760
}

Quotas

produce_bytes_per_sec and consume_bytes_per_sec are per-key budgets holding one second of traffic. A produce is checked against its size before anything is written. A consume is charged after the read, when its size is known; the budget can go into debt, and further reads get 429 until it recovers.

TLS and browsers

--tls-cert and --tls-key (PEM files) switch both the HTTP and the binary listener to TLS. Give both or neither; one alone is an error.

No CORS headers are sent unless you list origins with --cors-origin https://ui.example (repeatable).

Topics

  • POST/topicsCreate a topic. Admin.
  • GET/topicsList the topics the key can read.
  • GET/topics/:namePartitions, offsets, sizes and config.
  • GET/topics/:name/configThe topic's resolved config.
  • PUT/topics/:name/configChange config; omitted fields keep their value. Admin.
  • DELETE/topics/:nameDelete the topic and its data. Admin.
POST /topicsjson
{
  "name": "orders",
  "partitions": 3,
  "config": {
    "retention_ms": 86400000,
    "retention_bytes": 1073741824,
    "cleanup_policy": "delete",
    "segment_bytes": 67108864,
    "tombstone_retention_ms": 86400000
  }
}

Names use letters, digits, ., - and _, up to 200 characters. A topic has between 1 and 10,000 partitions. Every config field is optional.

delete
Drop sealed segments older than retention_ms or beyond retention_bytes.
compact
Keep only the newest record per key in sealed segments. An empty value is a tombstone, kept for tombstone_retention_ms.
compact,delete
Both. The segment being written to is never compacted or deleted.

A new segment_bytes applies from the next segment roll; policy and retention changes from the next cleanup pass.

Partitions

A partition is an ordered log with its own files, offsets and write lock: the unit of both parallelism and ordering. A produced record goes to the partition it names; otherwise to the hash of its key (so a key always lands in the same partition, in order); otherwise round-robin.

You needPartitionsBecause
One total order1Order only holds within a partition.
Order per key4 to 8Each key stays on one partition.
Throughput, no ordering4 to 16Producers write to partitions in parallel.
N consumers in a groupat least NA partition goes to one member at a time.

Produce

  • POST/topics/:name/produceAppend records. Needs write.
POST /topics/orders/producejson
{
  "records": [
    { "key": "k1", "value": "hello" },
    { "key": "k1", "value": "world", "partition": 0 }
  ]
}
200 OKjson
{ "results": [
    { "partition": 0, "offset": 142 },
    { "partition": 0, "offset": 143 } ] }

The response comes once the records are on disk (with the default --flush-every-records 1). One request is one write per partition, so batching records into a request is the main lever on throughput.

Idempotent producers

Add a producer_id to the request and a sequence to each record, counting up per partition. A retry is then recognised and answered with the original offset and "duplicate": true instead of being written again.

POST /topics/orders/producejson
{
  "producer_id": "checkout-svc-1",
  "records": [
    { "key": "abc", "value": "…", "partition": 0, "sequence": 42 },
    { "key": "abc", "value": "…", "partition": 0, "sequence": 43 }
  ]
}
Next sequence
Written.
One of the last 128
A retry: not written again; the original offset comes back with "duplicate": true.
Skips ahead
400, sequence gap: send the missing records first.
Older than that
400, sequence too low.

The whole request is checked before anything is written. A producer id belongs to the key that first used it (another key gets 403); a key may own 10,000 and ids idle for 7 days are forgotten. Accepted sequences are on disk before the response is sent.

Consume

  • GET/topics/:name/consumeRead a partition from an offset.
  • GET/groups/:group/consumeRead from the group's committed offset.
partition
Required.
offset
Required on the topic endpoint.
topic
Required on the group endpoint.
max_records
Default 100, at most 10,000 (--max-fetch-records).
max_bytes
Default 1 MiB, at most 8 MiB (--max-fetch-bytes).
encoding
utf8 (default) or base64. Use base64 to read values that are not text, such as anything produced over the binary protocol.
200 OKjson
{
  "records": [
    { "partition": 0, "offset": 142, "timestamp_ms": 1763500000000,
      "key": "k1", "value": "hello" }
  ],
  "next_offset": 143,
  "high_watermark": 200
}

Pass next_offset as the next request's offset. Offsets can have gaps (compaction, retention), so never compute it yourself. Consuming through a group does not commit; commit explicitly.

Consumer groups

A group belongs to the key that first joined it or committed to it: only that key, or a global admin, can use it.

  • POST/groups/:group/commitRecord the next offset to read.
  • GET/groups/:group/offsetsCommitted offsets on topics the key can read.
  • POST/groups/:group/joinJoin and receive partitions.
  • POST/groups/:group/heartbeatStay a member; learn about rebalances.
  • POST/groups/:group/leaveLeave; the group rebalances.
  • GET/groups/:group/assignmentFetch the member's current partitions.
POST /groups/analytics/commitjson
{ "topic": "orders", "partition": 0, "offset": 143 }

The topic and partition must exist. Commits are written to disk once a second and at shutdown; after a crash a consumer may read up to a second of records again.

Coordination

Members join with the topics they consume (1 to 100, all readable by the key) and get a share of their partitions. A member that misses heartbeats for coord_member_timeout is removed and the group rebalances. A member id is bound to the key that created it.

POST /groups/analytics/joinjson
{ "topics": ["orders", "billing"], "member_id": null }

// response
{ "member_id": "m_a1b2c3…", "generation": 3,
  "assignment": [{ "topic": "orders", "partition": 0 },
                 { "topic": "billing", "partition": 1 }] }

A heartbeat sends member_id and generation and gets one of:

ok
Nothing changed.
rebalance_required
Fetch the assignment again.
unknown_member
The member was removed; join again.

Assignment is sticky: a rebalance keeps each member's partitions where it can and moves only what it must, so every member ends up with an even share.

Health and metrics

  • GET/healthzLiveness, with the topic count. No key needed.
  • GET/readyzReadiness. No key needed.
  • GET/metricsPrometheus text format.
  • GET/raftRole, term, leader and replication state of every Raft partition. Admin.

With auth on, /metrics needs a key and reports only the topics it can read and the groups it owns; a global admin sees everything. Series:

SeriesTypeLabels
es_partition_start_offset, …_end_offset, …_size_bytes, …_segment_countgaugetopic, partition
es_topic_records_produced_total, …_bytes_produced_total, …_records_consumed_total, …_bytes_consumed_totalcountertopic
es_group_committed_offset, es_group_laggaugegroup, topic, partition
es_retention_segments_deleted_total, es_retention_bytes_reclaimed_totalcountertopic, partition
es_compaction_runs_total, es_compaction_records_dropped_totalcountertopic, partition

Administration

  • POST/admin/run-retentionRun a retention pass now; returns when it is done.
  • POST/admin/run-compactionRun a compaction pass now; returns when it is done.
  • GET/admin/producersList idempotent producers and their sequences.
  • DELETE/admin/producers/:idForget a producer's state.
  • POST/admin/reset-offsetsMove every group on a topic back to its first offset.
  • GET/admin/tiered/:topic/:partitionList segments in cold storage.

All admin endpoints need admin on *.

Schema registry

  • POST/schemasRegister a schema version for a subject. Admin.
  • GET/schemasList schemas the key can read.
  • GET/schemas/:idFetch a schema by id.
  • GET/subjects/:subject/versions/latestThe newest version of a subject.

Name subjects after their topic (orders-value, orders-key): reading a subject needs read on it, so a grant on orders covers its schemas. A new version of a JSON schema must not add required fields or remove properties.

shell
curl -s -X POST localhost:9000/schemas -H 'content-type: application/json' -d '{
  "subject": "orders-value", "type": "json_schema",
  "schema": "{\"type\":\"object\",\"properties\":{\"id\":{\"type\":\"integer\"}},\"required\":[\"id\"]}"
}'
{"id":1}

# Adding a required field breaks old records: refused
400 {"error":"backward-incompatible: new schema adds required field 'phone'"}

Binary protocol

With --bind-binary 127.0.0.1:9001 the broker also listens for a length-prefixed binary protocol over a long-lived connection. Keys and values travel as raw bytes, with the same keys, grants, quotas and idempotency as HTTP. The Rust clients are es_broker::binary::BinaryClient and the pipelined PipelinedClient.

  • TLS is on whenever the broker has a certificate; connect with ClientOptions { tls: Some(…), .. }.
  • Concurrency. Up to 32 requests per connection run at once; produces to one topic keep their order.
  • Limits. Connections per address (--binary-max-connections-per-ip), an idle timeout (--binary-idle-timeout, 10 minutes), and a shared memory budget for request bytes. A client that stops reading responses is disconnected.
  • Authorization is checked against the topic at the start of a frame, before a compressed frame is inflated, and again after every key revocation.

Handshake

handshakebig-endian
client
magic        [u8; 4]   "ES01"
features     u32       requested bits; 0x1 = gzip
token_len    u32       up to 1024
token        [u8]

server
magic        [u8; 4]
status       u8        0 ok, 1 bad magic, 2 auth failed, 3 auth required
features     u32       the requested bits it accepts
msg_len      u32
msg          [u8]

With gzip negotiated, each frame's payload is gzip-compressed; the length, request id and opcode are not.

Frames

framebig-endian
length       u32       bytes that follow
request_id   u32       echoed in the response
opcode       u8        0x01 produce, 0x02 consume, 0x03 ping
                       0x81, 0x82, 0x83 their responses; 0xFF error (UTF-8 message)
payload      [u8]

Produce

opcode 0x01big-endian
topic_len    u16
topic        [u8]
producer_id  i32 length (-1 = none), then bytes
count        u32
per record:
  key_len    i32 (-1 = no key), then key
  value_len  u32, then value
  partition  i32 (-1 = route by key)
  sequence   i64 (i64::MIN = none)

The response (0x81) is count: u32, then partition: u32, offset: u64 and duplicate: u8 per record.

Consume

opcode 0x02big-endian
topic_len    u16
topic        [u8]
partition    u32
offset       u64
max_records  u32
max_bytes    u32

The response (0x82) is next_offset: u64, high_watermark: u64, count: u32, then per record partition: u32, offset: u64, timestamp_ms: i64, key_len: i32, key, value_len: u32, value.