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
# 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:
Authorization: Bearer esk_AbCdEf… X-Es-Key: esk_AbCdEf…
| Status | Meaning |
|---|---|
401 | No key, or a key the broker does not know. |
403 | The key exists but has no grant for this topic, group or producer id. |
429 | Over 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
readallows, 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.
{
"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.
{
"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_msor beyondretention_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 need | Partitions | Because |
|---|---|---|
| One total order | 1 | Order only holds within a partition. |
| Order per key | 4 to 8 | Each key stays on one partition. |
| Throughput, no ordering | 4 to 16 | Producers write to partitions in parallel. |
| N consumers in a group | at least N | A partition goes to one member at a time. |
Produce
- POST
/topics/:name/produceAppend records. Needswrite.
{
"records": [
{ "key": "k1", "value": "hello" },
{ "key": "k1", "value": "world", "partition": 0 }
]
}
{ "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.
{
"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). encodingutf8(default) orbase64. Use base64 to read values that are not text, such as anything produced over the binary protocol.
{
"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.
{ "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.
{ "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:
| Series | Type | Labels |
|---|---|---|
es_partition_start_offset, …_end_offset, …_size_bytes, …_segment_count | gauge | topic, partition |
es_topic_records_produced_total, …_bytes_produced_total, …_records_consumed_total, …_bytes_consumed_total | counter | topic |
es_group_committed_offset, es_group_lag | gauge | group, topic, partition |
es_retention_segments_deleted_total, es_retention_bytes_reclaimed_total | counter | topic, partition |
es_compaction_runs_total, es_compaction_records_dropped_total | counter | topic, 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.
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
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
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
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
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.