公司动态
企业实战:数据图与状态定义
目录1 定义状态 (State)2 定义主图2.1 第一步创建节点骨架2.2 第二步编写主图代码2.3 第三步验证图流程1 定义状态 (State)所有节点共享同一个状态对象。我们需要定义它来存储处理过程中的数据如 PDF 路径、MD 内容、切片列表、向量等。node_entry - 读 local_file_path 因为要判断输入是 .md 还是 .pdf 。 - 读 task_id 任务追踪。 - 写 is_md_read_enabled / is_pdf_read_enabled 给路由函数决定下一跳。 - 写 md_path / pdf_path 把原始文件路径标准化到后续节点约定字段。 - 写 file_title 后续 item_name 识别失败时做兜底语义。 node_pdf_to_md - 读 pdf_path PDF 解析入口文件。 - 读 local_dir 解析产物落盘目录。 - 读 task_id 任务追踪。 - 写 md_path 告诉下游 markdown 文件在哪里。 - 写 md_content 下游可直接用文本不必重复读盘。 - 回写 local_dir 确保是规范化后的实际目录字符串。 node_md_img - 读 md_path 定位 markdown 和 images/ 目录。 - 读 md_content 可空有就用没有就从 md_path 读取。 - 读 task_id 任务追踪。 - 写 md_content 替换图片链接后的新内容。 - 写 md_path 可能变成 _new.md 下游必须跟新路径走。 node_document_split - 读 md_content 切片原材料。 - 读 file_title 给每个 chunk 补 metadata语义归属。 - 读 local_dir 把 chunks.json 备份到本地。 - 读 task_id 任务追踪。 - 写 chunks 这是后续识别、向量化、入库的核心载体。 node_item_name_recognition - 读 chunks 从前几段内容构建上下文给 LLM 识别主体。 - 读 file_title / md_path 识别失败时兜底出 item_name 。 - 读 task_id 任务追踪。 - 写 item_name 文档级主语后续幂等删除、检索过滤都依赖。 - 写 chunks 把 item_name 回填到每个 chunk保证 chunk 自包含。 node_bge_embedding - 读 chunks 拿 item_name content 生成向量输入。 - 读 task_id 可选任务追踪。 - 写 chunks 补 dense_vector / sparse_vector 供 Milvus 入库。 node_import_milvus - 读 chunks 真正要入库的数据。 - 读 task_id 任务追踪。 - 用 chunks[0].item_name 做幂等删除同主体重复导入先删旧数据。 - 写 chunks 把插入回显的 chunk_id 回填方便后续追踪与引用。文件:app/import_process/agent/state.pyfromtypingimportTypedDictimportcopyfromapp.core.loggerimportloggerclassImportGraphState(TypedDict): 图的状态定义包含所有节点产生和消费的数据字段。 TypedDict 让我们在代码中能有自动补全和类型检查。 使用字典式访问如state[session_id]、state.get(embedding_chunks) task_id:str# 任务唯一ID用于追踪日志# --- 流程控制标记 ---is_md_read_enabled:bool# 是否启用 Markdown 读取路径is_pdf_read_enabled:bool# 是否启用 PDF 读取路径# --- 路径相关 ---local_dir:str# 当前工作目录或输出目录local_file_path:str# 原始输入文件路径file_title:str# 文件标题文件名去后缀pdf_path:str# PDF 文件路径 (如果输入是PDF)md_path:str# Markdown 文件路径 (转换后或直接输入的)# --- 内容数据 ---md_content:str# Markdown 的全文内容chunks:list# 切片后的文本列表包含 metadataitem_name:str# 识别出的主体名称 (如: 万用表)用于增强检索# --- 数据库相关 ---embeddings_content:list# 包含向量数据的列表准备写入 Milvus# 建议定一个初始化对象方便后续使用# 定义图状态的默认初始值graph_default_state:ImportGraphState{task_id:,is_pdf_read_enabled:False,is_md_read_enabled:False,local_dir:,local_file_path:,pdf_path:,md_path:,file_title:,md_content:,chunks:[],item_name:,embeddings_content:[]}defcreate_default_state(**overrides)-ImportGraphState: 创建默认状态支持覆盖 Args: **overrides: 要覆盖的字段关键字参数解包 Returns: 新的状态实例 Examples: state create_default_state(task_idtask_001, local_file_pathdoc.pdf) # 默认状态statecopy.deepcopy(graph_default_state)# 用 overrides 覆盖默认值state.update(overrides)# 返回创建好的状态字典实例returnstatedefget_default_state()-ImportGraphState: 返回一个新的状态实例避免全局变量污染 returncopy.deepcopy(graph_default_state)if__name____main__: 测试 # 创建默认状态statecreate_default_state(local_file_path万用表RS-12的使用.pdf)logger.info(state)2 定义主图为了确保整体流程设计的科学性与执行连贯性我们采用“Top-Down”自顶向下的开发模式以 “总指挥部” 的全局视角统筹推进具体实施步骤如下搭建节点骨架Stubs优先定义全流程所需的所有功能节点仅保留核心日志打印能力如节点进入 / 退出日志暂不实现内部复杂业务逻辑快速搭建起流程的 “骨架结构”串联主图Graph基于预设的业务流转规则编写主图逻辑将所有节点骨架按序串联明确节点间的输入输出关系、分支判断条件如文件格式分流逻辑形成完整的流程链路验证流程通畅性启动端到端测试验证节点间的调用链路是否通顺、数据流转是否符合预期、分支跳转是否准确确保流程无阻塞、无逻辑漏洞填充节点核心逻辑在流程链路验证通过后再逐一聚焦每个节点的内部实现完成复杂业务逻辑的开发如 PDF 结构化转换、向量编码、重排序算法等实现 “骨架” 到 “完整系统” 的落地。该模式的核心优势在于先保障 “流程走得通”再聚焦 “功能做得好”避免因局部逻辑复杂导致整体流程设计偏差大幅提升开发效率与流程稳定性。2.1 第一步创建节点骨架我们需要先创建以下 8 个文件每个文件里只写一个最简单的“空函数”确保主图能导入它们。请在app/import_process/agent/nodes目录下创建以下文件并将对应代码复制进去。(1) 入口节点:node_entry.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_entry)defnode_entry(state:ImportGraphState)-ImportGraphState: 节点: 入口节点 (node_entry) 为什么叫这个名字: 作为图的 Entry Point负责接收外部输入并决定流程走向。 未来要实现: 1. 接收文件路径。 2. 判断文件类型 (PDF/MD)。 3. 设置 state 中的路由标记 (is_pdf_read_enabled / is_md_read_enabled)。 # 模拟简单的路由逻辑防止报错 (仅 node_entry 需要)iflocal_file_pathinstate:pathstate[local_file_path]ifpath.endswith(.pdf):state[is_pdf_read_enabled]Trueelifpath.endswith(.md):state[is_md_read_enabled]Truereturnstate(2) PDF转换节点:node_pdf_to_md.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_pdf_to_md)defnode_pdf_to_md(state:ImportGraphState)-ImportGraphState: 节点: PDF转Markdown (node_pdf_to_md) 为什么叫这个名字: 核心任务是将 PDF 非结构化数据转换为 Markdown 结构化数据。 未来要实现: 1. 调用 MinerU (magic-pdf) 工具。 2. 将 PDF 转换成 Markdown 格式。 3. 将结果保存到 state[md_content]。 returnstate(3) 图片处理节点:node_md_img.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_md_img)defnode_md_img(state:ImportGraphState)-ImportGraphState: 节点: 图片处理 (node_md_img) 为什么叫这个名字: 处理 Markdown 中的图片资源 (Image)。 未来要实现: 1. 扫描 Markdown 中的图片链接。 2. 将图片上传到 MinIO 对象存储。 3. (可选) 调用多模态模型生成图片描述。 4. 替换 Markdown 中的图片链接为 MinIO URL。 returnstate(4) 文档切分节点:node_document_split.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_document_split)defnode_document_split(state:ImportGraphState)-ImportGraphState: 节点: 文档切分 (node_document_split) 为什么叫这个名字: 将长文档切分成小的 Chunks (切片) 以便检索。 未来要实现: 1. 基于 Markdown 标题层级进行递归切分。 2. 对过长的段落进行二次切分。 3. 生成包含 Metadata (标题路径) 的 Chunk 列表。 returnstate(5) 主体识别节点:node_item_name_recognition.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_item_name_recognition)defnode_item_name_recognition(state:ImportGraphState)-ImportGraphState: 节点: 主体识别 (node_item_name_recognition) 为什么叫这个名字: 识别文档核心描述的物品/商品名称 (Item Name)。 未来要实现: 1. 取文档前几段内容。 2. 调用 LLM 识别这篇文档讲的是什么东西 (如: Fluke 17B 万用表)。 3. 存入 state[item_name] 用于后续数据幂等性清理。 returnstate(6) 向量化节点:node_bge_embedding.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_bge_embedding)defnode_bge_embedding(state:ImportGraphState)-ImportGraphState: 节点: 向量化 (node_bge_embedding) 为什么叫这个名字: 使用 BGE-M3 模型将文本转换为向量 (Embedding)。 未来要实现: 1. 加载 BGE-M3 模型。 2. 对每个 Chunk 的文本进行 Dense (稠密) 和 Sparse (稀疏) 向量化。 3. 准备好写入 Milvus 的数据格式。 returnstate(7) 写入Milvus节点:node_import_milvus.pyimportsysfromapp.core.loggerimportlogger,node_logfromapp.import_process.agent.stateimportImportGraphStatenode_log(node_import_milvus)defnode_import_milvus(state:ImportGraphState)-ImportGraphState: 节点: 导入向量库 (node_import_milvus) 为什么叫这个名字: 将处理好的向量数据写入 Milvus 数据库。 未来要实现: 1. 连接 Milvus。 2. 根据 item_name 删除旧数据 (幂等性)。 3. 批量插入新的向量数据。 returnstate2.2 第二步编写主图代码节点思路总结PDFMarkdown其他格式STARTnode_entrynode_pdf_to_mdnode_md_imgnode_document_splitnode_item_name_recognitionnode_bge_embeddingnode_import_milvusEND文件:app/import_process/agent/main_graph.py# 加载环境变量从 .env 文件读取配置如Milvus地址、KG服务地址、BGE模型路径等fromdotenvimportload_dotenv# 导入LangGraph核心依赖StateGraph(状态图)、START/END(内置起始/结束节点常量)fromlanggraph.graphimportStateGraph,END,STARTfromapp.core.loggerimportlogger# 导入自定义状态类统一管理工作流全程的所有数据各节点共享/修改fromapp.import_process.agent.stateimportImportGraphState,create_default_state# 导入所有自定义业务节点每个节点对应知识库导入的一个具体步骤fromapp.import_process.agent.nodes.node_entryimportnode_entry# 入口节点初始化参数、校验输入fromapp.import_process.agent.nodes.node_pdf_to_mdimportnode_pdf_to_md# PDF转MD解析PDF文件为markdown格式fromapp.import_process.agent.nodes.node_md_imgimportnode_md_img# MD图片处理提取/下载markdown中的图片、修复图片路径fromapp.import_process.agent.nodes.node_document_splitimportnode_document_split# 文档分块将长文档切分为符合模型要求的小片段fromapp.import_process.agent.nodes.node_item_name_recognitionimportnode_item_name_recognition# 项目名识别从分块中提取核心项目名称业务定制化fromapp.import_process.agent.nodes.node_bge_embeddingimportnode_bge_embedding# BGE向量化将文本分块转换为向量表示适配Milvus向量库fromapp.import_process.agent.nodes.node_import_milvusimportnode_import_milvus# 导入Milvus将向量数据写入Milvus向量数据库# 初始化环境变量必须在配置读取前执行确保后续节点能获取到环境变量中的配置信息load_dotenv()# 1. 初始化LangGraph状态图 # 核心StateGraph是LangGraph的核心类用于构建有状态的工作流# 参数ImportGraphState自定义TypedDict类型定义了工作流的**全量状态字段**# 作用所有节点的入参都是该状态对象节点返回的键值对会自动合并回状态实现节点间数据共享workflowStateGraph(ImportGraphState)# 2. 注册所有业务节点 # 语法add_node(节点唯一标识, 节点函数)# 要求节点函数必须接收「状态对象」作为入参返回字典用于更新状态# 所有节点按「知识库导入流程」先后顺序注册节点标识与函数名保持一致便于维护workflow.add_node(node_entry,node_entry)# 流程入口参数初始化、输入校验workflow.add_node(node_pdf_to_md,node_pdf_to_md)# PDF转MD非MD格式文件的前置处理workflow.add_node(node_md_img,node_md_img)# MD图片处理保证文档中图片的可访问性workflow.add_node(node_document_split,node_document_split)# 文档分块解决大文本无法向量化/推理的问题workflow.add_node(node_item_name_recognition,node_item_name_recognition)# 项目名识别业务定制化步骤提取核心业务标识workflow.add_node(node_bge_embedding,node_bge_embedding)# BGE向量化文本→向量为Milvus存储做准备workflow.add_node(node_import_milvus,node_import_milvus)# 向量入库将向量数据持久化到Milvus# 3. 设置工作流入口节点 # 语法set_entry_point(节点标识) → 推荐写法直接指定流程起始节点# 等效写法workflow.add_edge(START, node_entry)START是LangGraph内置起始常量# 作用指定工作流执行的第一个节点替代手动添加START到目标节点的边代码更简洁workflow.set_entry_point(node_entry)# 4. 定义条件路由函数入口节点后的分支逻辑 # 核心根据状态中的配置项动态决定后续执行路径实现「PDF导入」/「MD直接导入」分支# 要求接收状态对象为入参返回「目标节点标识」或END内置结束常量defroute_after_entry(state:ImportGraphState)-str: 入口节点后的条件路由逻辑 :param state: 工作流全量状态对象包含所有配置项和中间结果 :return: 目标节点标识/ENDLangGraph会自动跳转到对应节点 # 分支1开启MD直接导入 → 跳过PDF转MD直接执行MD图片处理ifstate.get(is_md_read_enabled):returnnode_md_img# 分支2开启PDF导入 → 执行PDF转MD再走后续流程elifstate.get(is_pdf_read_enabled):returnnode_pdf_to_md# 分支3未开启任何导入配置 → 直接终止工作流END是LangGraph内置结束常量else:returnEND# 注册条件边将入口节点与路由函数绑定# 语法add_conditional_edges(源节点标识, 路由函数)# 作用源节点执行完成后调用路由函数根据返回值动态跳转到目标节点workflow.add_conditional_edges(node_entry,route_after_entry,{node_md_img:node_md_img,node_pdf_to_md:node_pdf_to_md,END:END})# 5. 注册静态顺序边分支合并后的统一流程 # 核心所有分支最终合并为「固定顺序执行流程」从MD图片处理到知识图谱入库一步到底# 语法add_edge(源节点标识, 目标节点标识/END) → 静态边固定路由关系无分支逻辑workflow.add_edge(node_pdf_to_md,node_md_img)# PDF转MD完成 → MD图片处理workflow.add_edge(node_md_img,node_document_split)# MD处理完成 → 文档分块workflow.add_edge(node_document_split,node_item_name_recognition)# 分块完成 → 项目名识别workflow.add_edge(node_item_name_recognition,node_bge_embedding)# 项目名识别完成 → BGE向量化workflow.add_edge(node_bge_embedding,node_import_milvus)# 向量化完成 → 导入Milvus向量库workflow.add_edge(node_import_milvus,END)# Milvus入库完成 → 工作流执行结束END是内置结束节点# 6. 编译工作流为可执行对象 # 语法compile() → 将StateGraph构建的流程编译为LangGraph的可执行应用# 作用生成可调用的kb_import_app通过invoke()方法触发工作流执行# 特性编译后可重复调用支持传入不同的初始状态实现多任务执行kb_import_appworkflow.compile()2.3 第三步验证图流程在实现具体业务逻辑前我们先跑一个测试脚本看看图能不能跑通路线对不对。创建测试文件:test/04-test_graph_flow.pyimportjsonfromapp.import_process.agent.main_graphimportkb_import_appfromapp.import_process.agent.stateimportcreate_default_stateimportsysfromapp.core.loggerimportlogger logger.info( 开始测试 )initial_statecreate_default_state(local_file_path万用表RS-12的使用.pdf)final_stateNone# 只输出更最终的状态值字典形式不包含节点名称、执行日志、元数据等额外信息foreventinkb_import_app.stream(initial_state):forkey,valueinevent.items():logger.info(f节点:{key})final_statevalue# 格式化输出最终状态logger.info(f最终状态:{json.dumps(final_state,indent4,ensure_asciiFalse)})logger.info(图结构:)# uv add grandalfkb_import_app.get_graph().print_ascii()logger.info( 测试结束 )预期效果:你应该能看到控制台依次打印出每个节点的 [Stub] 执行节点: ...日志。PDF 流程应包含node_entry-node_pdf_to_md-node_md_img- … -node_import_milvusMarkdown 流程应包含node_entry-node_md_img- … -node_import_milvus(跳过了 node_pdf_to_md)其他格式的文档只包含node_entry直接结束如果能看到这些日志说明我们的图结构搭建成功接下来就可以放心地去填充每个节点的具体代码了。