
AI Agent 技术月度展望总结 7 月的收获并规划 8 月的学习与实践方向一、深度引言与场景痛点7 月收官了。10 篇文章写完回头看这个月在 Agent 开发上走过的路既有成就感也有遗憾。成就感是实打实的Agent 流水线从裸奔到分层架构上线RAG 系统延迟从 5 秒压到 200ms向量检索从单编码切换到混合检索Prompt 工程从靠直觉变成了靠量化测试。这些都是这个月的硬产出。遗憾同样真实Agent 的并行编排还没实现依然在串行跑RAG 的检索精度虽然提升了不少但在长文档场景仍然会丢失关键细节向量检索前沿追踪了一堆真正落地的只有 Scalar 量化。计划做的事比做完的事多——这是每个月的常态但 7 月尤其明显。展望 8 月最迫切的痛点有三个并行编排的空白。7 月的 Agent 是串行流水线路由→记忆→规划→执行→审计每个环节都在等前一个完成。但实际业务中很多工具调用是独立的——查天气和查新闻可以同时进行不需要等查天气完了才查新闻。串行编排白白浪费了时间P99 延迟被拖到了 500ms。多 Agent 协作的低效。三个 Agent意图识别、执行、总结之间通过硬编码的 pipeline 串联没有动态分工。意图识别 Agent 在简单查询上是多余的——直接调用执行 Agent 就够了但当前架构不允许跳步。需要一个更灵活的编排引擎根据任务复杂度动态决定走哪条路径。RAG 检索精度的天花板。混合检索把精度从 72% 提到 89%但剩下的 11% 的失败案例几乎都是长文档中的细节信息被截断丢失。智能截断策略虽然减少了 Token 数量但也丢掉了部分关键段落。需要一个更精细的段落级检索方案。下面这张图梳理了 7 月的成果闭环和 8 月的规划方向二、底层机制与原理深度剖析8 月的三个核心规划方向背后的原理分别对应 Agent 系统的三个进化维度DAG 并行编排——从串行流水线到 DAG 执行引擎。串行编排的问题是显而易见的每一步都要等前一步完成。但大部分 Agent 工具调用之间存在可并行的独立关系。比如一个旅行规划 Agent用户问明天北京天气怎么样、有什么好吃的餐厅查天气和查餐厅是完全独立的两个工具调用不需要串行等待。DAG 编排引擎的核心是将任务分解为 DAG有向无环图DAG 中的节点是任务步骤边是依赖关系。无依赖的节点可以并行执行有依赖的节点必须等前置节点完成。执行引擎按照拓扑排序调度任务每个层内的节点全部并行不同层之间的节点串行等待。这样既保证了依赖关系的正确性又最大化了并行度。动态路由引擎——从固定流水线到自适应编排。当前架构是固定的 Router→Memory→Planner→Executor→Auditor 五步流水线。但实际场景中简单查询北京今天天气只需要 Router→Executor 两步就够了不需要 Planner 拆解、不需要 Auditor 校验。复杂查询帮我规划下周的出差行程考虑天气、交通、酒店预订才需要完整的五步流水线。动态路由引擎的原理是在 Router 之后加一个复杂度评估器根据任务需要的工具数量、推理步骤数、上下文依赖深度等指标动态选择执行路径。简单任务走快速路径Router→Executor复杂任务走完整路径Router→Memory→Planner→Executor→Auditor。这样简单查询的延迟可以控制在 200ms 以内复杂查询走完整流程保证质量。段落级检索——从文档级到 sentence-level 的精度跃升。当前 RAG 系统的检索粒度是文档级——每条向量对应一篇完整文档。检索 Top-K 返回的是 K 篇文档然后在文档内做截断。这种方式的问题是一篇 5000 字的文档可能只有一小段与用户查询相关但检索时整篇文档作为一个向量语义被平均化了——相关段落和不相关段落混在一起向量表示不够精确。段落级检索的原理是把每篇文档按 sentence 或 paragraph 拆分每个段落单独生成 embedding 和向量索引。检索时直接匹配段落级向量返回的 Top-K 是 K 个最相关的段落而不是 K 篇文档。这样检索精度从文档级的 89% 可以提升到段落级的 95% 以上。代价是向量数量增加了一篇 100 段的文档变成 100 个向量但通过 metadata 关联段落→文档→来源可以轻松回溯完整上下文。三、生产级代码实现以下是一个 DAG 并行编排引擎的实现整合了动态路由和自适应路径选择import asyncio import json import time from dataclasses import dataclass, field from enum import Enum from typing import Any import structlog logger structlog.get_logger() # DAG 任务模型 class TaskStatus(Enum): PENDING pending RUNNING running SUCCESS success FAILED failed SKIPPED skipped class TaskComplexity(Enum): SIMPLE simple # 1-2 步无需规划 MEDIUM medium # 3-5 步需要规划 COMPLEX complex # 6 步需要完整流水线 dataclass class DAGTask: DAG 中的单个任务节点。 task_id: str name: str handler_name: str # 对应的处理函数名 args: dict[str, Any] field(default_factorydict) depends_on: list[str] field(default_factorylist) # 依赖的任务 ID status: TaskStatus TaskStatus.PENDING result: Any None error: str | None None latency_ms: float 0.0 retry_count: int 0 max_retries: int 2 # 动态路由引擎 class DynamicRouter: 动态路由引擎根据任务复杂度选择执行路径。 7 月痛点固定五步流水线对简单查询浪费步骤。 8 月方案复杂度评估 → 自适应路径选择。 def evaluate_complexity(self, query: str) - TaskComplexity: 评估任务复杂度。 # 规则 1工具数量推断 tool_indicators [ 查, 搜索, 找, # 单工具 帮我, 规划, 对比, 分析, # 多工具 同时, 还要, 并且, # 并行工具 ] single_tool_count sum( 1 for kw in tool_indicators[:3] if kw in query ) multi_tool_count sum( 1 for kw in tool_indicators[3:6] if kw in query ) parallel_count sum( 1 for kw in tool_indicators[6:] if kw in query ) # 规则 2推理深度推断 reasoning_indicators [ 为什么, 原因, 如何, # 需要推理 评估, 优化, 建议, # 需要深度推理 ] reasoning_count sum( 1 for kw in reasoning_indicators if kw in query ) # 规则 3查询长度推断 query_length len(query.split()) # 综合评分 complexity_score ( multi_tool_count * 2 parallel_count * 3 reasoning_count * 2 min(query_length / 5, 3) ) if complexity_score 2: return TaskComplexity.SIMPLE elif complexity_score 6: return TaskComplexity.MEDIUM else: return TaskComplexity.COMPLEX def select_pipeline( self, complexity: TaskComplexity ) - list[str]: 根据复杂度选择执行路径的步骤列表。 pipelines { TaskComplexity.SIMPLE: [ guardrail, router, executor, ], TaskComplexity.MEDIUM: [ guardrail, router, memory, executor, auditor, ], TaskComplexity.COMPLEX: [ guardrail, router, memory, planner, executor, auditor, ], } return pipelines[complexity] def build_dag( self, query: str, pipeline_steps: list[str], ) - list[DAGTask]: 根据执行路径和查询内容构建 DAG 任务列表。 tasks [] for i, step in enumerate(pipeline_steps): # 基本依赖每步依赖前一步完成 depends_on [tasks[i - 1].task_id] if i 0 else [] # 特殊处理executor 步骤可能有并行子任务 if step executor: # 检查是否需要并行工具调用 parallel_tools self._extract_parallel_tools(query) if parallel_tools: # 创建并行子任务 for j, tool_name in enumerate(parallel_tools): sub_task DAGTask( task_idfexec_{tool_name}, namef并行执行: {tool_name}, handler_nametool_name, args{query: query}, depends_on[tasks[i - 1].task_id], # 依赖前一步 ) tasks.append(sub_task) continue # 跳过串行 executor 节点 task DAGTask( task_idfstep_{step}, namef步骤: {step}, handler_namestep, args{query: query}, depends_ondepends_on, ) tasks.append(task) logger.info( dag_built, queryquery[:50], stepslen(tasks), pipelinepipeline_steps, ) return tasks def _extract_parallel_tools(self, query: str) - list[str]: 从查询中提取可并行执行的工具列表。 # 简化版检测并列关键词 parallel_keywords { 天气: weather_api, 餐厅: restaurant_api, 新闻: news_api, 酒店: hotel_api, 地图: map_api, 汇率: exchange_api, } tools [] for keyword, tool_name in parallel_keywords.items(): if keyword in query: tools.append(tool_name) return tools # DAG 执行引擎 class DAGExecutor: DAG 并行编排执行引擎。 核心逻辑按拓扑层级调度任务同一层级内的任务并行执行。 def __init__(self, max_parallel: int 5): self.max_parallel max_parallel self._handlers: dict[str, Any] {} self._semaphore asyncio.Semaphore(max_parallel) def register_handler(self, name: str, handler): 注册任务处理函数。 self._handlers[name] handler async def execute_dag( self, tasks: list[DAGTask], context: dict[str, Any] | None None, ) - dict[str, Any]: 执行 DAG 任务列表。 context context or {} task_map {t.task_id: t for t in tasks} results: dict[str, Any] {} total_start time.monotonic() # 计算拓扑层级每层的任务互不依赖可并行 layers self._compute_topological_layers(tasks, task_map) logger.info( dag_execution_start, total_taskslen(tasks), layerslen(layers), ) # 按层执行 for layer_idx, layer_tasks in enumerate(layers): logger.info( dag_layer_start, layerlayer_idx, taskslen(layer_tasks), ) # 同一层的任务并行执行 layer_results await asyncio.gather( *[self._execute_single_task(t, context, results) for t in layer_tasks], return_exceptionsTrue, ) # 收集结果 for task, result in zip(layer_tasks, layer_results): if isinstance(result, Exception): task.status TaskStatus.FAILED task.error str(result) results[task.task_id] {error: str(result)} logger.error( task_failed, task_idtask.task_id, errorstr(result), ) else: task.status TaskStatus.SUCCESS task.result result results[task.task_id] result total_latency (time.monotonic() - total_start) * 1000 success_count sum( 1 for t in tasks if t.status TaskStatus.SUCCESS ) logger.info( dag_execution_complete, total_taskslen(tasks), successsuccess_count, failedlen(tasks) - success_count, latency_msround(total_latency, 1), ) return { results: results, success_count: success_count, failed_count: len(tasks) - success_count, total_latency_ms: round(total_latency, 1), } async def _execute_single_task( self, task: DAGTask, context: dict[str, Any], completed_results: dict[str, Any], ) - Any: 执行单个 DAG 任务带限流和重试。 handler self._handlers.get(task.handler_name) if handler is None: # 没有注册的 handler用默认模拟处理 async with self._semaphore: start time.monotonic() await asyncio.sleep(0.05) # 模拟处理时间 task.latency_ms round((time.monotonic() - start) * 1000, 1) return {task: task.name, status: mock_success} # 将前置任务的结果注入上下文 task_args dict(task.args) for dep_id in task.depends_on: if dep_id in completed_results: task_args[fdep_{dep_id}] completed_results[dep_id] # 带限流的执行 for attempt in range(task.max_retries 1): try: async with self._semaphore: start time.monotonic() result await handler(task_args, context) task.latency_ms round( (time.monotonic() - start) * 1000, 1 ) return result except Exception as e: task.retry_count attempt 1 logger.warning( task_retry, task_idtask.task_id, attemptattempt 1, errorstr(e), ) if attempt task.max_retries: raise def _compute_topological_layers( self, tasks: list[DAGTask], task_map: dict[str, DAGTask], ) - list[list[DAGTask]]: 计算拓扑层级将 DAG 分解为可并行执行的层。 layers: list[list[DAGTask]] [] completed_ids: set[str] set() remaining list(tasks) while remaining: # 找出所有依赖已满足的任务即前置任务都已完成或无依赖 ready_tasks [ t for t in remaining if all( dep in completed_ids or dep not in task_map for dep in t.depends_on ) ] if not ready_tasks: # 存在循环依赖或无法解析的依赖 logger.error(dag_deadlock, remaininglen(remaining)) # 将剩余任务标记为 SKIPPED for t in remaining: t.status TaskStatus.SKIPPED t.error dependency_deadlock break layers.append(ready_tasks) for t in ready_tasks: completed_ids.add(t.task_id) remaining [t for t in remaining if t not in ready_tasks] return layers # 段落级检索方案 dataclass class Passage: 段落级检索单元。 passage_id: str content: str doc_id: str source: str position: int # 在文档中的位置第几个段落 embedding: list[float] field(default_factorylist) class PassageLevelSearcher: 段落级检索器8 月的检索精度优化方案。 解决痛点文档级检索丢失长文档中的关键段落。 def __init__(self, top_k_passages: int 10, top_k_docs: int 5): self.top_k_passages top_k_passages self.top_k_docs top_k_docs self._passages: list[Passage] [] def split_document_to_passages( self, doc_id: str, content: str, source: str, min_passage_length: int 100, ) - list[Passage]: 将文档按段落拆分为检索单元。 # 按自然段落分割 raw_paragraphs content.split(\n\n) passages [] position 0 for para in raw_paragraphs: # 过滤过短的段落 if len(para.strip()) min_passage_length: continue # 对超长段落再做 sentence 级拆分 if len(para) 500: sentences self._split_to_sentences(para) # 将连续的 sentence 合并为 200-400 token 的片段 chunks self._merge_sentences_to_chunks( sentences, max_chunk_length400 ) for chunk in chunks: p Passage( passage_idf{doc_id}_p{position}, contentchunk, doc_iddoc_id, sourcesource, positionposition, ) passages.append(p) position 1 else: p Passage( passage_idf{doc_id}_p{position}, contentpara, doc_iddoc_id, sourcesource, positionposition, ) passages.append(p) position 1 logger.info( document_split, doc_iddoc_id, total_passageslen(passages), original_lengthlen(content), ) return passages def _split_to_sentences(self, text: str) - list[str]: 按句号、问号、感叹号分割句子。 import re sentences re.split(r[。\.\?\!]\s*, text) return [s.strip() for s in sentences if len(s.strip()) 20] def _merge_sentences_to_chunks( self, sentences: list[str], max_chunk_length: int 400 ) - list[str]: 将句子合并为固定长度的片段。 chunks [] current_chunk for sentence in sentences: if len(current_chunk) len(sentence) max_chunk_length: if current_chunk: chunks.append(current_chunk.strip()) current_chunk sentence else: current_chunk sentence 。 if current_chunk.strip(): chunks.append(current_chunk.strip()) return chunks async def search( self, query_embedding: list[float], max_results: int 10, ) - list[dict[str, Any]]: 段落级检索返回最相关的段落按文档聚合。 # 实际项目中对接向量数据库的段落级搜索 # 这里用模拟数据演示 await asyncio.sleep(0.02) # 模拟检索延迟 # 模拟段落级检索结果 mock_results [ { passage_id: doc_1_p3, content: 关于 asyncio 的最佳实践段落..., doc_id: doc_1, source: tech_blog, score: 0.92, position: 3, }, { passage_id: doc_2_p1, content: asyncio 事件循环的核心机制..., doc_id: doc_2, source: official_doc, score: 0.88, position: 1, }, ] # 按文档聚合同一个文档的多个段落合并展示 aggregated self._aggregate_by_document(mock_results) logger.info( passage_search_complete, query_dimlen(query_embedding), passageslen(mock_results), documentslen(aggregated), ) return aggregated[:max_results] def _aggregate_by_document( self, passage_results: list[dict] ) - list[dict[str, Any]]: 将段落结果按文档聚合保留最佳段落作为代表。 doc_groups: dict[str, list[dict]] {} for p in passage_results: doc_id p[doc_id] if doc_id not in doc_groups: doc_groups[doc_id] [] doc_groups[doc_id].append(p) aggregated [] for doc_id, passages in doc_groups.items(): # 取该文档中 score 最高的段落 best_passage max(passages, keylambda p: p[score]) aggregated.append({ doc_id: doc_id, source: best_passage.get(source, ), best_passage: best_passage[content], best_score: best_passage[score], total_matching_passages: len(passages), all_passages: passages, }) # 按最高段落 score 降序排序 aggregated.sort(keylambda d: d[best_score], reverseTrue) return aggregated # 月度复盘工具 dataclass class MonthlyGoal: 月度目标追踪。 name: str category: str # agent, rag, infra, learning status: str # done, partial, not_started completion_pct: float # 0-100 lessons: str # 收获/教训 next_action: str # 8 月后续行动 class MonthlyReviewPlanner: 月度复盘与 8 月规划工具。 def __init__(self): self.july_goals: list[MonthlyGoal] [] self.august_goals: list[MonthlyGoal] [] self._load_july_review() def _load_july_review(self): 加载 7 月目标复盘。 self.july_goals [ MonthlyGoal( nameAgent 分层架构, categoryagent, statusdone, completion_pct100, lessons分层解耦审计节点是 Agent 生产化的基础, next_action升级为 DAG 并行编排, ), MonthlyGoal( nameRAG 延迟压缩, categoryrag, statusdone, completion_pct100, lessons逐段定位瓶颈逐轮定向爆破, next_action段落级检索提升精度, ), MonthlyGoal( name混合检索落地, categoryrag, statusdone, completion_pct90, lessons稀疏稠密双编码精度提升 17%, next_action自适应权重调参 DiskANN 评估, ), MonthlyGoal( nameAgent 并行编排, categoryagent, statuspartial, completion_pct30, lessons串行架构的并行化改造比预期复杂, next_action实现 DAG 执行引擎, ), MonthlyGoal( name动态路由引擎, categoryagent, statusnot_started, completion_pct0, lessons固定流水线对简单查询浪费步骤, next_action实现复杂度评估自适应路径, ), MonthlyGoal( nameasync 工具箱, categoryinfra, statusdone, completion_pct100, lessons限流Task管理HTTP单例是基础能力, next_action增加熔断器模式, ), ] self.august_goals [ MonthlyGoal( nameDAG 并行编排引擎, categoryagent, statusnot_started, completion_pct0, lessons, next_action实现拓扑排序调度 层级并行执行, ), MonthlyGoal( name动态路由引擎, categoryagent, statusnot_started, completion_pct0, lessons, next_action实现复杂度评估器 三级路径选择, ), MonthlyGoal( name段落级检索, categoryrag, statusnot_started, completion_pct0, lessons, next_action文档拆分为 passage 独立向量化, ), MonthlyGoal( nameAgent 记忆持久化, categoryagent, statusnot_started, completion_pct0, lessons, next_action跨 session 的长期记忆存储和检索, ), MonthlyGoal( nameScalar 量化落地, categoryrag, statusnot_started, completion_pct0, lessons, next_actionQdrant Scalar Quantization 部署和验证, ), MonthlyGoal( name全链路可观测性, categoryinfra, statusnot_started, completion_pct0, lessons, next_actionAgent 全链路 trace Token 成本追踪, ), ] def get_review_summary(self) - dict[str, Any]: 7 月复盘摘要。 done [g for g in self.july_goals if g.status done] partial [g for g in self.july_goals if g.status partial] not_done [g for g in self.july_goals if g.status not_started] avg_completion sum(g.completion_pct for g in self.july_goals) / len(self.july_goals) return { july_total_goals: len(self.july_goals), completed: len(done), partial: len(partial), not_started: len(not_done), avg_completion_pct: round(avg_completion, 1), top_lessons: [ g.lessons for g in done if g.lessons ], carry_over: [ g.name for g in partial not_done ], } def get_august_plan(self) - dict[str, Any]: 8 月规划摘要。 priorities { P0: [g for g in self.august_goals if g.category in [agent, rag]], P1: [g for g in self.august_goals if g.category infra], } return { august_total_goals: len(self.august_goals), p0_count: len(priorities[P0]), p1_count: len(priorities[P1]), p0_goals: [g.name for g in priorities[P0]], p1_goals: [g.name for g in priorities[P1]], focus_areas: list(set(g.category for g in self.august_goals)), } async def main(): 演示动态路由 DAG 执行 月度复盘。 # 1. 动态路由 router DynamicRouter() queries [ 北京今天天气怎么样, # simple 帮我对比 Redis 和 Milvus 的性能差异, # medium 规划下周出差行程查天气、订酒店、查交通, # complex ] for query in queries: complexity router.evaluate_complexity(query) pipeline router.select_pipeline(complexity) dag_tasks router.build_dag(query, pipeline) logger.info( dynamic_routing, queryquery[:50], complexitycomplexity.value, pipelinepipeline, dag_taskslen(dag_tasks), ) # 2. DAG 执行 executor DAGExecutor(max_parallel5) # 注册模拟 handler async def mock_handler(args: dict, ctx: dict) - dict: await asyncio.sleep(0.05) return {result: fmock_{args.get(query, )[:20]}} for step in [guardrail, router, executor, memory, planner, auditor]: executor.register_handler(step, mock_handler) # 构建并执行 DAG tasks router.build_dag(queries[2], router.select_pipeline(TaskComplexity.COMPLEX)) result await executor.execute_dag(tasks) logger.info(dag_result, **result) # 3. 月度复盘 planner MonthlyReviewPlanner() review planner.get_review_summary() august_plan planner.get_august_plan() logger.info(july_review, **review) logger.info(august_plan, **august_plan) if __name__ __main__: asyncio.run(main())四、边界分析与架构权衡DAG 编排 vs 固定流水线DAG 编排灵活度高但调试复杂度也高。固定流水线串行执行行为可预测、日志易追踪DAG 编排并行执行任务间依赖关系动态变化一个节点失败可能触发多条路径的异常处理。我的策略是第一版 DAG 只支持两层并行预计算层 执行层不搞复杂的 DAG 嵌套。等两层并行稳定了再增加层级。动态路由的误判风险复杂度评估器是规则驱动的不是模型驱动的。规则简单可靠但覆盖不了所有场景。帮我查一下 Python 的装饰器用法——规则判断为 SIMPLE单工具但实际上用户可能期待详细的代码示例和对比分析需要 MEDIUM 级别的路径。更精细的方案是用小模型做意图分类但增加了延迟和成本。段落级检索的向量膨胀一篇 100 段的文档变成 100 个向量10 万篇文档就是 1000 万个向量。向量数量增加 10 倍索引构建时间、内存占用、检索延迟都受影响。解决方案是段落级索引只对高频查询的热门文档启用冷门文档保留文档级索引。这种分级索引策略在向量数量和检索精度之间做了折中。8 月目标数量控制规划了 6 个目标但 8 月只有 4 周。我的原则是 P0 目标最多 3 个DAG 编排、动态路由、段落级检索P1 目标按精力弹性推进记忆持久化、量化落地、可观测性。宁可少做几个但做完不要贪多但每个都半成品。五、总结7 月收官10 篇文章写完。回头看这个月在 Agent 和 RAG 上走了很远的路但还有更长的路要走。7 月的硬产出是六件事Agent 分层架构、RAG 延迟压缩、混合检索、Prompt 量化方法论、Redis 向量运维、async 工具箱。每一件都是从痛点出发、有代码落地、有数据验证的实打实的成果。没有值得关注式发现没有推荐阅读级突破只有一个个具体的问题被一个个具体的方案解决。7 月的遗憾也是三个Agent 并行编排只实现了 30%、动态路由引擎没启动、长文档检索精度仍有瓶颈。这些遗憾不是失败是未完成——它们已经从痛点变成了明确的 8 月规划。8 月的方向很清晰DAG 并行编排是 Agent 进化的下一步动态路由引擎是效率优化的关键段落级检索是精度跃升的路径。三个方向都是从 7 月的遗憾中生长出来的——每一个未解决的问题都指向了下一步的明确方向。做技术规划的本质不是列清单是排优先级。6 个目标看起来很多但 P0 只有 3 个。3 个 P0 目标在 4 周内做完是有把握的——每个目标 1 周的核心开发 1 周的验证打磨。P1 目标看精力弹性推进做不完就顺延到 9 月。不焦虑不贪多一步步走。7 月充实收官。8 月继续迭代。资料说明本文中的协议、版本、性能、成本和行业趋势应以可核验的一手资料为准。未标注统计口径的比例、时间表和预测仅作工程讨论不应视为行业事实。可参考 0731 资料来源索引并在发布前将具体来源贴到对应断言之后。