В традиционных архитектурах с удалёнными вызовами процедур (RPC) существует фундаментальная проблема: производители и потребители должны быть синхронизированы по масштабу и времени. Если производители отправляют слишком много данных, которые потребители не могут обработать, или если потребители или последующие сервисы становятся недоступны, события теряются. Проблема усугубляется, когда несколько потребителей должны независимо обрабатывать данные. Например, бэкенд электронной коммерции может генерировать события при завершении транзакций, которые должны читать система аналитики и сервис обнаружения мошенничества.

Решить это можно путём разделения связи между производителями и потребителями — вставки посередине сервиса, который абсорбирует записи и позволяет независимым читателям потреблять данные в своём темпе.

Сегодня Cloudflare запускает K2 в публичной бете для решения этой проблемы. K2 — это примитив надёжной потоковой передачи событий на платформе разработчиков. События отправляются в поток K2, который хранит их в виде упорядоченного лога. Потребители могут читать их различными способами: разделяя чтение между группой потребителей или доставляя все сообщения всем потребителям. Система полностью бессерверная, масштабируется на огромные объёмы данных и поддерживает длительное хранение, поэтому даже длительные периоды недоступности потребителей не приводят к потере данных.

Под капотом K2 реализует разделённый надёжный лог на основе R2 object storage, что позволяет ему масштабироваться на огромные объёмы хранилища.

Если вы готовы начать, можете создать свой первый поток за несколько секунд, следуя руководству.

Потоки на edge

K2 был создан потому, что нам нужен был надёжный буфер на edge, изначально для использования в качестве уровня приёма для Basin Pipelines. Pipelines работает на stream processing engine, который использует модель на основе pull, что означает, что какая-то другая система должна хранить события перед их чтением, преобразованием и записью в R2. А поскольку мы берём на себя обязательство никогда не терять события после их принятия в Pipelines Stream, это хранилище должно быть надёжным — значит, оно не должно терять данные — в течение потенциально длительных периодов времени.

Здесь большинство компаний развёртывают Apache Kafka. Однако Pipelines работает на edge Cloudflare, который охватывает огромное количество серверов в более чем 335 городах. Наша уникальная архитектура означает, что мы часто не можем запускать традиционное распределённое ПО, такое как Kafka, и должны переосмыслить, как эти системы строятся и работают.

Для stateful сервисов глобальная инфраструктура Cloudflare представляет некоторые вызовы: мы получаем относительно небольшие слои машин, эти машины относительно эфемерны, и сеть часто работает через общественный Интернет. Но наша инфраструктура также имеет несколько суперсил: она находится близко к пользователям, где бы они ни находились, и обладает невероятной способностью к горизонтальному масштабированию.

При разработке системы надёжной буферизации, которая стала K2, мы решили полагаться на мощный примитив состояния, который уже есть у нас: R2. Системы object storage, такие как R2, объединяют чрезвычайно надёжное хранилище (11 9s!) с сильно консистентными API. Перенос репликации и консенсуса на уровень хранилища позволяет нам сделать уровень приложения (в данном случае K2) радикально проще, дешевле и производительнее. Вторичное преимущество — разделение compute и storage, что означает, что каждый может масштабироваться независимо. Это позволяет нам хранить огромные объёмы исторических данных с низкой стоимостью.

Как построить лог на основе object storage? Непосредственная проблема в том, что R2 — как и другие object stores — не поддерживает append, стандартную операцию над логом. Вместо этого мы должны писать полные файлы или сегменты, которые достаточно большие, чтобы оправдать стоимость записи и чтения каждого из них. Мы делаем это, сначала накапливая записи в памяти на edge сервисе. После ожидания короткого периода поступления данных мы записываем все события как файл сегмента. Мы достигаем упорядочивания и строго возрастающих смещений, используя атомарные операции R2 без необходимости в отдельном сервисе координации.

Хотя построение на R2 имеет много преимуществ, есть один недостаток: более высокие latency при производстве. Запись в object storage медленнее, чем на локальный диск, и мы должны ждать, пока локальный batch накопится перед началом записи. В нашем первоначальном выпуске K2 это добавляет примерно 1 секунду latency при производстве на 99-м процентиле времён ответа.

Мы поделимся более подробной информацией о дизайне K2 в предстоящем техническом глубоком погружении.

Потоки, Queues или Pipelines?

Cloudflare имеет несколько существующих асинхронных примитивов доставки, включая Queues и Basin Pipelines. Когда следует использовать K2 вместо этих существующих продуктов?

Есть некоторые поверхностные сходства между Queues и K2 Streams: оба получают события, надёжно хранят их и доставляют потребителям. Queues разработаны для отслеживания отдельных элементов дорогостоящей или трудоёмкой работы, которая должна быть асинхронно завершена. Например, приложение обработки изображений может поставить в очередь запрос пользователя для обработки фактическим сервисом обработки изображений. Они поддерживают сложную логику на уровне отдельного рабочего элемента, такую как повторные попытки, задержки и dead-letter queues для неудачных попыток.

K2, напротив, разработан для массового движения данных, долгосрочного хранения и fan-out потребления. Сообщения производятся и потребляются партиями — обеспечивая эффективную обработку за счёт retry на уровне сообщений. Эта партионизация также приводит к более высокому latency при производстве, чем для queues.

Basin Pipelines — это бессерверный сервис приёма. Вы можете отправлять JSON события в Pipeline, которые можно преобразовать и записать в R2 или Basin Catalog. Мы рекомендуем Pipelines, когда конечный результат — запись событий в object storage или Iceberg таблицы, и K2 для пользовательской обработки или записи в другие места назначения.

Начало работы

Использование K2 начинается с создания потока. Вы можете иметь много потоков в своём аккаунте для разных вариантов использования или типов событий. Потоки можно создавать через cf, Wrangler, dashboard или API.

Рассмотрим пример сбора и обработки аналитики продукта. Сначала создадим поток с помощью cf:

$ cf k2 streams create --name app_events --http-enabled

{
  "id": "d78b09ee1f50430e9ec92a8af92b0231",
  "name": "app_events",
  "retention_seconds": 604800,
  "endpoint": "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com",
  "http": {
    "enabled": true,
    "authentication": false
  },
  "worker_binding": {
    "enabled": true
  },
  "created_at": "2026-09-28T15:14:39.053Z",
  "modified_at": "2026-09-28T15:14:39.053Z"
}

После создания потока можно начать его использование через HTTP API или Worker binding. Например, из Worker:

const result = await env.EVENTS.send([
  {
    content: new TextEncoder().encode(
      JSON.stringify({
        event: "page_view",
        path: new URL(request.url).pathname,
        timestamp: Date.now(),
      }),
    ),
    headers: { "content-type": "application/json" },
  },
]);

if (!result.success) {
  console.error(`Produce failed: ${result.error.message}`);
  return new Response("Failed to record event", {
    status: result.error.retryable ? 503 : 500,
  });
}

K2 представляет данные как байты, поэтому можете использовать любой формат или кодировку, который имеет смысл для вашего приложения.

После размещения событий в потоке можно создать подписку. Подписки распределяют работу между потребителями, обеспечивая параллелизм чтения — масштабирование к нескольким читателям для обработки большей нагрузки, чем может справиться один сервер.

Подписку можно создать через HTTP API.

$ curl -X POST "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com/subscriptions" \
    -H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
    -H "Content-Type: application/json" \
    --data '{
      "name": "analytics_processor",
      "start_at": { "type": "earliest" }
    }'
{
  "result": { "id": "ee13f761783d3823a447a47b572ebf76" },
  "success": true,
  "errors": [],
  "messages": []
}

После создания подписки можно начать опрашивать её из каждого из потребителей:

$ curl -X POST \"https://4d8f5394e3e733debdeca9c65c5b7439.k2.cloudflarestorage.com/subscriptions/ee13f761783d3823a447a47b572ebf76/consume" \
    -H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
    -H "Content-Type: application/json" \
    --data '{
      "worker_id": "analytics-1",
      "max_records": 100
    }' 
{
  "result": {
    "batch_id": "b7e4c9210a3f468d95c2e1068fdb734a",
    "leased_until_ms": 1790633929437,
    "records": [
      {
        "timestamp_ms": 1790633629168,
        "content": "eyJldmVudCI6InBhZ2VfdmlldyJ9",
        "headers": {
          "content-type": "application/json"
        }
      },
      ...
    ]
  },
  "success": true,
  "errors": [],
  "messages": []
}

Когда клиент вызывает consume, он получает lease на конкретный batch событий на 5 минут. Клиент может сделать одно из трёх:

  • ack batch, что означает его обработку и гарантирует, что он не будет переотправлен
  • nack (отрицательный ack), что означает, что обработка не удалась и нужна переотправка
  • extend его lease, если требуется больше времени для завершения обработки
$ curl -X POST "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com/subscriptions/ee13f761783d3823a447a47b572ebf76/batches/b7e4c9210a3f468d95c2e1068fdb734a/ack" \
  -H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
  -H "Content-Type: application/json" \
  --data '{ "worker_id": "analytics-1" }'

Это один из способов потребления из K2: разделение работы между несколькими потребителями таким образом, что каждый получает часть данных. Другой способ чтения — отдельная подписка для каждого потребителя — шаблон pub/sub — в этом случае каждый потребитель видит все сообщения. Или можно комбинировать оба подхода, имея несколько независимых пулов потребителей.

Подробные детали API см. в документации K2.

Цены и доступность

K2 доступен сегодня в публичной бете для аккаунтов с подписками Workers Paid со следующими лимитами:

  • Максимум 10 GB используемого хранилища
  • 30 MB/s produce per stream

Если требуются более высокие лимиты, обратитесь к команде в Discord или заполните форму запроса на увеличение лимитов.

Использование K2 не будет взиматься во время бета-периода. После начала выставления счетов предполагается следующее ценообразование:

Цена

Произведённые данные

$0.04 / GB

Потреблённые данные

$0.04 / GB

Сохранённые данные

$0.02 / GB / месяц

Дальнейшие планы

На предстоящие месяцы у нас есть амбициозная дорожная карта для K2:

  • Более высокий параллелизм записи, до потоков multi-GB/s
  • Ключи сообщений и гарантии упорядочивания на основе ключа
  • Push-based worker потребители
  • Express tier с более низким produce и сквозным latency
  • Drop-in поддержка клиентов Apache Kafka

Нам не терпится увидеть, что вы построите на K2! Поделитесь отзывами в Cloudflare Discord.