一个分支挂掉,答案整个没了
目标
把一个问题同时撒向多个数据源,再汇总回来。亲自确认同时写入同一格会抛什么异常、合并的顺序由什么决定,并做出数量在运行时确定的分支、宽度上限和部分成功。
为什么重要
LangGraph 的执行以超级步为单位,一个超级步里能去的节点有多个时,这些节点就一起运行。所以增加分支只需要几行边。难的是合并——当多个值同时进入同一个键,而没有 Reducer 时,图会因 InvalidUpdateError 而停下。在依次运行的图里悄悄被覆盖的 bug,会在改成并行的那一刻暴露出来,所以这个异常其实算是友善的。
合并的顺序也容易误解。不是“先结束的排在前面”。本实验中要通过改名字来亲自确认。
数据源不固定的话,就无法提前画出边。Send 是“用这个输入运行这个节点”的指令,列表有多长,那个节点就运行多少次。而且那个节点看到的不是完整状态,而是 Send 传来的片段。
最后,如果一个分支抛出异常,整个超级步就会崩溃,已经成功的分支的结果也会一起消失。把失败作为值返回,就能做出部分成功;留下失败的数据源的名称,部分结果就不会冒充成全部结果。
评分器不会相信你写下的说明。它会真正导入你的模块,用任意的主题和任意的节点名称运行,并把结果与评分器另外计算的值对照。
步骤
- 在 /root/work/agpara/fanout.py 中创建
SOURCES、hits、State、每个数据源对应的节点和build_graph()。从 START 分出与数据源数量一样多的分支,再流向 END。findings使用接续拼接的 Reducer。 - 增加
conflict_demo(topic),捕获没有 Reducer 的键被两个节点在同一个超级步写入时抛出的异常,并以{"error": ..., "message": ..., "note": ...}返回。 - 增加
merge_order(names),通过改变节点名称来亲自确认合并的顺序。 - 增加
combine节点,把各分支送往combine而不是 END。combine生成ranked(按条数降序,条数相同时按名称升序)和total。 - 增加
plan、fan_out、probe、build_dynamic()、search_dynamic(topic, picked),用Send做出数量在运行时确定的分支。 - 增加
MAX_FANOUT = 2和pick_sources(topic, limit=None),给宽度设上限。先挑条数多的数据源,条数相同时按名称顺序挑选。 - 增加
failures键和run_partial(topic),让一个分支找不到时,其余答案也能保留下来。 - 在 /root/work/agpara/fanout_report.json 和 /root/work/agpara/fanout_report.md 中记录确认的内容。
参考
- 执行契约:评分器会把
/root/work/agpara/fanout.py当作 Python 模块导入,直接使用上面列出的名称。不会作为脚本运行。 SOURCES是{자료원이름: {주제: 건수}}(占位符依次为数据源名称、主题、条数)。wiki为退款 3、配送 5、换货 2、包装 4,ticket为退款 7、配送 1、换货 4、包装 4,manual为退款 2、配送 6、换货 9、包装 1。包装在 wiki 和 ticket 中条数相同,所以会暴露出用什么来区分条数相同的情况。hits(source, topic)对不存在的主题返回 0。- 数据源节点返回
{"findings": [{"source": 이름, "hits": 건수}]}(占位符依次为名称与条数)。节点名称要与数据源名称相同。不能与状态键重名——重名的话,编译时会出现ValueError: 'x' is already being used as a state key。 conflict_demo另外建一个小图来用,里面有一个没有 Reducer 的键和一个有 Reducer 的键。捕获到异常时,在error中放入异常名称,在message中原样放入str(예외)(占位符为异常对象),note为空字符串。没有抛异常时,error和message为空,note中放入剩下的值。merge_order(names)用这些名称创建节点,并排放置,运行一次,返回合并后的列表。combine是ranked = [건수 내림차순, 동점이면 이름 오름차순으로 정렬한 자료원 이름](韩文,意为“按条数降序、条数相同时按名称升序排列的数据源名称”),total = 건수의 합(韩文,意为“条数之和”)。- 用
Send调用的节点接收的不是完整状态,而是Send传来的字典。search_dynamic原样返回findings列表。 pick_sources(topic, limit=None)在没有limit时使用MAX_FANOUT。0 以下则为空列表。run_partial(topic)的答案:{"ranked": [...], "total": 정수, "failures": [...정렬됨], "sources": 찾아낸 자료원 수}(占位符依次为整数、“已排序”、找到的数据源数量)。ticket对自己不知道的主题,不要抛异常,而是返回{"failures": ["ticket"]}。- 这个 Pod 没有互联网。langgraph 0.2.60 已经装好了。
- 官方文档:Use the graph API · Graph API overview · Types 参考
- 常见错误:让每个分支各自流向 END,再另设汇总节点(那样就无法知道汇总节点什么时候运行)、以为合并的顺序就是完成的顺序、在分支里抛异常、在宽度限制中随意区分条数相同的情况。
一次查三个地方
在 /root/work/agpara/fanout.py 中创建 SOURCES、hits、State、每个数据源对应的节点和 build_graph()。从 START 分出与数据源数量一样多的分支再流向 END,每个节点向 findings 写入自己的一条结果。
创建分支只要调用与数据源数量一样多次的 add_edge(START, 이름)(占位符为名称)就行了。重要的是给 findings 加上接续拼接的 Reducer——没有的话,下一步才会看到的异常会在这里先冒出来。节点名称与数据源名称相同,但不要与状态键重名。
同时写入同一格就会停下
增加 conflict_demo(topic),做一个让两个节点在同一个超级步写入没有 Reducer 的键的小图,捕获抛出的异常,并以 {"error": 예외이름, "message": 예외메시지, "note": ""}(占位符依次为异常名称与异常消息)返回。
是 langgraph.errors.InvalidUpdateError。不要编造消息,把 str(예외)(占位符为异常对象)原样放进去——哪个键有问题,上面写着。如果是依次运行的图,这本来会被悄悄覆盖,这一点就是这一步的要点。
合并的顺序不是偶然的
增加 merge_order(names)。用这些名称创建节点并排放置,运行一次,然后返回汇集在接续拼接的键中的列表。
把同样的名称只改变顺序再传进去试试。结果如何,就是这一步的答案。用循环创建节点时,要小心 lambda 只抓住最后一个名称的陷阱——要用默认参数或外层函数把名称绑定住。
分支全部结束后汇总一次
增加 combine 节点,把数据源节点送往 combine 而不是 END。combine 生成 ranked(按条数降序,条数相同时按名称升序)和 total(条数之和)。
把分支汇聚到汇总节点,那个节点就会在分支全部结束之后只运行一次。如果让每个分支各自流向 END,再另设汇总节点,就无法知道它什么时候运行。用名称区分条数相同的情况,是为了让答案是确定的。
数量在运行时确定
增加 plan、fan_out、probe、build_dynamic()、search_dynamic(topic, picked)。fan_out 返回 Send 列表,probe 只看 Send 传来的片段。
形式是 Send("노드이름", 넘길딕셔너리)(占位符依次为节点名称与要传递的字典)。条件边的第三个参数要给出能去的节点名称列表。在 probe 里用 state["source"] 可能会觉得陌生,那是因为该节点接收的不是完整状态,而是 Send 传来的内容。
给宽度设上限
增加 MAX_FANOUT = 2 和 pick_sources(topic, limit=None)。先挑条数多的数据源,条数相同时按名称升序挑选,没有 limit 时使用 MAX_FANOUT。0 以下则为空列表。
数据源变成二十个,就会一次发出二十次调用。比设上限更重要的是挑选规则必须是确定的。条数相同时随意挑选,同一个问题就会得到不同的答案,那样就无法比较昨天和今天。
一个失败,其余的也要保住
增加 failures 键和 run_partial(topic)。ticket 对自己不知道的主题不要抛异常,而是返回 {"failures": ["ticket"]}。答案是 {"ranked": [...], "total": 정수, "failures": [...정렬됨], "sources": 정수}(占位符依次为整数、“已排序”、整数)。
一个分支抛出异常,整个超级步就会崩溃,已经成功的分支的结果也会一起消失。把失败作为值返回,汇总节点就能知道“三个里有两个查到了,一个失败了”。不留下失败的数据源的名称,部分结果就会冒充成全部结果。
记录确认的内容
在 /root/work/agpara/fanout_report.json 中写入 conflict_error、merge_order、ranked、total、dynamic、picked、partial,在 /root/work/agpara/fanout_report.md 中用 ## 누가 같은 칸에 쓰는가(韩文,意为“谁在写入同一格”)、## 합쳐지는 순서는 무엇으로 정해지나(韩文,意为“合并的顺序由什么决定”)、## 갯수가 실행 시점에 정해질 때(韩文,意为“数量在运行时确定时”)、## 하나가 실패하면(韩文,意为“一个失败时”)四节来写。
merge_order 是 {"input": [...], "output": [...]},同时放入你传入的名称和得到的顺序。ranked、total 是以 환불(韩文,意为“退款”)主题运行图得到的结果,dynamic 是 search_dynamic 的答案,picked 是 pick_sources('배송')(韩文,意为“配送”)的答案,partial 是 run_partial('쿠폰')(韩文,意为“优惠券”)的答案中只取 failures 和 sources。全部运行后得到。