Skip to content

消息队列(Message queuing)让用户无需等待智能体处理完当前消息,即可连续快速发送多条消息。每条消息都会被立即接受、加入活动线程的队列,并按顺序处理,让你对等待中的任务拥有完整的可见性与控制权。

import { PatternEmbed } from "/snippets/pattern-embed.jsx"

INFO

This feature requires the LangGraph Agent Server. Run your agent locally with langgraph dev or deploy it to LangSmith to use this pattern.

为什么要用消息队列?

在典型的聊天界面中,用户必须等待智能体完成回复后才能发送下一条消息。这在多种场景下会造成不便:

  • 批量提问:用户希望一次性提出五个相关问题,而不是等待每个答案
  • 追问链:在智能体仍在工作时提交澄清或补充上下文
  • 自动化测试序列:以编程方式发送一系列提示词来验证智能体行为
  • 数据录入工作流:逐个送入结构化输入以进行处理

消息队列通过立即接受所有提交并按顺序处理来解决这个问题。

这是智能体界面的基础能力(agent UX primitive),而非单纯的聊天装饰功能。SDK 将队列作为流控制器的一部分来维护,因此你的界面可以显示待处理的任务、取消过期的请求,并在当前运行继续时保持输入框可用。

工作原理

当你想让一次提交等待在当前正在运行的请求之后时,传入 multitaskStrategy: "enqueue"。在智能体处理期间,排队的提交会被添加到活动线程的队列中。当前运行完成后,下一条排队消息会自动派发。

使用对应框架附带的队列辅助函数读取队列状态:

属性类型描述
queue.entriesSubmissionQueueEntry[]所有待处理队列条目的数组
queue.sizenumber当前队列中的条目数量
queue.cancel(id)(id: string) => Promise<void>按 ID 取消特定的排队条目
queue.clear()() => Promise<void>取消所有排队条目

每个 SubmissionQueueEntry 对象包含:

字段类型描述
idstring此队列条目的唯一标识符
valuesobject提交的输入值(包括消息)
optionsobject随提交传入的任何其他选项
createdAtstring条目创建时的 ISO 时间戳

设置 useStream

useStream 连接到你的智能体,然后将其与对应框架的提交队列辅助函数搭配使用。调用 stream.submit() 在运行进行中发送消息;在应该等待当前活动请求之后的提交上传入 multitaskStrategy: "enqueue"。读取 queue.entriesqueue.size 以渲染待处理任务,并使用 queue.cancel()queue.clear() 在条目开始处理之前将其移除。

INFO

The code examples use useStream<typeof myAgent> for type-safe stream state. See Type inference for Python or JavaScript backends.

tsx
import { useStream, useSubmissionQueue } from "@langchain/react";

function Chat() {
  const stream = useStream<typeof myAgent>({
    apiUrl: "http://localhost:2024",
    assistantId: "simple_agent",
  });
  const queue = useSubmissionQueue(stream);

  const handleSubmit = (text: string) => {
    stream.submit({
      messages: [{ type: "human", content: text }],
    });
  };

  const pendingCount = queue.size;
  const entries = queue.entries;

  return (
      <MessageList messages={stream.messages} />
      {pendingCount > 0 && <QueueList entries={entries} queue={queue} />}
      <ChatInput onSubmit={handleSubmit} />
  );
}
vue
<script setup lang="ts">
import { useStream, useSubmissionQueue } from "@langchain/vue";
import { computed } from "vue";

const stream = useStream<typeof myAgent>({
  apiUrl: "http://localhost:2024",
  assistantId: "simple_agent",
});
const queue = useSubmissionQueue(stream);

function handleSubmit(text: string) {
  stream.submit({
    messages: [{ type: "human", content: text }],
  });
}

const pendingCount = computed(() => queue.size.value);
const entries = computed(() => queue.entries.value);
</script>

<template>
    <MessageList :messages="stream.messages" />
    <QueueList v-if="pendingCount > 0" :entries="entries" :queue="queue" />
    <ChatInput @submit="handleSubmit" />
</template>
svelte
<script lang="ts">
  import { useStream, useSubmissionQueue } from "@langchain/svelte";

  const stream = useStream<typeof myAgent>({
    apiUrl: "http://localhost:2024",
    assistantId: "simple_agent",
  });
  const queue = useSubmissionQueue(stream);

  function handleSubmit(text: string) {
    stream.submit({
      messages: [{ type: "human", content: text }],
    });
  }
</script>

  <MessageList messages={stream.messages} />
  {#if queue.size > 0}
    <QueueList entries={queue.entries} {queue} />
  {/if}
  <ChatInput on:submit={(e) => handleSubmit(e.detail)} />
ts
import { Component } from "@angular/core";
import { injectStream, injectSubmissionQueue } from "@langchain/angular";

@Component({
  selector: "app-chat",
  template: `
    <message-list [messages]="stream.messages()" />
    @if (queue.size() > 0) {
      <queue-list [entries]="queue.entries()" [queue]="queue" />
    }
    <chat-input (onSubmit)="handleSubmit($event)" />
  `,
})
export class ChatComponent {
  stream = injectStream<typeof myAgent>({
    apiUrl: "http://localhost:2024",
    assistantId: "simple_agent",
  });
  queue = injectSubmissionQueue(this.stream);

  handleSubmit(text: string) {
    this.stream.submit({
      messages: [{ type: "human", content: text }],
    });
  }
}

展示队列

构建一个 QueueList 组件,为每条待处理消息显示一个取消按钮。这可以让用户看到等待中的内容,并能移除不再需要的条目。

tsx
function QueueList({ entries, queue }) {
  return (
        Queued messages ({entries.length})
        <button onClick={() => queue.clear()}>Clear all</button>
        {entries.map((entry) => {
          const text = entry.values?.messages?.at(-1)?.content ?? "Pending...";
          return (
              {text}
                {new Date(entry.createdAt).toLocaleTimeString()}
              <button
                className="queue-cancel"
                onClick={() => queue.cancel(entry.id)}
              >
                Cancel
              </button>
          );
        })}
  );
}

TIP

将每条排队消息的前几个字符作为预览显示,这样用户无需阅读完整消息即可快速判断要取消哪些条目。

取消排队中的消息

你有两个层级的取消方式:

取消单个条目

按 ID 从队列中移除特定消息。智能体会跳过它并处理下一条。

ts
await queue.cancel(entryId);

清空整个队列

一次移除所有待处理消息。当用户更换上下文或想要重新开始时很有用。

ts
await queue.clear();

INFO

取消队列条目只会影响尚未开始处理的消息。如果智能体已经在处理某条消息,从队列中取消它不会有任何效果。请使用 stream.stop() 来中断当前的运行。

使用 onCreated 链接后续提交

onCreated 回调在新运行创建时触发,为你提供一个以编程方式提交后续消息的钩子。这在构建多步骤工作流时很有用,因为下一步的问题依赖于上一步的提交已被接受。

ts
stream.submit(
  { messages: [{ type: "human", content: "What is quantum computing?" }] },
  {
    onCreated(run) {
      console.log("Run created:", run.runId);
      // 链接一个后续提交
      stream.submit({
        messages: [{ type: "human", content: "Give me a simple analogy." }],
      });
    },
  }
);

这种模式会自然地填满队列。第一条消息立即开始处理,后续消息则排在它后面。

开启新线程

当用户想要开始一段全新对话时,更新你传给 stream 的响应式 threadId。传入 null 会清除当前的线程绑定;下一次提交会创建一个新线程。

tsx
function NewThreadButton() {
  const [threadId, setThreadId] = useState<string | null>(null);
  const stream = useStream<typeof myAgent>({ threadId, onThreadId: setThreadId });

  return (
    <button onClick={() => setThreadId(null)}>
      New conversation
    </button>
  );
}
vue
<script setup lang="ts">
const threadId = ref<string | null>(null);
const stream = useStream<typeof myAgent>({
  threadId,
  onThreadId: (id) => (threadId.value = id),
});
</script>

<template>
  <button @click="threadId = null">New conversation</button>
</template>
svelte
<script lang="ts">
  let threadId = $state<string | null>(null);
  const stream = useStream<typeof myAgent>({
    threadId: () => threadId,
    onThreadId: (id) => (threadId = id),
  });
</script>

<button onclick={() => (threadId = null)}>New conversation</button>
ts
threadId = signal<string | null>(null);
stream = injectStream<typeof myAgent>({
  threadId: this.threadId,
  onThreadId: (id) => this.threadId.set(id),
});

// 在模板中:
// <button (click)="threadId.set(null)">New conversation</button>

最佳实践

  • 限制队列大小:虽然客户端对队列大小没有硬性限制,但请注意非常大的队列会降低用户体验。当队列超过合理阈值(例如 10 个条目)时,考虑显示警告。
  • 显示队列位置:为每个排队条目编号,让用户了解处理顺序。
  • 保持输入焦点:提交后保持输入框处于聚焦状态,以便用户可以立即输入下一条消息。
  • 动画过渡:当条目开始处理时,将其从队列面板平滑地移动到消息列表中。
  • 优雅地处理错误:如果某条排队消息失败,在不阻塞后续队列条目的前提下呈现错误。
  • 对快速提交进行防抖:对于自动化或编程方式的提交,在消息之间加入少量延迟,以避免压垮服务器。