Skip to content

Consume records

Last updated View as MarkdownAgent setup

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:

  1. Create a subscription on the stream.
  2. Request a batch of records. K2 leases the batch to your consumer.
  3. Process the records, then acknowledge the batch. If processing takes longer than the lease, extend the lease.
  4. Repeat from step 2.

K2 delivers records at least once. Your consumer can receive the same record more than once.

Before you begin

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>

Create a subscription

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.

Manage subscriptions

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.

Request a batch

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.

No 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.

Leases

When K2 returns a batch, it leases the batch to the worker_id that requested it. The lease lasts five minutes.

  • Each worker_id can 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 429 error with code 10216.
  • 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.

Acknowledge 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.

Extend a lease

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.

Release a batch

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.

Handle errors

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.

Next steps

  • Review consumer limits in Limits.

Was this helpful?