Consumers read records from a K2 stream through a subscription, which tracks which records have been processed. Each subscription on a stream reads independently, meaning multiple applications can read the same records.
To consume records:
- Create a subscription on the stream.
- Request a batch of records. K2 leases the batch to your consumer.
- Process the records, then acknowledge the batch. If processing takes longer than the lease, extend the lease.
- Repeat from step 2.
K2 delivers records at least once. Your consumer can receive the same record more than once.
All consumer requests use the stream endpoint,
https://<STREAM_ID>.k2.cloudflarestorage.com, and require an API token with
the K2 Consume permission in the account that owns the stream.
The examples on this page use the following shell variables:
export K2_ENDPOINT=https://<STREAM_ID>.k2.cloudflarestorage.com
export CLOUDFLARE_API_TOKEN=<YOUR_API_TOKEN>After you create a subscription, export its ID. After you request a batch, export the batch_id from the response:
export SUBSCRIPTION_ID=<SUBSCRIPTION_ID>
export BATCH_ID=<BATCH_ID>curl "$K2_ENDPOINT/subscriptions" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{
"name": "orders-ingestion",
"start_at": { "type": "earliest" }
}'{
"result": { "id": "e2f747f8bccf453eb2767496779cd135" },
"success": true,
"errors": [],
"messages": []
}| Field | Type | Required | Description |
|---|---|---|---|
name |
string | Yes | 1 to 128 letters, numbers, underscores, or hyphens. Must be unique within the stream. Not case-sensitive. |
start_at.type |
string | Yes | Where the subscription starts reading. earliest starts at the oldest retained record. latest starts after the newest record in the stream. |
Subscriptions are immutable once created.
If you create a subscription with the same name and settings as an existing
one, K2 returns the existing subscription ID. If the name matches but the
settings differ, K2 returns a 422 error.
| Action | Request |
|---|---|
| List subscriptions | GET /subscriptions |
| Find a subscription by name | GET /subscriptions?name=<SUBSCRIPTION_NAME> |
| Get a subscription | GET /subscriptions/<SUBSCRIPTION_ID> |
| Delete a subscription | DELETE /subscriptions/<SUBSCRIPTION_ID> |
List requests return an array of subscriptions, oldest first. A request with name returns an array with at most one subscription.
A get request returns a single subscription:
{
"result": {
"id": "e2f747f8bccf453eb2767496779cd135",
"name": "orders-processor",
"start_at": { "type": "earliest" },
"created_at": "2026-09-23T12:00:00.000Z",
"modified_at": "2026-09-23T12:00:00.000Z"
},
"success": true,
"errors": [],
"messages": []
}A stream can have up to 100 subscriptions. If you create a subscription on a stream that already
has 100, K2 returns a 422 error with code 10219. Delete an existing subscription to create a new one.
Send a POST request to /subscriptions/<SUBSCRIPTION_ID>/consume:
curl "$K2_ENDPOINT/subscriptions/$SUBSCRIPTION_ID/consume" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{
"worker_id": "worker-1",
"max_records": 100
}'| Field | Type | Required | Description |
|---|---|---|---|
worker_id |
string | Yes | An identifier for your consumer, 1 to 256 characters. Use a different value for each concurrent consumer. |
max_records |
integer | Yes | The maximum number of records to return, from 1 to 10000. |
{
"result": {
"batch_id": "4f1c2a9e8b7d4c6f9a0e1d2c3b4a5968",
"leased_until_ms": 1790165100000,
"records": [
{
"timestamp_ms": 1790164800000,
"content": "eyJvcmRlcl9pZCI6MTAwMSwic3RhdHVzIjoiY3JlYXRlZCJ9",
"headers": { "event-type": "order.created" }
}
]
},
"success": true,
"errors": [],
"messages": []
}| Field | Description |
|---|---|
batch_id |
The ID of the leased batch. Use it to acknowledge the batch. |
leased_until_ms |
When the lease expires, in milliseconds since the Unix epoch. |
records[].timestamp_ms |
When K2 received the record, in milliseconds since the Unix epoch. |
records[].content |
The record payload, as standard base64. |
records[].headers |
The record headers. Omitted if the record was produced without any headers. |
max_records is a ceiling for the number of records that will be returned, but a batch may contain
fewer records. This does not necessarily imply that there are no more records available.
If there are no new records, the response contains an empty batch:
{
"result": { "batch_id": null, "leased_until_ms": null, "records": [] },
"success": true,
"errors": [],
"messages": []
}Wait before you send another request. Use a backoff interval to avoid polling continuously.
When K2 returns a batch, it leases the batch to the worker_id that requested it. The lease lasts five minutes.
- Each
worker_idcan hold one lease at a time. - K2 never leases the same records to two workers at the same time.
- A subscription can have up to 128 active leases at a time. If every lease is in use, K2 returns a
429error with code10216. - K2 does not guarantee the order in which records are processed across workers that share a subscription.
If a worker requests a batch while it already holds a lease, K2 returns the
same batch and records again, and refreshes the lease. K2 ignores max_records
when it returns the same batch. Use this to recover a
batch if your consumer loses the response, for example after a network error.
If a lease expires before the batch is acknowledged, K2 delivers the records again to the next worker that requests a batch.
After you process a batch, acknowledge it. This marks the records as processed and releases the lease, so the worker can request the next batch.
curl "$K2_ENDPOINT/subscriptions/$SUBSCRIPTION_ID/batches/$BATCH_ID/ack" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{ "worker_id": "worker-1" }'{ "result": {}, "success": true, "errors": [], "messages": [] }Acknowledge the batch before leased_until_ms. If the lease has expired and K2
has already delivered the records in a new batch, the acknowledgement has no
effect.
Acknowledging the same batch more than once is safe. K2 returns a success response for batches that are unknown or already acknowledged.
If your consumer needs more than five minutes to process a batch, extend the lease before it expires:
curl "$K2_ENDPOINT/subscriptions/$SUBSCRIPTION_ID/batches/$BATCH_ID/extend" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{ "worker_id": "worker-1" }'{
"result": { "leased_until_ms": 1790165400000 },
"success": true,
"errors": [],
"messages": []
}A successful extension sets the lease to expire five minutes after the request. It never shortens a lease.
Unlike ack and nack, the worker_id must match the worker that holds the lease. If the lease has
expired, has been released, or the batch has been acknowledged or delivered to another worker, K2
returns a 409 error with code 10218. Retrying the same request cannot succeed. Stop processing the
batch and request a new batch. K2 delivers the records again with a new batch_id.
If your consumer cannot process a batch, send a negative acknowledgement (nack) to release it immediately instead of waiting for the lease to expire:
curl "$K2_ENDPOINT/subscriptions/$SUBSCRIPTION_ID/batches/$BATCH_ID/nack" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{ "worker_id": "worker-1" }'{ "result": {}, "success": true, "errors": [], "messages": [] }K2 delivers the released records again in the next batch requested by any
worker, with a new batch_id. The subscription does not move past the records
until that new batch is acknowledged.
Subscription and consume errors use the Cloudflare API format, with an errors array:
{
"result": null,
"success": false,
"errors": [
{
"code": 10216,
"message": "All parallel read slots for this subscription are in use, please retry"
}
],
"messages": []
}This is different from the format of produce errors.
Retry requests that fail with codes 10211, 10214, 10216, or 10217 after a backoff interval. Do not retry other errors without changing the request.
Subscription and consume error codes
| Code | HTTP status | Description |
|---|---|---|
10200 |
404 |
The stream does not exist. |
10201 |
422 |
A subscription with this name exists with different settings. |
10204 |
400 |
The request is invalid. For example, the JSON is malformed or a field is out of range. |
10205 |
415 |
The Content-Type header is not application/json. |
10208 |
401 |
The request has no Authorization header. |
10209 |
401 |
The Authorization header is malformed or the API token is invalid. |
10210 |
403 |
The API token does not have permission to consume from this stream. |
10211 |
503 |
K2 is temporarily unavailable. Retry the request. |
10213 |
500 |
An internal error occurred. |
10214 |
503 |
K2 could not read records from storage. Retry the request. |
10215 |
404 |
The subscription does not exist. |
10216 |
429 |
The maximum number of concurrent consumers has been exceeded for this subscription. Retry after a lease is acknowledged, released, or expires. |
10217 |
409 |
K2 is still reading a batch for this worker_id from an earlier request. Retry the request to receive the batch. |
10218 |
409 |
The lease is no longer held by this worker. Do not retry. Request a new batch instead. |
10219 |
422 |
The stream already has the maximum of 100 subscriptions. |
- Review consumer limits in Limits.