ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Hindsight 异步操作(Operations)全指南:任务队列、生命周期与 Worker 调优

Hindsight 异步操作(Operations)全指南:任务队列、生命周期与 Worker 调优 Hindsight 异步操作Operations全指南任务队列、生命周期与 Worker 调优【免费下载链接】hindsightHindsight: Agent Memory That Learns项目地址: https://gitcode.com/GitHub_Trending/hindsight2/hindsightHindsight 将多项维护与摄入任务retain、consolidation、graph maintenance 等从 HTTP 请求处理中剥离改为异步执行所有任务共享同一张async_operations队列表和同一个 worker 池并通过统一的 REST 端点list / status / cancel / retry进行管理。本文基于官方 Operations 文档结合 Hindsight 开源仓库的源码实现完整讲解每种操作类型、触发时机、状态机生命周期以及如何用 API、CLI 和官方 SDK 检查与管理这些操作最后深入 worker 并发槽位slot调优原理。异步操作的基本原理当一次 API 调用需要后台工作时请求处理器会先向async_operations表写入一行statuspending的记录并立即返回不会阻塞调用方。随后worker 会轮询该表、认领claim待处理行、执行对应的 handler最后将行标记为completed或failed。从源码看认领机制的核心是FOR UPDATE SKIP LOCKED保证多个 worker 并发认领时互不冲突。在 worker/poller.py 中WorkerPoller的claim_batch方法先通过轻量EXISTS探测哪些 schema 有待处理工作再对活跃 schema 执行昂贵的FOR UPDATE SKIP LOCKED认领查询认领发生在单条连接上的单次事务里行锁只保持到该 schema 的认领提交为止不会拖住整个轮询周期每个 schema 的认领会记录next_bank_cursor下一轮从该 bank 之后继续避免单个繁忙 bank 独占轮询顺序对应 issue #3861 的公平性修复。默认情况下所有操作都在 API 进程内运行不需要外部队列也不需要额外部署进程。当吞吐量要求提高时同一套代码路径可以横向扩展出独立的 worker 进程见 Services - Worker Service。独立 worker 的入口在 worker/main.py启动命令为hindsight-workerworker 进程通过WorkerPoller每HINDSIGHT_API_WORKER_POLL_INTERVAL_MS毫秒默认 500ms轮询一次并使用MemoryEngine.execute_task作为执行器。它以WorkerTaskBackend作为任务后端——submit_task是空操作因为行已经存在于async_operations表中由 retain 触发的子任务如 consolidation会在下一轮轮询时被拾取而不是在父任务内联执行从而避免阻塞父任务。生命周期每行操作记录会经历以下状态Status含义pending已入队。可能尚未被任何 worker 拾取也可能被扩展通过next_retry_at推迟到未来例如背压场景。processing已被某 worker 认领正在执行 handler。completedhandler 成功返回。failedhandler 抛出了异常。error_message记录原因可通过POST /…/retry重新入队。cancelled已通过DELETE /…/operations/{id}取消。对pending和processing状态均有效。worker 会在最终落定为failed之前最多按HINDSIGHT_API_WORKER_MAX_RETRIES默认 3见 config.py重试失败操作。确定性失败例如无效的 embedding 维度、完整性约束冲突会跳过重试——重跑不可能成功。操作历史的保留与清理默认情况下completed、failed和cancelled的历史记录无限期保留。设置HINDSIGHT_API_OPERATION_RETENTION_DAYS为正整数可以限制保留时长后台维护循环会按照自己的调度而非任务处理的副作用以有界批次清理过期终态行。相关配置在 config.py 中HINDSIGHT_API_OPERATION_RETENTION_DAYS默认 0禁用清理保留全部历史HINDSIGHT_API_OPERATION_CLEANUP_BATCH_SIZE每轮每 schema 清理的行数默认 1000HINDSIGHT_API_OPERATION_CLEANUP_INTERVAL_SECONDS清理循环间隔默认 900 秒见 config.py。需要说明的是维护循环仅在 PostgreSQL 上运行——Oracle 后端没有该循环因此 Oracle 上的操作历史无界增长。整行数据共享同一个 TTL因此在保留期内失败/已取消的操作仍可重试已完成的操作可通过include_payloadtrue检查其载荷。pending和processing的操作永远不会被保留清理移除。操作类型总览每个操作在数据库中都有operation_type在载荷中有task_type两者通常一致。worker 启动时会把每种类型的槽位保留情况打印出来见 worker/main.py。以下逐一介绍六种核心操作类型。retain由POST /v1/default/banks/{bank_id}/memories携带asynctrue提交或由多条目retain_batch调用提交。handler 执行的流水线与同步 retain 完全一致事实抽取LLM→ embedding 生成 → 实体解析 → 链接创建时间/语义链接。在需要摄入数千条数据、又不想让 HTTP 调用挂起数分钟的场景下应使用异步 retain。响应中的operation_id可用于轮询完成状态。其提交实现位于 memory_engine.py 的 submit_async_retain支持调用方自行传入operation_id作为幂等键重复提交相同 id 会返回原始操作、不产生新工作从而避免丢失确认后的客户端重试造成重复入队。父操作retain_batch对于大批量提交Hindsight 会自动把输入拆分为子批次并创建一个retain_batch父操作来跟踪所有子操作父状态反映聚合结果——pending直到至少一个子任务开始运行、processing子任务执行中、completed所有子任务完成、failed任一子任务失败每个子任务本身就是一个retain操作并通过result_metadata中的parent_operation_id关联到父操作因此可以下钻查看每个批次的错误信息。父/子状态汇总逻辑有两处实现一处位于 worker/poller.py 的_maybe_update_parent_operation处理未经引擎失败路径、直接以未捕获异常失败的场景另一处在MemoryEngine内处理任务执行事务内的正常路径。关键细节父行会被FOR UPDATE锁定防止并发子任务更新互相覆盖子任务出现cancelled也视为完成否则父任务会永远卡在processingissue #4131真实失败优先于取消父任务失败时会继承子任务中最常见的错误信息_summarise_child_error_messages避免父任务只剩一句泛化的 One or more sub-batches failed。列表操作默认同时显示父操作与子操作传exclude_parentstrue可隐藏聚合行、只显示独立的retain任务。file_convert_retain由文件上传端点提交。handler 先执行 MIME 特定的格式转换PDF → 文本、DOCX → 文本等再把提取的文本送入 retain 流水线。此类失败默认不可重试——损坏的 PDF 或缺失的 OCR 重跑也不会变好所以操作会直接进入failed。使用哪个解析器markitdown、iris或llama_parse由部署级配置HINDSIGHT_API_FILE_PARSER决定默认markitdown见 config.py客户端也可以在单次请求中覆盖。相关配置还包括HINDSIGHT_API_FILE_PARSER_ALLOWLIST允许客户端请求的解析器白名单、HINDSIGHT_API_FILE_PARSER_MARKITDOWN_OCR_ENABLED及各解析器所需的 API Key 等见 config.py。consolidation从新的世界/经验记忆生成观察observations。触发时机每次 retain 新增了世界/经验事实之后由每 bank 的enable_auto_consolidation与enable_observations控制删除使既有观察失效之后源记忆消失 → 派生观察过期 → 用幸存的共同源记忆重跑手动触发POST /v1/default/banks/{bank_id}/consolidate可传observation_scopes只合并匹配特定标签组合的记忆。按 bank 去重当一个 bank 的consolidation任务处于pending时重复提交会返回已有的operation_id而不是叠加新任务任务开始处理后下一次提交才成为新的 pending 槽位。refresh_mental_modelmental model 有source_query定义其汇总哪些记忆。handler 会重新执行该查询、重新汇总结果并在原地更新模型内容。触发方式手动调用POST /v1/default/banks/{bank_id}/mental-models/{id}/refresh或由配置了自动刷新计划的 mental model 按计划自动触发。graph_maintenance在删除之后调和过期的派生状态。每次调用排空两个队列均由删除本身填充因此一次运行只处理该删除触及的内容链接补齐Link top-up排空出向时间/语义链接失去邻居的单元。若某单元未达上限20 条时间链接、50 条语义链接Hindsight 会重跑 retain 使用的相同探测并补插缺失链接。没有这一步retain 流水线的 top-K 截断会让幸存的单元在每次删除后永久低于上限降低图扩展召回率。实体清理Entity prune排空删除可能遗留的实体。已无任何unit_entities引用的实体被删除外键ON DELETE CASCADE会连带删除其entity_cooccurrences行对于幸存实体清理那些两端仍存在但已无任何memory_unit同时引用两端的共现行——该共现关系记录时是真实的但见证它的所有单元都已被删除。提交时按 bank 去重因此针对同一 bank 的并发触发会合并为一次排空。每次运行以有界批次在墙钟预算内提交单次运行无法处理完的积压如批量删除不算错误——本次运行报告已完成的部分下一次运行从断点继续。触发条件任何删除memory_units的操作——DELETE /documents/{id}、DELETE /memories/{id}以及重 retain 已有document_idupsert 路径。整库清空delete_bank是空操作bank 里已没有需要维护的内容。webhook_delivery某些操作完成后例如某个注册了 webhook 的 bank 完成 consolidationHindsight 会入队一个webhook_delivery任务。handler 向配置的 URL POST 载荷并在瞬时失败时重试。webhook 表结构及投递相关实现在 webhooks/manager.py。管理端 API 端点以下所有路径均以bank_id为作用域。列出操作GET /v1/default/banks/{bank_id}/operations查询参数Param说明status按pending、processing、completed、failed、cancelled过滤。type按retain、file_convert_retain、consolidation、refresh_mental_model、graph_maintenance、webhook_delivery过滤。limit1–100默认 20。offset分页偏移量。exclude_parents从结果中排除父批次操作大型retain_batch调用会创建 1 个父操作 N 个子操作。端点的 HTTP 实现位于 api/http.py参数校验与文档一致limit有ge1, le100约束。Node.js SDK 示例// 列出某 bank 最近的操作默认最近 20 条。 const { data: recent } await sdk.listOperations({ client: apiClient, path: { bank_id: my-bank }, }); for (const op of recent.operations) { console.log(op.id, op.task_type, op.status); } // 按状态和类型过滤。 const { data: pendingRecompute } await sdk.listOperations({ client: apiClient, path: { bank_id: my-bank }, query: { status: pending, type: graph_maintenance }, }); // 隐藏 retain_batch 父行只显示独立的子 retain 任务。 const { data: flat } await sdk.listOperations({ client: apiClient, path: { bank_id: my-bank }, query: { exclude_parents: true }, });CLIhindsight operation list my-bankPython SDK 提供OperationsApi见 operations_api.py包含list_operations、get_operation_status、cancel_operation、retry_operation等异步方法。items_count是操作专属字段——仅 retain 形态的操作非零统计提交中的内容条目数。查询操作状态GET /v1/default/banks/{bank_id}/operations/{operation_id}Node.js SDK 示例const { data: status } await sdk.getOperationStatus({ client: apiClient, path: { bank_id: my-bank, operation_id: 550e8400-e29b-41d4-a716-446655440000 }, }); console.log(status.status, status.error_message); // 附带提交载荷retain 批次可能很大。 const { data: detailed } await sdk.getOperationStatus({ client: apiClient, path: { bank_id: my-bank, operation_id: 550e8400-e29b-41d4-a716-446655440000 }, query: { include_payload: true }, });CLIhindsight operation get my-bank $OPERATION_ID查询参数Param说明include_payload在响应中以task_payload返回原始任务载荷提交参数。默认false可能很大。值得关注的响应字段Field说明updated_at操作行最后一次变更时间——认领、进度心跳或完成。progress运行中操作最近一次进度快照未记录则为null瞬时完成或该功能上线前的行。task_payload原始提交参数仅在include_payloadtrue时填充。progress在粗粒度阶段/批次边界写入consolidation、批量 retain可用来区分健康的长任务与卡死的任务若轮询时processed持续前进说明任务存活若数字不变且at不再更新说明卡住了。其结构Field说明stage操作最近上报的粗粒度阶段如processing_batch。at快照写入时的 ISO-8601 时间戳。processed目前完成的工作单元数子批次、记忆已知时提供。total操作的总工作单元数已知时提供。detail操作专属计数器如observations_created、round、items_in_sub_batch。取消操作DELETE /v1/default/banks/{bank_id}/operations/{operation_id}可取消pending或processing状态的操作若已达终态completed、failed、cancelled则返回409。该逻辑在 memory_engine.py 的 cancel_operation 中实现UPDATE语句中再次校验状态因此SELECT 之后、UPDATE 之前恰好到达终态写入的任务会赢得竞态不会被误报为已取消。取消是协作式cooperative且非即时的行会被立即标记为cancelled但正在执行它的 worker 只会在下一个检查点注意到——retain 在子批次/文档边界检查consolidation 在 LLM 批次之间检查——并在那里停下。因此已提交的工作保持已提交正在进行的批次可能仍会跑完。没有检查点的操作类型会一直执行到结束无论如何行都保持cancelled因为任何 worker 写入都不会覆盖该状态_mark_failed中WHERE status cancelled的守卫即为此设计。这也正是清理被杀死 worker 卡在 processing的操作的方法已无任何进程在运行取消会立即生效。之后用POST /…/operations/{id}/retry重新入队即可。Node.js SDK 示例// 取消一个尚未被 worker 认领的 pending 操作。 // 若操作已 processing/completed/failed 则返回 409。 await sdk.cancelOperation({ client: apiClient, path: { bank_id: my-bank, operation_id: 550e8400-e29b-41d4-a716-446655440000 }, });CLIhindsight operation cancel my-bank $OPERATION_ID重试失败操作POST /v1/default/banks/{bank_id}/operations/{operation_id}/retry行的状态重置为pendingworker 会再次拾取。若操作不在failed或cancelled状态则返回409。Node.js SDK 示例// 重新入队一个失败或已取消的操作。 // 若操作不在 failed/cancelled 状态则返回 409。 await sdk.retryOperation({ client: apiClient, path: { bank_id: my-bank, operation_id: 550e8400-e29b-41d4-a716-446655440000 }, });CLIhindsight operation retry my-bank $OPERATION_ID异步 retain 完整示例提交一个大批量异步任务并轮询直至完成Node.js// 异步提交大批量——调用立即返回带 operation_id 的结果可用于轮询。 const submission await client.retainBatch(my-bank, [ { content: Alice joined Google in 2023 }, { content: Bob prefers Python over JavaScript }, ], { async: true }); const operationId submission.operation_id; while (true) { const { data: s } await sdk.getOperationStatus({ client: apiClient, path: { bank_id: my-bank, operation_id: operationId }, }); if ([completed, failed, cancelled].includes(s.status)) { console.log(finished: ${s.status}); break; } await new Promise((r) setTimeout(r, 2000)); }CLIbash jq# 提交异步 retain并从 JSON 响应中捕获 operation_id。 OPERATION_ID$( hindsight memory retain my-bank Alice joined Google in 2023 --async -o json \ | jq -r .operation_id ) # 轮询直至 worker 完成——completed/failed/cancelled 均为终态。 while true; do STATUS$(hindsight operation get my-bank $OPERATION_ID -o json | jq -r .status) if [ $STATUS completed ] || [ $STATUS failed ] || [ $STATUS cancelled ]; then echo finished: $STATUS break fi sleep 2 doneWorker 并发调优每个 worker 有一个跨所有操作类型共享的并发预算HINDSIGHT_API_WORKER_MAX_SLOTS默认 10见 config.py。按类型的槽位保留HINDSIGHT_API_WORKER_TYPE_RESERVED_SLOTS在该预算内划出保证容量剩余槽位形成任何类型都可用的共享池。从源码看槽位模型在 poller.py 的_get_available_slots中实现核心语义是保留槽位是下限floor而非上限ceiling——某类型超出其保留数时超出部分占用共享池默认保留仅consolidation: 2其余类型为 0见 config.pyHINDSIGHT_API_WORKER_TYPE_MAX_SLOTS是旧名deprecated alias历史上一直是设置保留下限而非上限仍可用但会打警告日志见 config.py。对大多数部署默认值已足够。仅当某类型被另一类型的洪流饿死时才需要为其保留槽位——例如在删除密集型负载下长时间运行的file_convert_retain阻塞了graph_maintenance。槽位还会跨 bank 轮转。每次认领先服务下一个 bank每次一个操作然后用剩余槽位按最旧优先从任何位置填充。因此批量摄入的 bank 无法独占整个池让另一个 bank 的单个写入在积压之后干等也不会被限流——当没有其他 bank 等待时它仍然会占用所有槽位。轮转实现在claim_batch的两遍扫描中第一遍公平轮转每个活跃 schema 每池最多 1 个认领第二遍容量回填只从第一遍有工作的 schema 填充剩余槽位同时通过next_bank_cursor记录每个 schema 内上次服务的 bank。此外worker 还支持HINDSIGHT_API_WORKER_CONSOLIDATION_BANK_PRIORITY按 bank 名称模式的优先级映射*为通配符*裸键为兜底默认用于让 consolidation 任务按优先级分层认领而非纯created_at顺序。wall-clock 兜底超时也值得一提retain 类任务受HINDSIGHT_API_RETAIN_WALL_TIMEOUT绝对上限约束consolidation 受HINDSIGHT_API_CONSOLIDATION_WALL_TIMEOUT的无进展上限约束每次提交批次会重置时钟见 poller.py——这能把卡死的任务从永久占用槽位变为失败且可重试。下一步阅读Documents — 追踪文档来源Memory Banks — 配置 bank 设置【免费下载链接】hindsightHindsight: Agent Memory That Learns项目地址: https://gitcode.com/GitHub_Trending/hindsight2/hindsight创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表