Spec-Zone.ru › Deno 2

Использование очередей

В среде выполнения Deno реализован API очередей, который поддерживает перегрузку более объемных задач для асинхронной обработки, гарантируя доставку сообщений в очереди не менее одного раза. Очереди могут использоваться для перегрузки задач в веб-приложении или для планирования задач на будущее.

Основные API, которые вы будете использовать с очередями, находятся в пространстве имен Deno.Kv как enqueue и listenQueue.

Добавление сообщения в очередь

Для добавления сообщения в очередь для обработки используйте метод enqueue на экземпляре Deno.Kv. В примере ниже показано, как можно добавить уведомление для доставки.

queue_example.ts
// Describe the shape of your message object (optional)
interface Notification {
  forUser: string;
  body: string;
}

// Get a reference to a KV instance
const kv = await Deno.openKv();

// Create a notification object
const message: Notification = {
  forUser: "alovelace",
  body: "You've got mail!",
};

// Enqueue the message for immediate delivery
await kv.enqueue(message);

Вы можете добавить сообщение для последующей доставки, указав опцию delay в миллисекундах.

// Enqueue the message for delivery in 3 days
const delay = 1000 * 60 * 60 * 24 * 3;
await kv.enqueue(message, { delay });

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

// Configure a key where a failed message would be sent
const backupKey = ["failed_notifications", "alovelace", Date.now()];
await kv.enqueue(message, { keysIfUndelivered: [backupKey] });

// ... disaster strikes ...

// Get the unsent message
const r = await kv.get<Notification>(backupKey);
// This is the message that didn't get sent:
console.log("Found failed notification for:", r.value?.forUser);

Прослушивание сообщений

Вы можете настроить JavaScript-функцию, которая будет обрабатывать элементы, добавленные в вашу очередь, с помощью метода listenQueue на экземпляре Deno.Kv.

listen_example.ts
// Define the shape of the object we expect as a message in the queue
interface Notification {
  forUser: string;
  body: string;
}

// Create a type guard to check the type of the incoming message
function isNotification(o: unknown): o is Notification {
  return (
    ((o as Notification)?.forUser !== undefined &&
      typeof (o as Notification).forUser === "string") &&
    ((o as Notification)?.body !== undefined &&
      typeof (o as Notification).body === "string")
  );
}

// Get a reference to a KV database
const kv = await Deno.openKv();

// Register a handler function to listen for values - this example shows
// how you might send a notification
kv.listenQueue((msg: unknown) => {
  // Use type guard - then TypeScript compiler knows msg is a Notification
  if (isNotification(msg)) {
    console.log("Sending notification to user:", msg.forUser);
    // ... do something to actually send the notification!
  } else {
    // If the message is of an unknown type, it might be an error
    console.error("Unknown message received:", msg);
  }
});

API очередей с атомарными транзакциями KV

Вы можете объединить API очередей с атомарными транзакциями KV, чтобы атомарно добавлять сообщения в очередь и изменять ключи в одной транзакции.

kv_transaction_example.ts
const kv = await Deno.openKv();

kv.listenQueue(async (msg: unknown) => {
  const nonce = await kv.get(["nonces", msg.nonce]);
  if (nonce.value === null) {
    // This messaged was already processed
    return;
  }

  const change = msg.change;
  const bob = await kv.get(["balance", "bob"]);
  const liz = await kv.get(["balance", "liz"]);

  const success = await kv.atomic()
    // Ensure this message was not yet processed
    .check({ key: nonce.key, versionstamp: nonce.versionstamp })
    .delete(nonce.key)
    .sum(["processed_count"], 1n)
    .check(bob, liz) // balances did not change
    .set(["balance", "bob"], bob.value - change)
    .set(["balance", "liz"], liz.value + change)
    .commit();
});

// Modify keys and enqueue messages in the same KV transaction!
const nonce = crypto.randomUUID();
await kv
  .atomic()
  .check({ key: ["nonces", nonce], versionstamp: null })
  .enqueue({ nonce: nonce, change: 10 })
  .set(["nonces", nonce], true)
  .sum(["enqueued_count"], 1n)
  .commit();

Поведение очереди

Гарантии доставки сообщений

Система гарантирует доставку сообщений не менее одного раза. Это означает, что для большинства сообщений, помещенных в очередь, обработчик listenQueue будет вызываться один раз для каждого сообщения. В некоторых ситуациях сбоя обработчик может вызываться несколько раз для одного и того же сообщения, чтобы гарантировать доставку. Важно спроектировать ваши приложения так, чтобы дублирующиеся сообщения обрабатывались корректно.

Вы можете использовать очереди в сочетании с атомарными транзакциями KV, чтобы гарантировать, что обновления ключей в обработчике очереди KV выполняются ровно один раз на сообщение. См. API очередей с атомарными транзакциями KV.

Автоматические повторы

Обработчик listenQueue вызывается для обработки сообщений из очереди, когда они готовы к доставке. Если ваш обработчик генерирует исключение, среда выполнения автоматически повторит попытку вызова обработчика до тех пор, пока он не увенчается успехом или не будет достигнуто максимальное количество попыток повтора. Сообщение считается успешно обработанным, когда вызов обработчика listenQueue завершается успешно. Сообщение будет проигнорировано, если обработчик постоянно терпит неудачу при повторах.

Порядок доставки сообщений

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

Очереди в Deno Deploy

Deno Deploy предлагает глобальную, серверless, распределенную реализацию API очереди, предназначенную для высокой доступности и пропускной способности. Вы можете использовать ее для создания приложений, масштабируемых для обработки больших объемов работы.

Динамическое создание изолированных процессов

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

Предел размера очереди

Максимальное количество непереданных сообщений в очереди ограничено 100 000. Метод enqueue вернет ошибку, если очередь заполнена.

Детали ценообразования и ограничения

  • enqueue рассматривается как любая другая операция записи в Deno.Kv. Сообщения в очереди занимают место в хранилище KV и потребляют ресурсы записи.
  • Сообщения, переданные через listenQueue, потребляют запросы и ресурсы записи KV.
  • См. Детали ценообразования для получения дополнительной информации.

Сценарии использования

Очереди могут быть полезны во многих сценариях, но есть несколько распространённых сценариев при создании веб-приложений.

Перегрузка асинхронных процессов

Иногда задача, инициированная клиентом (например, отправка уведомления или API-запрос), может занять достаточно много времени, поэтому вы не хотите заставлять клиентов ждать завершения этой задачи перед возвратом ответа. В других случаях клиентам вообще не нужен ответ, например, когда клиент отправляет вашему приложению Webhook-запрос, поэтому нет необходимости ждать завершения задачи перед возвратом ответа.

В этих случаях вы можете перегрузить работу в очередь, чтобы сохранить отзывчивость вашего веб-приложения и отправлять немедленные ответы клиентам. Чтобы увидеть пример использования этого сценария, ознакомьтесь с нашим примером обработки webhook.

Планирование задач на будущее

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

Чтобы посмотреть пример планирования уведомления на будущее, ознакомьтесь с нашим примером уведомления.

© 2018–2024 the Deno authors
Licensed under the MIT License.
https://docs.deno.com/deploy/kv/manual/queue_overview

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API