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_factory=dict)
    depends_on: list[str] = field(default_factory=list)  # 依赖的任务 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_id=f"exec_{tool_name}",
                            name=f"并行执行: {tool_name}",
                            handler_name=tool_name,
                            args={"query": query},
                            depends_on=[tasks[i - 1].task_id],  # 依赖前一步
                        )
                        tasks.append(sub_task)
                    continue  # 跳过串行 executor 节点

            task = DAGTask(
                task_id=f"step_{step}",
                name=f"步骤: {step}",
                handler_name=step,
                args={"query": query},
                depends_on=depends_on,
            )
            tasks.append(task)

        logger.info(
            "dag_built",
            query=query[:50],
            steps=len(tasks),
            pipeline=pipeline_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_tasks=len(tasks),
            layers=len(layers),
        )

        # 按层执行
        for layer_idx, layer_tasks in enumerate(layers):
            logger.info(
                "dag_layer_start",
                layer=layer_idx,
                tasks=len(layer_tasks),
            )

            # 同一层的任务并行执行
            layer_results = await asyncio.gather(
                *[self._execute_single_task(t, context, results) for t in layer_tasks],
                return_exceptions=True,
            )

            # 收集结果
            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_id=task.task_id,
                        error=str(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_tasks=len(tasks),
            success=success_count,
            failed=len(tasks) - success_count,
            latency_ms=round(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[f"dep_{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_id=task.task_id,
                    attempt=attempt + 1,
                    error=str(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", remaining=len(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_factory=list)

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_length=400
                )
                for chunk in chunks:
                    p = Passage(
                        passage_id=f"{doc_id}_p{position}",
                        content=chunk,
                        doc_id=doc_id,
                        source=source,
                        position=position,
                    )
                    passages.append(p)
                    position += 1
            else:
                p = Passage(
                    passage_id=f"{doc_id}_p{position}",
                    content=para,
                    doc_id=doc_id,
                    source=source,
                    position=position,
                )
                passages.append(p)
                position += 1

        logger.info(
            "document_split",
            doc_id=doc_id,
            total_passages=len(passages),
            original_length=len(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_dim=len(query_embedding),
            passages=len(mock_results),
            documents=len(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, key=lambda 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(key=lambda d: d["best_score"], reverse=True)
        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(
                name="Agent 分层架构",
                category="agent",
                status="done",
                completion_pct=100,
                lessons="分层解耦+审计节点是 Agent 生产化的基础",
                next_action="升级为 DAG 并行编排",
            ),
            MonthlyGoal(
                name="RAG 延迟压缩",
                category="rag",
                status="done",
                completion_pct=100,
                lessons="逐段定位瓶颈,逐轮定向爆破",
                next_action="段落级检索提升精度",
            ),
            MonthlyGoal(
                name="混合检索落地",
                category="rag",
                status="done",
                completion_pct=90,
                lessons="稀疏+稠密双编码,精度提升 17%",
                next_action="自适应权重调参 + DiskANN 评估",
            ),
            MonthlyGoal(
                name="Agent 并行编排",
                category="agent",
                status="partial",
                completion_pct=30,
                lessons="串行架构的并行化改造比预期复杂",
                next_action="实现 DAG 执行引擎",
            ),
            MonthlyGoal(
                name="动态路由引擎",
                category="agent",
                status="not_started",
                completion_pct=0,
                lessons="固定流水线对简单查询浪费步骤",
                next_action="实现复杂度评估+自适应路径",
            ),
            MonthlyGoal(
                name="async 工具箱",
                category="infra",
                status="done",
                completion_pct=100,
                lessons="限流+Task管理+HTTP单例是基础能力",
                next_action="增加熔断器模式",
            ),
        ]

        self.august_goals = [
            MonthlyGoal(
                name="DAG 并行编排引擎",
                category="agent",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="实现拓扑排序调度 + 层级并行执行",
            ),
            MonthlyGoal(
                name="动态路由引擎",
                category="agent",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="实现复杂度评估器 + 三级路径选择",
            ),
            MonthlyGoal(
                name="段落级检索",
                category="rag",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="文档拆分为 passage + 独立向量化",
            ),
            MonthlyGoal(
                name="Agent 记忆持久化",
                category="agent",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="跨 session 的长期记忆存储和检索",
            ),
            MonthlyGoal(
                name="Scalar 量化落地",
                category="rag",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="Qdrant Scalar Quantization 部署和验证",
            ),
            MonthlyGoal(
                name="全链路可观测性",
                category="infra",
                status="not_started",
                completion_pct=0,
                lessons="",
                next_action="Agent 全链路 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",
            query=query[:50],
            complexity=complexity.value,
            pipeline=pipeline,
            dag_tasks=len(dag_tasks),
        )

    # 2. DAG 执行
    executor = DAGExecutor(max_parallel=5)

    # 注册模拟 handler
    async def mock_handler(args: dict, ctx: dict) -> dict:
        await asyncio.sleep(0.05)
        return {"result": f"mock_{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 资料来源索引,并在发布前将具体来源贴到对应断言之后。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐