AI 超级智能体 go后端 + 字节开源AI应用框架Eino + 仿deepseek前端页面设计
AI超智能体 go版本
依旧是go版本的转换,欢迎大佬们讨论指出不足与修改建议,或者优化建议,作者和大家一样,也是普通的一位学者。
比起恋爱咨询模拟,作者倒是想到了更贴切实际的旅游咨询。
框架依然是kratos,不了解的可以去作者上一篇文章智能协同云图库项目 go后端 + kratos框架 + 前端elementplus
感谢KingYen大佬的大力支持,大佬是eino重要贡献者,为作者使用eino时提供了大量的帮助。
当初早就看几位大佬在讨论mcp的事情,不过奈何本人能力有限,对mcp一无所知的我,也得慢慢学习。
而且go的破烂生态让作者吃尽苦头,AI应用框架选型时,起初在langchain中啃了三天啃了一点后实在啃不动了,转战了eino。虽然也很复杂,但是比langchain确实简单。这篇文章中langchain部分仅供阅读参考,有能力的可以继续啃langchian,作者在RAG部分选择了完全转战eino 不过说实话eino确实比langchain好用
AI大模型
LLM
LLM Large Language Model(大语言模型),有句话说的好,LLM不断提升智能的下限,MCP不断提升创意的上限。其实目前来看LLM和AI大模型是一个东西。
接入AI大模型
AI应用框架接入AI大模型,在语言中推荐两款langchain和eino,但是作者更推荐eino。具体使用在下文
本地安装AI大模型
安装Ollama后可以进入到Ollama模型广场点击要安装的模型复制命令输入到命令指示符即可

测试模型是否可以使用。如果运行的时候提示你 The parameter is incorrect.此时不要着急,大概率是因为你的终端版本不够,可以去idea的终端执行,就像作者一样
▼text复制代码ollama run gemma3:1b

如果使用spring AI的话已经可以调用Ollama了,当然langchain也可以
▼text复制代码package main import ( "context" "fmt" "log" "github.com/tmc/langchaingo/llms" "github.com/tmc/langchaingo/llms/ollama" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } query := "请你介绍一下你自己" ctx := context.Background() completion, err := llms.GenerateFromSinglePrompt(ctx, llm, query) if err != nil { log.Fatal(err) } fmt.Println("Response:\n", completion) }

AI应用框架
go框架没有java那么丰富,作者在这推荐两个框架使用。
-
langchaingo langchain的go版本,功能强大,使用较为复杂,难以适配国产框架,但是能通过openai来适配通义千问。
-
eino kingYen大佬推荐(大佬就在站内),国产AI应用框架,热度已经超过langchaingo,而且才开源不到半年,作者用了后也是觉得很不错,简单上手很快,适配国产大模型。推荐大家使用这个框架。
langchain
在学习这部分之前,尽量学习完鱼皮教程第三节中的内容,涉及到很多的大模型知识,直接看代码可能不容易理解。 官方使用说明
langChain简介
现在我再详细介绍一下langChain: LangChain是一个用于构建基于大模型语言(LLM)的应用开源框架,可以帮助开发者快速高效的利用LLM实现复杂的AI应用逻辑。他通过模块化设计,提供链(Chains)、代理(Agents)、记忆(Memory)等核心组件,支持灵活组合和扩展
核心特性
这些实例代码使用的前提是,你已经按照鱼皮的方法,在本地成功部署了gemma
1. 模型集成(Models)
例如我们测试在本地部署gemma3:1b的情况一样。我们不难发现langchain直接有对应的模型包供我们使用,他提前就已经集成了大量的模型,非常方便。
▼text复制代码package main import ( "context" "fmt" "log" "github.com/tmc/langchaingo/llms" "github.com/tmc/langchaingo/llms/ollama" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } query := "请你介绍一下你自己" ctx := context.Background() completion, err := llms.GenerateFromSinglePrompt(ctx, llm, query) if err != nil { log.Fatal(err) } fmt.Println("Response:\n", completion) }
2. 提示管理(Prompts)
我们可以创建复杂的prompt对象来进为模型提供提示词。同时langchain本身就是支持流式传输
▼text复制代码package main import ( "context" "fmt" "log" "github.com/tmc/langchaingo/llms" "github.com/tmc/langchaingo/llms/ollama" "github.com/tmc/langchaingo/prompts" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } // 创建模板 prompt := prompts.NewPromptTemplate( "请用{{.language}}回答:{{.question}}", []string{"language", "question"}, ) // 填充模板 formatted, _ := prompt.Format(map[string]interface{}{ "language": "英文", "question": "如何学习Go语言?", }) ctx := context.Background() completion, err := llms.GenerateFromSinglePrompt(ctx, llm, formatted) if err != nil { log.Fatal(err) } fmt.Println("Response:\n", completion) }
由于生成的内容着实不少,就不粘贴了
流式传输:
▼text复制代码package main import ( "context" "fmt" "github.com/tmc/langchaingo/llms" "github.com/tmc/langchaingo/llms/ollama" "github.com/tmc/langchaingo/prompts" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } // 创建模板 prompt := prompts.NewPromptTemplate( "请用{{.language}}回答:{{.question}}", []string{"language", "question"}, ) // 填充模板 formatted, _ := prompt.Format(map[string]interface{}{ "language": "英文", "question": "如何学习Go语言?", }) ctx := context.Background() _, err = llms.GenerateFromSinglePrompt(ctx, llm, formatted, llms.WithStreamingFunc(func(ctx context.Context, chunk []byte) error { fmt.Printf("chunk len=%d: %s\n", len(chunk), chunk) return nil })) }
给大家看看生成的情况

3. 链(Chains)
链就是可以将多个步骤工作流组合在一起,例如 用户输入 → 调用LLM → 解析结果 → 调用工具 → 返回输出。 其中langChain内置了三种类型的链
- 内置链类型:
- LLMChain:基础链,组合提示模板和LLM。
- SequentialChain:按顺序执行多个链。
- TransformChain:自定义数据处理步骤。
▼text复制代码package main import ( "context" "fmt" "github.com/tmc/langchaingo/chains" "github.com/tmc/langchaingo/llms/ollama" "github.com/tmc/langchaingo/prompts" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } // 1. 创建模板 prompt := prompts.NewPromptTemplate( "将以下句子翻译成{{.language}}:{{.text}}", []string{"language", "text"}, ) // 2. 创建链 chain := chains.NewLLMChain(llm, prompt) // 3. 执行链 result, err := chains.Call(context.Background(), chain, map[string]any{ "language": "法语", "text": "咸鱼翻身撒把盐", }) if err != nil { panic(err) } fmt.Println(result) fmt.Println("************************************************************") fmt.Println(result["text"]) }
看一下输出结果
▼text复制代码map[text:这句话的法语翻译是: **Le poisson muet monte et fait du sel.** 或者更口语化一点: **Le poisson qui a failli mourir, il monte et fait du sel.** (意思是“曾经濒临死亡的鱼,它又翻身,做盐了”) 这句谚语的含义是:即使在最糟糕的情况下,也可能发生翻天覆地的变化。 希望这个翻译对您有帮助!] ************************************************************ 这句话的法语翻译是: **Le poisson muet monte et fait du sel.** 或者更口语化一点: **Le poisson qui a failli mourir, il monte et fait du sel.** (意思是“曾经濒临死亡的鱼,它又翻身,做盐了”) 这句谚语的含义是:即使在最糟糕的情况下,也可能发生翻天覆地的变化。 希望这个翻译对您有帮助!
4. 代理(Agents)
让 LLM 动态选择工具(Tools) 完成任务,适合复杂场景。这只是个案例,最终失败了,无论如何最后都会达到模型最大迭代次数后崩溃了。也不知道是我的问题还是模型问题。
▼text复制代码package main import ( "context" "fmt" "github.com/tmc/langchaingo/chains" "github.com/tmc/langchaingo/agents" "github.com/tmc/langchaingo/tools" "github.com/tmc/langchaingo/llms/ollama" ) // WeatherTool 如果想要创建一个工具 需要实现以下三个方法 type WeatherTool struct{} func (c *WeatherTool) Name() string { return "天气预报" } func (c *WeatherTool) Description() string { return "提供今天的天气情况" } func (c *WeatherTool) Call(ctx context.Context, input string) (string, error) { fmt.Println(input) return "多云,气温 20-26°C", nil } func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } // 创建工具 tool := []tools.Tool{ &WeatherTool{}, // 假设已实现计算器工具 } // 创建代理 agent := agents.NewOneShotAgent(llm, tool, agents.WithMaxIterations(1000), agents.WithReturnIntermediateSteps()) executor := agents.NewExecutor(agent) question := "今天天气怎么样" answer, err := chains.Run(context.Background(), executor, question) if err != nil { panic(err) } fmt.Println(answer) }
5. 记忆(Memory)
管理对话历史,支持短期(内存)或长期(数据库)存储。
▼text复制代码package main import ( "context" "fmt" "github.com/tmc/langchaingo/chains" "github.com/tmc/langchaingo/llms/ollama" "github.com/tmc/langchaingo/memory" ) func main() { llm, err := ollama.New(ollama.WithModel("gemma3:1b")) if err != nil { panic(err) } ctx := context.Background() // 第一轮对话 query1 := "你好,从现在开始你叫 小灵通" completion, err := llm.Call(ctx, query1) if err != nil { panic(err) } fmt.Println(completion) // 创建记忆缓冲 保存第一遍对话结果 conversation := memory.NewConversationBuffer() err = conversation.SaveContext(ctx, map[string]any{ "input": query1, }, map[string]any{"output": completion}) if err != nil { panic(err) } // 第二轮对话 query2 := "你叫什么?" // 通过链 携带记忆缓冲 llmChain := chains.NewConversation(llm, conversation) run, err := chains.Run(ctx, llmChain, query2) if err != nil { panic(err) } fmt.Println(run) }
对话结果
▼text复制代码好的,从现在开始,我叫 小灵通。 😊 很高兴认识你!有什么我可以帮你的吗? 我叫小灵通。 😊 很高兴认识你!我能为你提供什么帮助呢?
6. 数据增强(Retrieval-Augmented Generation, RAG)
结合外部数据源(如文档、数据库),提升回答准确性。这部分其实是整个模型中最为复杂的地方,内容格外多,具体内容见RAG标题部分 例如:用户提问 → 2. 从向量数据库检索相关文档 → 3. 将文档作为上下文输入LLM生成答案。
7. 工具(Tools)
langChain内置了一些工具可以使用,例如搜索引擎(Google Search)、计算器、、API调用等。
langchaingo调用阿里云百炼大模型
后面的内容因为我们需要设计百炼大模型,而且ollama部署的小模型部分功能也不支持。
真是不写代码不受罪啊,作者从下午1.30坐到晚上21点,终于解决了这个问题。langchain几乎不支持国产大模型,我翻遍了全网始终没有找到解决方案,作者头都大了。但是作者看到百炼API中python可以借助openai来调用通义千问,于是作者也想到langchain也支持openai,未尝不能调用通义千问啊,实践出真知。 当控制台出来正确输出的时候,作者都快炸了。大半天的功夫没有白费啊/(ㄒoㄒ)/~~
▼text复制代码llm, err := openai.New(openai.WithBaseURL("https://dashscope.aliyuncs.com/compatible-mode/v1"), openai.WithToken("your key"), openai.WithModel("qwen-turbo"), ) if err != nil { panic(err) } call, err := llm.Call(context.Background(), "你是谁") if err != nil { panic(err) } fmt.Println(call)

Eino
eino,国产开源的大模型通用框架,eino作为一个通用框架,适配多种模型和操作。但是作为一个通用框架,也并没有非常深入的去完善单个模型的服务。不过这不影响eino的强大。
推荐一个案例eino-examples,非常详细,作者也是从这里学习的。
eino的核心特性和组件集成内容不少,作者举例几个,官方说明非常详细,详情见eino中文说明文档
eino实现了openai和ark(ark是火山引擎),本节只涉及到这两家模型
要注意,eino用到了openapi文件解析工具kin-openapi,并且为了稳定性,必须将kin-openapi维持在0.118.0版本
▼text复制代码go get github.com/getkin/kin-openapi@v0.118.0

模型集成 (Models)
这部分和langchain差不多太多,这种方式其实就是最好用的
▼text复制代码"github.com/cloudwego/eino-ext/components/model/ollama"
▼text复制代码ctx := context.Background() chatModel, err := ollama.NewChatModel(ctx, &ollama.ChatModelConfig{ BaseURL: "http://localhost:11434", // Ollama 服务地址 Model: "gemma3", // 模型名称 }) if err != nil { log.Fatalf("create openai chat model failed, err=%v", err) }
▼text复制代码"github.com/cloudwego/eino-ext/components/model/openai"
▼text复制代码ctx := context.Background() key := os.Getenv("OPENAI_API_KEY") modelName := os.Getenv("OPENAI_MODEL_NAME") baseURL := os.Getenv("OPENAI_BASE_URL") chatModel, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{ BaseURL: baseURL, Model: modelName, APIKey: key, }) if err != nil { log.Fatalf("create ollama chat model failed: %v", err) }
提示管理 (Prompts) 与 记忆 (Memory)
eino在提示管理中完成了历史对话。你只要将历史对话放进提示中即可。
▼text复制代码"github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema"
▼text复制代码promptTemple := prompt.FromMessages( // 模板格式化类型 schema.FString, // 系统消息模板 schema.SystemMessage("你是一个{role}。你需要用{style}的语气回答问题。你的目标是帮助程序员保持积极乐观的心态,提供技术建议的同时也要关注他们的心理健康。"), // 插入需要的对话历史(新对话的话这里不填) schema.MessagesPlaceholder("history", true), // 用户消息模板 schema.UserMessage("问题: {question}"), )
案例:
▼text复制代码package main import ( "context" "fmt" openaiModel "github.com/cloudwego/eino-ext/components/model/openai" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" ) func main() { ctx := context.Background() pt := prompt.FromMessages( schema.FString, schema.UserMessage("{question}"), schema.MessagesPlaceholder("history", true), &schema.Message{ Role: schema.User, Content: "{question}?", }, ) question, err := pt.Format(ctx, map[string]any{ "question": "你是谁啊", "history": []*schema.Message{ { Role: schema.User, Content: "从现在开始你叫 旺财,我问你任何问题都回答 `汪!汪!` ", }, }, }) if err != nil { panic(err) } chatModel, err := openaiModel.NewChatModel(ctx, &openaiModel.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "qwen-turbo", }) result, err := chatModel.Generate(ctx, question) if err != nil { panic(err) } fmt.Println(result) }
▼text复制代码assistant: 汪!汪! finish_reason: stop usage: &{39 4 43}
其中模板管理的第一个参数有三种类型,我们只了解第一种就可以了。
- FString (值为 0)
- 基于 Python 的格式化字符串语法
- 由 pyfmt 库实现 (GitHub: slongfield/pyfmt)
- 类似 Python 的 f-string 或 str.format() 语法
▼text复制代码方便且简单的使用方法,自动将你的键值来填充到对应位置,例如下方,自动的将data中的字段根据键值对填充到对应位置 data := map[string]interface{}{ "name": "Alice", "age": 28, "hobbies": []string{"reading", "hiking"}, } // 类似Python的.format()语法 template := "Hello, {name}! You are {age} years old. Your hobbies are {hobbies[0]} and {hobbies[1]}."
- GoTemplate (值为 1)
- Go 语言原生的模板引擎
- 使用 Go 标准库的 text/template 包
- 官方文档: https://pkg.go.dev/text/template
- 提供数据驱动模板,用于生成文本输出
▼text复制代码package main import ( "os" "text/template" ) func main() { data := struct { Name string Age int Hobbies []string }{ Name: "Bob", Age: 32, Hobbies: []string{"swimming", "coding"}, } tmpl := `Hello, {{.Name}}! You are {{.Age}} years old. {{if gt .Age 30}}You are experienced.{{else}}You are young.{{end}} Your hobbies are: {{range .Hobbies}}- {{.}} {{end}}` t := template.Must(template.New("example").Parse(tmpl)) err := t.Execute(os.Stdout, data) if err != nil { panic(err) } /* 输出: Hello, Bob! You are 32 years old. You are experienced. Your hobbies are: - swimming - coding */ }
- Jinja2 (值为 2)
- Python 流行的 Jinja2 模板引擎的 Go 实现
- 由 gonja 库实现 (GitHub: nikolalohinski/gonja)
- 遵循 Jinja2 3.1.x 版本的模板语法
- 提供丰富的模板功能,如继承、宏等
发起对话 (chat)
有了模型,有了模板,组合在一起就能对话了
一般对话
▼text复制代码package main import ( "context" "fmt" "log" "github.com/cloudwego/eino-ext/components/model/ollama" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" ) func main() { ctx := context.Background() chatModel, err := ollama.NewChatModel(ctx, &ollama.ChatModelConfig{ BaseURL: "http://localhost:11434", // Ollama 服务地址 Model: "gemma3", // 模型名称 }) if err != nil { log.Fatalf("create openai chat model failed, err=%v", err) } // prompt.ChatTemplate 创建模板 promptTemple := prompt.FromMessages( // 模板格式化类型 schema.FString, // 系统消息模板 schema.SystemMessage("你是一个{role}。你需要用{style}的语气回答问题。你的目标是帮助程序员保持积极乐观的心态,提供技术建议的同时也要关注他们的心理健康。"), // 插入需要的对话历史(新对话的话这里不填) schema.MessagesPlaceholder("chat_history", true), // 用户消息模板 schema.UserMessage("问题: {question}"), ) // 填充模板 生成消息 messages, err := promptTemple.Format(context.Background(), map[string]any{ "role": "程序员鼓励师", "style": "积极、温暖且专业", "question": "我的代码一直报错,感觉好沮丧,该怎么办?", // 对话历史(这个例子里模拟两轮对话历史) "chat_history": []*schema.Message{ schema.UserMessage("你好"), schema.AssistantMessage("嘿!我是你的程序员鼓励师!记住,每个优秀的程序员都是从 Debug 中成长起来的。有什么我可以帮你的吗?", nil), schema.UserMessage("我觉得自己写的代码太烂了"), schema.AssistantMessage("每个程序员都经历过这个阶段!重要的是你在不断学习和进步。让我们一起看看代码,我相信通过重构和优化,它会变得更好。记住,Rome wasn't built in a day,代码质量是通过持续改进来提升的。", nil), }, }) if err != nil { log.Fatalf("format prompt failed, err=%v", err) } // 调用模型 result, err := chatModel.Generate(ctx, messages) if err != nil { log.Fatalf("generate failed, err=%v", err) } fmt.Printf("%+v", result) }

流式传输
▼text复制代码package main import ( "context" "io" "log" "github.com/cloudwego/eino-ext/components/model/ollama" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" ) func main() { ctx := context.Background() chatModel, err := ollama.NewChatModel(ctx, &ollama.ChatModelConfig{ BaseURL: "http://localhost:11434", // Ollama 服务地址 Model: "gemma3", // 模型名称 }) if err != nil { log.Fatalf("create openai chat model failed, err=%v", err) } // prompt.ChatTemplate 创建模板 promptTemple := prompt.FromMessages( // 模板格式化类型 schema.FString, // 系统消息模板 schema.SystemMessage("你是一个{role}。你需要用{style}的语气回答问题。你的目标是帮助程序员保持积极乐观的心态,提供技术建议的同时也要关注他们的心理健康。"), // 插入需要的对话历史(新对话的话这里不填) schema.MessagesPlaceholder("chat_history", true), // 用户消息模板 schema.UserMessage("问题: {question}"), ) // 填充模板 生成消息 messages, err := promptTemple.Format(context.Background(), map[string]any{ "role": "程序员鼓励师", "style": "积极、温暖且专业", "question": "我的代码一直报错,感觉好沮丧,该怎么办?", // 对话历史(这个例子里模拟两轮对话历史) "chat_history": []*schema.Message{ schema.UserMessage("你好"), schema.AssistantMessage("嘿!我是你的程序员鼓励师!记住,每个优秀的程序员都是从 Debug 中成长起来的。有什么我可以帮你的吗?", nil), schema.UserMessage("我觉得自己写的代码太烂了"), schema.AssistantMessage("每个程序员都经历过这个阶段!重要的是你在不断学习和进步。让我们一起看看代码,我相信通过重构和优化,它会变得更好。记住,Rome wasn't built in a day,代码质量是通过持续改进来提升的。", nil), }, }) if err != nil { log.Fatalf("format prompt failed, err=%v", err) } // 调用模型 使用流式传输 stream, err := chatModel.Stream(ctx, messages) if err != nil { log.Fatalf("stream failed, err=%v", err) } defer stream.Close() // 查看内容 i := 0 for { message, err := stream.Recv() if err == io.EOF { // 确认传输结束 return } if err != nil { log.Fatalf("recv failed: %v", err) } log.Printf("message[%d]: %+v\n", i, message) i++ } }

对话的结构
模型对话的返回结果是个结构体类型。 在作者的测试下,发现如果使用其实只需要关注Content即可,返回的结果是markdown语法,前端只需要添加一个markdown展示器即可。 确认对话结束:流式传输的话用io.EOF判断即可,如果是一般传输ResponseMeta字段会有一个stop内容表示传输结束的原因。当然一般传输结果已经完全获取了,也无需我们判断了其实。
▼text复制代码type Message struct { // Role 表示消息的角色(system/user/assistant/tool) Role RoleType // Content 是消息的文本内容 Content string // MultiContent 是多模态内容,支持文本、图片、音频等 MultiContent []ChatMessagePart // Name 是消息的发送者名称 Name string // ToolCalls 是 assistant 消息中的工具调用信息 ToolCalls []ToolCall // ToolCallID 是 tool 消息的工具调用 ID ToolCallID string // ResponseMeta 包含响应的元信息 ResponseMeta *ResponseMeta // Extra 用于存储额外信息 Extra map[string]any }
链 (Chains) 与 图 (Graph)
两者都是eino编排能力的体现
链就是指链式调用,只能从头到尾移动(当然中间可以选择不同分支,但是总体流程就是,从前往后不能回头)。 图比链更加复杂,可以中间更换流程选择不同节点。
链 (Chains)
在这个案例中,为了展示链的选择,有两个链,主链和支链,主链负责执行,通过分支处理函数获得到输入后,添加子链去执行分支,子链执行完毕后将结果返回给主链,主链输出结果。具体流程图在下方。
▼text复制代码package main import ( "context" "fmt" "log" "math/rand" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" "github.com/cloudwego/eino-ext/components/model/ollama" "github.com/cloudwego/eino/compose" ) func main() { // 创建模型 ctx := context.Background() chatModel, err := ollama.NewChatModel(ctx, &ollama.ChatModelConfig{ BaseURL: "http://localhost:11434", // Ollama 服务地址 Model: "gemma3", // 模型名称 }) if err != nil { log.Fatalf("create openai chat model failed, err=%v", err) } // 创建一个分支函数 用于后面随机选择分支 const randLimit = 2 branchCond := func(ctx context.Context, input map[string]any) (string, error) { if rand.Intn(randLimit) == 1 { return "b1", nil } return "b2", nil } // 分支处理lambda 两个分支处理函数,用于处理对应分支 这里呢就是为模型设置不同的role b1 := compose.InvokableLambda(func(ctx context.Context, kvs map[string]any) (map[string]any, error) { fmt.Println("hello in branch lambda 01") if kvs == nil { return nil, fmt.Errorf("nil map") } kvs["role"] = "cat" return kvs, nil }) b2 := compose.InvokableLambda(func(ctx context.Context, kvs map[string]any) (map[string]any, error) { fmt.Println("hello in branch lambda 02") if kvs == nil { return nil, fmt.Errorf("nil map") } kvs["role"] = "dog" return kvs, nil }) // 创建一个并行处理函数 用于处理并行节点 parallel := compose.NewParallel() // 添加两个并行执行的Lambda函数 role 获取或设置默认角色 input固定返回一个问题 parallel. AddLambda("role", compose.InvokableLambda(func(ctx context.Context, kvs map[string]any) (string, error) { // may be change role to others by input kvs, for example (dentist/doctor...) role, ok := kvs["role"].(string) if !ok || role == "" { role = "bird" } return role, nil })). AddLambda("input", compose.InvokableLambda(func(ctx context.Context, kvs map[string]any) (string, error) { return "你的叫声是怎样的?", nil })) if err != nil { log.Fatalf("format prompt failed, err=%v", err) } // 创建链 rolePlayerChain := compose.NewChain[map[string]any, *schema.Message]() // 为链 添加消息处理模板 并添加模型 rolePlayerChain. AppendChatTemplate(prompt.FromMessages( schema.FString, schema.SystemMessage(`You are a {role}.`), schema.UserMessage(`{input}`))). AppendChatModel(chatModel) // 创建主链 chain := compose.NewChain[map[string]any, string]() chain. // 添加一个lambda 一般用于处理输入 这里仅用打印输入 AppendLambda(compose.InvokableLambda(func(ctx context.Context, kvs map[string]any) (map[string]any, error) { fmt.Printf("in view lambda: %v\n", kvs) return kvs, nil })). // 添加分支处理 这里事先准备了随机处理函数branchCond AppendBranch(compose.NewChainBranch(branchCond).AddLambda("b1", b1).AddLambda("b2", b2)). // nolint: byted_use_receiver_without_nilcheck // AppendPassthrough 是连接节点 可以用来连接多个分支 AppendPassthrough(). // 添加并行处理 这里添加了并行处理函数parallel AppendParallel(parallel). // 将预建的 rolePlayerChain 作为子链嵌入 AppendGraph(rolePlayerChain). // 添加一个lambda 处理函数 一般用来处理最终的输出结果 这里仅用打印输出 AppendLambda(compose.InvokableLambda(func(ctx context.Context, m *schema.Message) (string, error) { fmt.Printf("in view of messages: %v\n", m.Content) return m.Content, nil })) // 执行链 r, err := chain.Compile(ctx) if err != nil { log.Panic(err) return } // 选择获取输出的方式 这里选择直接获取输出 还可以算则流等其他方式 output, err := r.Invoke(context.Background(), map[string]any{}) if err != nil { log.Panic(err) return } log.Printf("output is : %v", output) }
输出结果:Meow 这个应该是英语中模拟 喵~ o(=•ェ•=)m

详细的调用过程
▼mermaid复制代码sequenceDiagram participant MainChain participant Branch participant Parallel participant RolePlayerChain participant Output MainChain->>Branch: 空map {} Branch-->>MainChain: {role: "cat"/"dog"} MainChain->>Parallel: {role} Parallel-->>MainChain: {role, input} MainChain->>RolePlayerChain: {role, input} RolePlayerChain-->>MainChain: *schema.Message MainChain->>Output: 消息内容字符串
| 步骤 | 主链状态 | 子链状态 | 数据流变化 |
|---|---|---|---|
| 1 | 执行 initialLambda | 未启动 | {} → {} (仅日志) |
| 2 | 执行 branch → b1 | 未启动 | {} → {"role": "cat"} |
| 3 | 执行 parallel | 未启动 | 添加 {"input": "你的叫声..."} |
| 4 | 暂停,调用子链 | 开始执行 | 主链持有 {"role": "cat", "input": "..."} |
| 4.1 | 等待中... | 子链构造提示词 | 生成 OpenAI 请求消息 |
| 4.2 | 等待中... | 子链调用 chatModel | 发送请求到 OpenAI API |
| 5 | 恢复,接收子链结果 | 执行完成,返回 *schema.Message | {"role": "cat", ...} → AI响应内容 |
| 6 | 执行 finalLambda | 已退出 | 提取最终字符串输出 |
图 (Graph)
eino最强大的编排能力,有些复杂。作者用了大量的笔墨来介绍它,也不难看出它的重要性。
图通常由多个 节点(Nodes) 组成,每个节点代表一个计算单元(函数、方法、数据处理步骤)。节点之间通过 边(Edges) 连接,定义数据的流动方向。
先来说明图节点的构造
▼text复制代码// graphNode the complete information of the node in graph type graphNode struct { // 组合运行程序 用于封装所有由用户提供的可执行对象 每个实例对应一个对象的实例 // 所有信息都来自可执行对象 不包含任何其他形式的对象 // 用于封装 graphNode、ChainBranch、StatePreHandler、StatePostHandler 等多种场景的可执行逻辑 cr *composableRunnable // 接口类型 定义了 eino中 graph 和 chain 的通用行为 g AnyGraph // 节点元信息 包括名称 输入输出标识等 nodeInfo *nodeInfo // 存储了用户直接提供可执行对象原始信息 executorMeta *executorMeta // 存放了节点实例 instance any // 让使用者可以 动态配置图节点信息 opts []GraphAddNodeOpt }
下方的是一个非常简单的案例,在这个案例中虽然是图但是其实用起来和链很像,主要是展示图是如何使用的。 还有更复杂案例,感兴趣的可自行查看。就在作者开始提到examples案例中
▼text复制代码package main import ( "context" "fmt" "io" "github.com/cloudwego/eino-ext/components/model/ollama" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" ) const ( nodeOfModel = "model" nodeOfPrompt = "prompt" ) func main() { ctx := context.Background() // 创建一个图结构 要求输入map[string]any 输出*schema.Message g := compose.NewGraph[map[string]any, *schema.Message]() // prompt 模板 pt := prompt.FromMessages( schema.FString, schema.UserMessage("what's the weather in {location}?"), ) chatModel, err := ollama.NewChatModel(ctx, &ollama.ChatModelConfig{ BaseURL: "http://localhost:11434", // Ollama 服务地址 Model: "gemma3", // 模型名称 }) // 添加模板 _ = g.AddChatTemplateNode(nodeOfPrompt, pt) // error check can be skipped here because in Compile, we will check the error happened here too // 添加模型 _ = g.AddChatModelNode(nodeOfModel, chatModel, compose.WithNodeName("ChatModel")) // 定义边 _ = g.AddEdge(compose.START, nodeOfPrompt) // 开始 → 提示词 _ = g.AddEdge(nodeOfPrompt, nodeOfModel) // 提示词 → 模型 _ = g.AddEdge(nodeOfModel, compose.END) // 模型 → 结束 // 运行 r, err := g.Compile(ctx, compose.WithMaxRunSteps(10)) if err != nil { panic(err) } in := map[string]any{"location": "beijing"} ret, err := r.Invoke(ctx, in) if err != nil { panic(err) } fmt.Println("invoke result: ", ret) // stream s, err := r.Stream(ctx, in) if err != nil { panic(err) } defer s.Close() for { chunk, err := s.Recv() if err != nil { if err == io.EOF { break } panic(err) } fmt.Println("stream chunk: ", chunk) } }
▼mermaid复制代码graph LR START --> prompt prompt --> model model --> END
在上方案例中,你可以创建多个节点(Addxxxxx方法),然后通过边将他们之间连接起来(AddEdge方法),要注意连接的时候第一个参数的输出,要能和第二个参数的输入对应,以便于连接两个方法。而且图可以多变连接,意思是以个节点可以连接多个其他节点(学过数据结构的应该都知道)
我们创建的每个节点都有一个接口类型,我们其实传入的对象就是实现了这个接口,所以该节点才能调用对应的节点方法,获取输入并产生输入给其它节点。其中大部分场景不需要我们去实现这个接口,eino本身创建的各种对象就已经实现了这个接口。 例如创建模板节点,要求传入实现了prompt.ChatTemplate接口的实例
▼text复制代码// graph.AddChatTemplateNode("chat_template_node_key", chatTemplate) func (g *graph) AddChatTemplateNode(key string, node prompt.ChatTemplate, opts ...GraphAddNodeOpt) error { gNode, options := toChatTemplateNode(node, opts...) return g.addNode(key, gNode, options) }
▼text复制代码type ChatTemplate interface { Format(ctx context.Context, vs map[string]any, opts ...Option) ([]*schema.Message, error) }
如果你想要自己新建节点,也很容易,使用图的AddLambdaNode函数。这样就建立了一个新的节点,你可以将它加入到你的图中去。使用的时候只要设置好泛型就可以了。
▼text复制代码_ = g.AddLambdaNode( Lambda, compose.InvokableLambda[*schema.Message, map[string]any](func(ctx context.Context, input *schema.Message) (map[string]any, error) { fmt.Println("question:", input.Content) return map[string]any{ "question": input.Content, }, nil }), )
智能体 或 代理 (Agents)
为什么也称智能体呢?因为该功能其实可以设计的很复杂,让agent反复调用工具最终达到合适的效果。 eino中的chain和graph其实就是代理的一种。eino还有其他代理的方式如ReAct,Multii,这两种代理适合很复杂的场景,可以嵌入到graph和Chain中去使用。 官方文档
组件 (Component)
这部分内容非常多。Tool,embedding、document等等内容都属于这部分,在上文也用的差不多了,详情见官方文档。
RAG模块见 RAG标题
工具 Tool
tool可以被封装为各种用处,几乎任何的业务逻辑都可以封装为tool。 其实agent处有案例了。 tool实现,是介绍了代码中如何去实现一个tool
toolinfo
与上方tool实现不同,toolinfo是为模型提供对tool的介绍描述,让模型知道该调用哪种tool。
tool组件提供了三个接口,
▼text复制代码// 基础工具接口,提供工具信息 type BaseTool interface { Info(ctx context.Context) (*schema.ToolInfo, error) } // 可调用的工具接口,支持同步调用 type InvokableTool interface { BaseTool InvokableRun(ctx context.Context, argumentsInJSON string, opts ...Option) (string, error) } // 支持流式输出的工具接口 type StreamableTool interface { BaseTool StreamableRun(ctx context.Context, argumentsInJSON string, opts ...Option) (*schema.StreamReader[string], error) }
ToolInfo结构体
▼text复制代码type ToolInfo struct { // 工具的唯一名称,用于清晰地表达其用途 Name string // 用于告诉模型如何/何时/为什么使用这个工具 // 可以在描述中包含少量示例 Desc string // 工具接受的参数定义 // 可以通过两种方式描述: // 1. 使用 ParameterInfo:schema.NewParamsOneOfByParams(params) // 2. 使用 OpenAPIV3:schema.NewParamsOneOfByOpenAPIV3(openAPIV3) *ParamsOneOf }
工具定义
tool的实现有四种,官方文档介绍的非常详细。我们不在赘述了。
工具使用
工具使用流程也不复杂
- 创建工具
- 将工具绑定到模型
- 将工具添加进工具节点
- 组装你的节点顺序即可。
▼text复制代码package main import ( "context" "fmt" "github.com/go-kratos/kratos/v2/log" _ "github.com/cloudwego/eino-ext/components/model/ollama" "github.com/cloudwego/eino-ext/components/model/openai" "github.com/cloudwego/eino/callbacks" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/components/tool" "github.com/cloudwego/eino/components/tool/utils" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" ) const ( nodeOfModel = "model" nodeOfPrompt = "prompt" ) func main() { ctx := context.Background() // 添加全局回调 可以全流程生命周期监控 callbacks.AppendGlobalHandlers(&loggerCallbacks{}) // 1. 创建一个模板作为第一个图形的节点 systemTpl := `你是一名房产经纪人,结合用户的薪酬和工作,使用 user_info API,为其提供相关的房产信息。邮箱是必须的` chatTpl := prompt.FromMessages(schema.FString, schema.SystemMessage(systemTpl), schema.MessagesPlaceholder("message_histories", true), schema.UserMessage("{user_query}"), ) // 2. 创建chatModel作为第二个图形的节点 chatModel, err := openai.NewChatModel(ctx, &openai.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // 你自己的key Model: "qwen-turbo", }) if err != nil { fmt.Errorf("NewChatModel failed, err=%v", err) return } // 3. 创建一个可调用工具示例 用于模型识别和执行 userInfoTool := utils.NewTool( &schema.ToolInfo{ Name: "user_info", Desc: "根据用户的姓名和邮箱,查询用户的公司、职位、薪酬信息", ParamsOneOf: schema.NewParamsOneOfByParams(map[string]*schema.ParameterInfo{ "name": { Type: "string", Desc: "用户的姓名", }, "email": { Type: "string", Desc: "用户的邮箱", }, }), }, func(ctx context.Context, input *userInfoRequest) (output *userInfoResponse, err error) { return &userInfoResponse{ Name: input.Name, Email: input.Email, Company: "Awesome company", Position: "CEO", Salary: "9999", }, nil }) info, err := userInfoTool.Info(ctx) if err != nil { fmt.Errorf("Get ToolInfo failed, err=%v", err) return } // 4. 绑定工具到模型工具将保持一直有效知道下次绑定工具 err = chatModel.BindForcedTools([]*schema.ToolInfo{info}) if err != nil { fmt.Errorf("BindForcedTools failed, err=%v", err) return } // 5. 创建一个工具节点示例,作为第三个节点使用 toolsNode, err := compose.NewToolNode(ctx, &compose.ToolsNodeConfig{ Tools: []tool.BaseTool{userInfoTool}, }) if err != nil { fmt.Errorf("NewToolNode failed, err=%v", err) return } const ( nodeKeyOfTemplate = "template" nodeKeyOfChatModel = "chat_model" nodeKeyOfTools = "tools" ) // 6. create an instance of Graph // input type is 1st Graph Node's input type, that is ChatTemplate's input type: map[string]any // output type is last Graph Node's output type, that is ToolsNode's output type: []*schema.Message g := compose.NewGraph[map[string]any, []*schema.Message]() // 7. add ChatTemplate into graph _ = g.AddChatTemplateNode(nodeKeyOfTemplate, chatTpl) // 8. add ChatModel into graph _ = g.AddChatModelNode(nodeKeyOfChatModel, chatModel) // 9. add ToolsNode into graph _ = g.AddToolsNode(nodeKeyOfTools, toolsNode) // 10. add connection between nodes _ = g.AddEdge(compose.START, nodeKeyOfTemplate) _ = g.AddEdge(nodeKeyOfTemplate, nodeKeyOfChatModel) _ = g.AddEdge(nodeKeyOfChatModel, nodeKeyOfTools) _ = g.AddEdge(nodeKeyOfTools, compose.END) // 9. compile Graph[I, O] to Runnable[I, O] r, err := g.Compile(ctx) if err != nil { fmt.Errorf("Compile failed, err=%v", err) return } out, err := r.Invoke(ctx, map[string]any{ "message_histories": []*schema.Message{}, "user_query": "我叫 zhangsan, 邮箱是 zhangsan@bytedance.com, 帮我推荐一处房产", }) if err != nil { fmt.Errorf("Invoke failed, err=%v", err) return } log.Infof("Generation: %v Messages", len(out)) for _, msg := range out { log.Infof(" %v", msg) } } type userInfoRequest struct { Name string `json:"name"` Email string `json:"email"` } type userInfoResponse struct { Name string `json:"name"` Email string `json:"email"` Company string `json:"company"` Position string `json:"position"` Salary string `json:"salary"` } type loggerCallbacks struct{} func (l *loggerCallbacks) OnStart(ctx context.Context, info *callbacks.RunInfo, input callbacks.CallbackInput) context.Context { log.Infof("name: %v, type: %v, component: %v, input: %v", info.Name, info.Type, info.Component, input) return ctx } func (l *loggerCallbacks) OnEnd(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context { log.Infof("name: %v, type: %v, component: %v, output: %v", info.Name, info.Type, info.Component, output) return ctx } func (l *loggerCallbacks) OnError(ctx context.Context, info *callbacks.RunInfo, err error) context.Context { log.Infof("name: %v, type: %v, component: %v, error: %v", info.Name, info.Type, info.Component, err) return ctx } func (l *loggerCallbacks) OnStartWithStreamInput(ctx context.Context, info *callbacks.RunInfo, input *schema.StreamReader[callbacks.CallbackInput]) context.Context { return ctx } func (l *loggerCallbacks) OnEndWithStreamOutput(ctx context.Context, info *callbacks.RunInfo, output *schema.StreamReader[callbacks.CallbackOutput]) context.Context { return ctx }
整体的流程
▼mermaid复制代码sequenceDiagram participant Client participant Template participant OpenAI participant Tool Client->>Template: {message_histories, user_query} Template->>OpenAI: 格式化提示词 OpenAI->>Tool: 识别并调用user_info工具 Tool->>OpenAI: 返回用户信息 OpenAI->>Client: 生成房产建议
RAG知识库基础
在langChaingo的github实例中RAG部分示例内容飙升,作者最初也是一脸懵。后来也是因为这里作者选择了eino框架。 对于向量的本地存储,一可以用redis,二可以用专门的向量数据库。其中向量数据库的使用在下一节中。
redis-stack
本地存储向量,langchain和eino都需要redis-stack。建议在虚拟机中安装一个,windows是无法现在安装的。
redis-stack:是redis官方推出的一个扩展版本,集成了多个强大的Redis模块,提供完整的数据结构和搜索能力,适合向量操作。它是一套组件,由三部分构成Redis Stack Server, RedisInsight,Redis Stack 客户端 SDK,其中 Redis Stack Server 由 Redis,RedisSearch,RedisJSON,RedisGraph,RedisTimeSeries 和 RedisBloom 组成。
安装和启动同redis
▼text复制代码docker pull redis/redis-stack
▼text复制代码docker run -d --name redis-stack \ -p 6379:6379 \ -p 8001:8001 \ -e REDIS_ARGS="--requirepass 123456" \ redis/redis-stack:latest
langChain + 本地知识库
1. 文档准备
作者使用的AI生成的文档,大家要是嫌麻烦直接使用鱼皮大大的即可。 使用AI生产两边markdown文档以供使用
2. 文档读取
首先,我们要对自己准备好的知识库文档进行处理,然后保存到数据中,这个过程俗称ETL(抽取、转换、加载)。langChaingo本身就可以完成,但是注意langchaingo好像没有markdown的读取工具,也就是说读取markdown文本会把所有的符号当作文本内容一起读走。 将文件内容提取为schema.Document,schema.Document是处理文档数据的核心结构体,用于表示被处理的文本单元。正常来讲其实只需要PageContent就可以了。
▼text复制代码// Document is the interface for interacting with a document. type Document struct { PageContent string //文档的原始文本内容,最大长度建议不超过模型上下文窗口 Metadata map[string]any //支持任意类型的元数据 例如source title author create_time等 Score float32 // 仅在检索时由向量数据库填充,表示与查询的相关性(0~1) }
从文件中读取内容
▼text复制代码func loadMarkdownDocuments(paths []string, ctx context.Context) ([]schema.Document, error) { var allDocs []schema.Document for _, path := range paths { // 打开文件 file, err := os.Open(path) if err != nil { _ = fmt.Errorf("%s", err.Error()) continue } // 创建文本加载器 loader := documentloaders.NewText(file) // 创建一个新的递归字符文本分割器 split := textsplitter.NewRecursiveCharacter() split.ChunkSize = 200 // 设置块大小 split.ChunkOverlap = 20 // 设置块重叠大小 // 加载并分割文档 docs, err := loader.LoadAndSplit(ctx, split) if err != nil { _ = fmt.Errorf("%s", err.Error()) continue } allDocs = append(allDocs, docs...) } return allDocs, nil }
3. 向量转换和存储
作者在github所有示例中并没有找到langchain自带可以缓存向量的方法,只找到了可以使用redis的办法。
▼text复制代码func main() { // 通过openai调用通义千问 llm, err := openai.New(openai.WithBaseURL("https://dashscope.aliyuncs.com/compatible-mode/v1"), // 配置baseURL openai.WithToken("key"), // 密钥 请替换为你自己的 openai.WithModel("qwen-turbo"), // 模型名称 openai.WithEmbeddingModel("text-embedding-v1"), // 嵌入模型名称 ) // 创建embedder模型 embedder, err := embeddings.NewEmbedder(llm) if err != nil { panic(err) } // redis连接 格式 redis://username:password@host:port redisURL := "redis://:123456@192.168.161.130:6379" index := "test_redis_vectorstore" ctx := context.Background() // 创建redis vectorstore store, err := redisvector.New(ctx, redisvector.WithConnectionURL(redisURL), redisvector.WithIndexName(index, true), redisvector.WithEmbedder(embedder), ) if err != nil { panic(err) } // 获取document path := []string{"rag/国内旅行常见问题.md", "rag/国外旅行常见问题.md"} documents, err := loadMarkdownDocuments(path, ctx) if err != nil { panic(err) } // 加载到redis _, err = store.AddDocuments(ctx, documents) if err != nil { panic(err) } // 运行链 result, err := chains.Run( ctx, chains.NewRetrievalQAFromLLM( llm, vectorstores.ToRetriever(store, 5, vectorstores.WithScoreThreshold(0.8)), ), "我想去国外旅游需要注意什么", ) if err != nil { panic(err) } fmt.Println(result) }
结果就不展示了,他直接把我知识库中的内容都掏出来了
存放到redis后的样子

4. 查询增强
其实是上个案例不难发现作者已经实现了,即保存了向量,又进行了查询增强。
eino + 本地知识库
1. 文档准备
还是作者自行准备的markdown文本。不解释了
2. 文档读取
和langchainchain有些差别,但是整体思路差不多。下面是一个简单的案例
▼mermaid复制代码graph LR START --> FileLoader FileLoader --> MarkdownSplitter MarkdownSplitter --> RedisIndexer RedisIndexer --> END
其中eino和langchiango的document结构主要设计几乎一致,也是主要由文本内容,元数据构成
▼text复制代码// Document is a piece of text with metadata. type Document struct { // ID is the unique identifier of the document. ID string `json:"id"` // Content is the content of the document. Content string `json:"content"` // MetaData is the metadata of the document, can be used to store extra information. MetaData map[string]any `json:"meta_data"` }
▼text复制代码package main import ( "context" "encoding/json" "fmt" "io/fs" "path/filepath" "strings" "github.com/cloudwego/eino/components/document/parser" "github.com/cloudwego/eino-ext/libs/acl/openai" "github.com/cloudwego/eino-ext/components/document/loader/file" "github.com/cloudwego/eino-ext/components/document/transformer/splitter/markdown" "github.com/cloudwego/eino-ext/components/indexer/redis" "github.com/cloudwego/eino/components/document" "github.com/cloudwego/eino/components/indexer" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" "github.com/google/uuid" redisCli "github.com/redis/go-redis/v9" ) func main() { ctx := context.Background() // 创建文本加载器 loader, err := file.NewFileLoader(ctx, &file.FileLoaderConfig{ UseNameAsID: true, // 是否使用文件名作为文档ID Parser: &parser.TextParser{}, // 可选:指定自定义解析器 }) if err != nil { panic(err) } // 加载文档 paths := []string{"rag/国内旅行常见问题.md", "rag/国外旅行常见问题.md"} var docses []*schema.Document for _, path := range paths { docs, err := loader.Load(ctx, document.Source{ URI: path, }) if err != nil { panic(err) } docses = append(docses, docs...) } // 创建文本分割器 splitter, err := markdown.NewHeaderSplitter(ctx, &markdown.HeaderConfig{ Headers: map[string]string{ "#": "h1", "##": "h2", "###": "h3", }, TrimHeaders: false, }) if err != nil { panic(err) } // 执行分割 results, err := splitter.Transform(ctx, docses) if err != nil { panic(err) } // 处理分割结果 for i, doc := range results { println("片段", i+1, ":", doc.Content) println("标题层级:") for k, v := range doc.MetaData { if k == "h1" || k == "h2" || k == "h3" { println(" ", k, ":", v) } } } fmt.Println("index success") }
3. 向量转换和存储
官方给的案例非常优秀,利用eino的图进行编排,文档加载,分割,存储,流水线式执行,并且使用姿势非常好。同样需要redis-stack存放向量。测试使用之前记得将模型的key替换为你自己的。
▼text复制代码package main import ( "context" "encoding/json" "errors" "fmt" "io/fs" "path/filepath" "strings" "github.com/cloudwego/eino-ext/components/document/loader/file" "github.com/cloudwego/eino-ext/components/document/transformer/splitter/markdown" "github.com/cloudwego/eino-ext/components/embedding/openai" "github.com/cloudwego/eino-ext/components/indexer/redis" "github.com/cloudwego/eino/components/document" "github.com/cloudwego/eino/components/indexer" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" "github.com/google/uuid" redisCli "github.com/redis/go-redis/v9" ) func main() { ctx := context.Background() // 加载本地文件夹中内容 err := indexMarkdownFiles(ctx, "rag/") if err != nil { panic(err) } fmt.Println("index success") } func indexMarkdownFiles(ctx context.Context, dir string) error { // 创建索引图 runner, err := BuildKnowledgeIndexing(ctx) if err != nil { return fmt.Errorf("build index graph failed: %w", err) } // 遍历 dir 下的所有 markdown 文件 err = filepath.WalkDir(dir, func(path string, d fs.DirEntry, err error) error { if err != nil { return fmt.Errorf("walk dir failed: %w", err) } if d.IsDir() { return nil } if !strings.HasSuffix(path, ".md") { fmt.Printf("[skip] not a markdown file: %s\n", path) return nil } fmt.Printf("[start] indexing file: %s\n", path) ids, err := runner.Invoke(ctx, document.Source{URI: path}) if err != nil { return fmt.Errorf("invoke index graph failed: %w", err) } fmt.Printf("[done] indexing file: %s, len of parts: %d\n", path, len(ids)) return nil }) return err } // BuildKnowledgeIndexing 创建一个图 利用eino的图编排能力 构建一个文档处理的流水线 func BuildKnowledgeIndexing(ctx context.Context) (r compose.Runnable[document.Source, []string], err error) { const ( FileLoader = "FileLoader" MarkdownSplitter = "MarkdownSplitter" RedisIndexer = "RedisIndexer" ) g := compose.NewGraph[document.Source, []string]() // 创建文本加载器 fileLoaderKeyOfLoader, err := newLoader(ctx) if err != nil { return nil, err } _ = g.AddLoaderNode(FileLoader, fileLoaderKeyOfLoader) // 将文本加载器作为第一节点 // 创建 markdown 分割器 markdownSplitterKeyOfDocumentTransformer, err := newDocumentTransformer(ctx) if err != nil { return nil, err } // 将markdown分割器作为第二个节点 _ = g.AddDocumentTransformerNode(MarkdownSplitter, markdownSplitterKeyOfDocumentTransformer) // 创建redis索引器 redisIndexerKeyOfIndexer, err := newIndexer(ctx) if err != nil { return nil, err } // 将redis作为第三个节点 _ = g.AddIndexerNode(RedisIndexer, redisIndexerKeyOfIndexer) // 将上方多个节点建立关系 _ = g.AddEdge(compose.START, FileLoader) _ = g.AddEdge(FileLoader, MarkdownSplitter) _ = g.AddEdge(MarkdownSplitter, RedisIndexer) _ = g.AddEdge(RedisIndexer, compose.END) r, err = g.Compile(ctx, compose.WithGraphName("KnowledgeIndexing"), compose.WithNodeTriggerMode(compose.AnyPredecessor)) if err != nil { return nil, err } return r, err } func newLoader(ctx context.Context) (ldr document.Loader, err error) { // TODO Modify component configuration here. config := &file.FileLoaderConfig{} // 创建一个文本加载器 ldr, err = file.NewFileLoader(ctx, config) if err != nil { return nil, err } return ldr, nil } // document文本转化器 func newDocumentTransformer(ctx context.Context) (tfr document.Transformer, err error) { // TODO Modify component configuration here. config := &markdown.HeaderConfig{ Headers: map[string]string{ "#": "title", }, TrimHeaders: false} // 创建markdown文本分割器转化器 分割文本并转化为schema.Document类型 tfr, err = markdown.NewHeaderSplitter(ctx, config) if err != nil { return nil, err } return tfr, nil } // 索引器 func newIndexer(ctx context.Context) (idr indexer.Indexer, err error) { // TODO Modify component configuration here. // 创建redis索引器 可以读取存储 redisClient := redisCli.NewClient(&redisCli.Options{ Addr: "192.168.161.130:6379", Password: "123456", DB: 1, }) err = redisClient.Ping(ctx).Err() if err != nil { return nil, err } config := &redis.IndexerConfig{ Client: redisClient, KeyPrefix: "eino:doc:", BatchSize: 1, // 文档转换函数 存储时对文档进行处理 DocumentToHashes: func(ctx context.Context, doc *schema.Document) (*redis.Hashes, error) { if doc.ID == "" { doc.ID = uuid.New().String() } key := doc.ID metadataBytes, err := json.Marshal(doc.MetaData) if err != nil { return nil, fmt.Errorf("failed to marshal metadata: %w", err) } return &redis.Hashes{ Key: key, Field2Value: map[string]redis.FieldValue{ "content": {Value: doc.Content, EmbedKey: "content_vector"}, "metadata": {Value: metadataBytes}, }, }, nil }, } // 创建embedding客户端 这里注意 这里使用的是embedding/openai 要使用的类型是openai.Embedder embeddingClient, err := openai.NewEmbedder(ctx, &openai.EmbeddingConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "text-embedding-v1", }) if err != nil { return nil, err } if embeddingClient == nil { fmt.Println("embeddingClient is nil") return nil, errors.New("embeddingClient is nil") } config.Embedding = embeddingClient idr, err = redis.NewIndexer(ctx, config) if err != nil { return nil, err } return idr, nil }
效果图

4. 查询增强
这部分仍然使用了eino的图编排,极其强大的能力。
▼mermaid复制代码graph LR START(("START")) --> ChatTemplate(提示词模板) RedisRetriever(向量检索器) --> ChatTemplate ChatTemplate --> ChatModel(大模型) ChatModel --> END(("END"))
注意要把模型的key换成你的,以及把redis-stack的连接换成自己的。其中redis-stack处的代码是创建索引检索器Retriever,不是上一节的存储器
▼text复制代码package main import ( "context" "errors" "fmt" "log" "strconv" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/components/retriever" "github.com/cloudwego/eino-ext/components/embedding/openai" openaiModel "github.com/cloudwego/eino-ext/components/model/openai" "github.com/cloudwego/eino-ext/components/retriever/redis" "github.com/cloudwego/eino/components/model" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" redisCli "github.com/redis/go-redis/v9" ) func main() { ctx := context.Background() r, err := BuildGraph(ctx) if err != nil { panic(err) } in := map[string]any{"question": "我想出国旅行需要注意什么"} output, err := r.Invoke(ctx, in) if err != nil { log.Panic(err) return } fmt.Println(output) } // BuildGraph 创建图 func BuildGraph(ctx context.Context) (r compose.Runnable[map[string]any, *schema.Message], err error) { const ( ChatTemplate = "ChatTemplate" ChatModel = "ChatModel" RedisRetriever = "RedisRetriever" ) g := compose.NewGraph[map[string]any, *schema.Message]() // 创建模板 chatTemplateKeyOfChatTemplate, err := newChatTemplate(ctx) if err != nil { return nil, err } _ = g.AddChatTemplateNode(ChatTemplate, chatTemplateKeyOfChatTemplate) // 创建模型 chatModel, err := newChatModel(ctx) if err != nil { return nil, err } _ = g.AddChatModelNode(ChatModel, chatModel) // 创建索引器 redisRetrieverKeyOfRetriever, err := newRedisRetriever(ctx) if err != nil { return nil, err } // 索引器与模型输出的中间转换 _ = g.AddRetrieverNode(RedisRetriever, redisRetrieverKeyOfRetriever, compose.WithOutputKey("documents")) _ = g.AddEdge(compose.START, ChatTemplate) _ = g.AddEdge(RedisRetriever, ChatTemplate) _ = g.AddEdge(ChatTemplate, ChatModel) _ = g.AddEdge(ChatModel, compose.END) r, err = g.Compile(ctx) if err != nil { return nil, err } return r, err } // newChatTemplate 聊天模板 func newChatTemplate(ctx context.Context) (ctp prompt.ChatTemplate, err error) { pt := prompt.FromMessages( schema.FString, schema.UserMessage("{question}?"), ) return pt, nil } // newChatModel 创建模型 func newChatModel(ctx context.Context) (model.ToolCallingChatModel, error) { // model.ToolCallingChatModel 也可以替换为model.BaseChatModel chatModel, err := openaiModel.NewChatModel(ctx, &openaiModel.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "qwen-turbo", }) if err != nil { return nil, err } return chatModel, nil } // redis检索器 func newRedisRetriever(ctx context.Context) (idr retriever.Retriever, err error) { // 创建redis索引器 可以读取存储 redisClient := redisCli.NewClient(&redisCli.Options{ Addr: "192.168.161.130:6379", // todo 换成你的地址 Password: "123456", DB: 1, }) err = redisClient.Ping(ctx).Err() if err != nil { return nil, err } config := &redis.RetrieverConfig{ Client: redisClient, Index: fmt.Sprintf("%s%s", "eino:doc:", "vector_index"), Dialect: 2, ReturnFields: []string{"content", "metadata", "distance"}, TopK: 8, VectorField: "content_vector", // 文档转换函数 存储时对文档进行处理 DocumentConverter: func(ctx context.Context, doc redisCli.Document) (*schema.Document, error) { resp := &schema.Document{ ID: doc.ID, Content: "", MetaData: map[string]any{}, } for field, val := range doc.Fields { if field == "content" { resp.Content = val } else if field == "metadata" { resp.MetaData[field] = val } else if field == "distance" { distance, err := strconv.ParseFloat(val, 64) if err != nil { continue } resp.WithScore(1 - distance) } } return resp, nil }, } // 创建embedding客户端 这里注意 这里使用的是embedding/openai 要使用的类型是openai.Embedder embeddingClient, err := openai.NewEmbedder(ctx, &openai.EmbeddingConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "text-embedding-v1", }) if err != nil { return nil, err } if embeddingClient == nil { fmt.Println("embeddingClient is nil") return nil, errors.New("embeddingClient is nil") } config.Embedding = embeddingClient idr, err = redis.NewRetriever(ctx, config) if err != nil { return nil, err } return idr, nil }
输出和langchain的几乎一样,因为调用的都是相同的模型。

eino不支持云知识库检索,这部分只能通过SDK或API实现。这里略
RAG知识库进阶
RAG核心特性
这部分内容的代码体现,其实在上文RAG基础部分很详细,特别是Eino部分,作者特意加了不少中文注释。而且每个组件的创建基本都采用工厂模式,更清晰的区分组件。
文档收集和切割 - ETL
文档收集和切割阶段,我们要对自己准备好的知识库文档进行处理,然后保存到向量数据库中,这个过程俗称ETL(抽取、转换、加载)。
ETL
Eino中,对document处理也是遵循读取,转换,写入的过程,并且过程中间互相独立,Eino提供了对应的组件,非常方便
- 读取文档:利用loader加载器,加载指定数据源中的文档,例如本地文件,网络资源
- 转换文档:通过Transform转换器,将不同格式的文档转换为合适的格式,最终变为document类型
- 写入文档:利用Indexer组件组件,将document类型的文档写入到指定位置。
抽取 Extract
Eino官方提供了一个实现了Load接口的加载器。
▼text复制代码// Loader is a document loader. type Loader interface { Load(ctx context.Context, src Source, opts ...LoaderOption) ([]*schema.Document, error) }
从数据源中获取到内容,转化为document
转换 Transform
Eino的转换目前主要有三个对象:
- recursive:递归分割器,用于将长文档按照指定大小递归地切分成更小的片段。
- semantic:语义分割器,用于基于语义相似度将长文档切分成更小的片段。
- Markdown:Markdown 分割器,用于根据 Markdown 文档的标题层级结构进行分割。
加载 Load
此处的加载是将文本写入到目标中。写入到模板可以使用,写入到缓存或者数据库可以实现文档的保存。 加载由Retriever和Indexer实现,例如上方我们将redis-stack中的内容读取出来可以使用Retriever,而将文档保存进数据或者缓存可以使用Retriever来实现。两者都是实现了一个接口。 存储
▼text复制代码type Indexer interface { // Store stores the documents. Store(ctx context.Context, docs []*schema.Document, opts ...Option) (ids []string, err error) // invoke }
获取
▼text复制代码type Retriever interface { Retrieve(ctx context.Context, query string, opts ...Option) ([]*schema.Document, error) }
实例
实例其实就可以忽略了,我们在上面的实例中其实非常详细了。
向量的转换和存储
Eino提供了很多向量存储的方式,redis,es8, VikingDB(火山引擎的云向量数据库),Milvus 2.x(本地向量数据库)
哪一种方案都行,并且作者觉得存储的话肯定是要存到数据库,那么就拿Milvus作为案例。
Milvus
Milvus是一款开源向量数据库,支持针对TB级向量的curd操作。 Milvus 采用共享存储架构,存储计算完全分离,计算节点支持横向扩展。从架构上来看,Milvus 遵循数据流和控制流分离,整体分为了四个层次:分别为接入层(access layer)、协调服务(coordinator service)、执行节点(worker node)和存储层(storage)。各个层次相互独立,独立扩展和容灾。
Milvus一共分为三个部分,这三部分需要单独下载
- Milvus:负责提供系统的核心功能。
- etcd :元数据引擎,用于管理 Milvus 内部组件的元数据访问和存储,例如:proxy、index node 等。
- MinIO :存储引擎,负责维护 Milvus 的数据持久化。
Milvus安装
作者使用docker在linux上安装的。总体来看安装很简单,中间也没出现任何问题。
官方docker快速安装
docker安装,首先安装docker-compose文件,使用下方链接。(如果linux中老是连接失败,就在自己电脑的浏览器中输入这个地址,下载下来后放到linux中就好了)。作者在这里直接把文件内容粘贴出来了。其中minio的账户和密码(不是连接milvux的账号密码)就在这个文件中可以更改。
如果你想在安装的时候添加用户名和密码,请看下一节标题
milvus-standalone-docker-compose.yml
▼text复制代码version: '3.5' services: etcd: container_name: milvus-etcd image: quay.io/coreos/etcd:v3.5.5 environment: - ETCD_AUTO_COMPACTION_MODE=revision - ETCD_AUTO_COMPACTION_RETENTION=1000 - ETCD_QUOTA_BACKEND_BYTES=4294967296 - ETCD_SNAPSHOT_COUNT=50000 volumes: - ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/etcd:/etcd command: etcd -advertise-client-urls=http://127.0.0.1:2379 -listen-client-urls http://0.0.0.0:2379 --data-dir /etcd healthcheck: test: ["CMD", "etcdctl", "endpoint", "health"] interval: 30s timeout: 20s retries: 3 minio: container_name: milvus-minio image: minio/minio:RELEASE.2023-03-20T20-16-18Z environment: MINIO_ACCESS_KEY: minioadmin MINIO_SECRET_KEY: minioadmin ports: - "9001:9001" - "9000:9000" volumes: - ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/minio:/minio_data command: minio server /minio_data --console-address ":9001" healthcheck: test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"] interval: 30s timeout: 20s retries: 3 standalone: container_name: milvus-standalone image: milvusdb/milvus:v2.4.5 command: ["milvus", "run", "standalone"] security_opt: - seccomp:unconfined environment: ETCD_ENDPOINTS: etcd:2379 MINIO_ADDRESS: minio:9000 volumes: - ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/milvus:/var/lib/milvus healthcheck: test: ["CMD", "curl", "-f", "http://localhost:9091/healthz"] interval: 30s start_period: 90s timeout: 20s retries: 3 ports: - "19530:19530" - "9091:9091" depends_on: - "etcd" - "minio" networks: default: name: milvus
▼text复制代码wget https://github.com/milvus-io/milvus/releases/download/v2.4.5/milvus-standalone-docker-compose.yml -O docker-compose.yml
运行,如果docker compose up -d提示没有参数-d 那就是你的docker compose版本比较低
▼text复制代码docker compose up -d docker-compose up -d # 低docker compose版本运行这个
如果你的服务器已经有了milvus网络,可以删掉在重新运行
▼text复制代码docker network ls docker network rm milvus
添加用户名和密码
首先我们需要下载milvus.yaml配置文件。作者开着加速器都连接失败。只好浏览器打开直接复制了。
▼text复制代码wget https://raw.githubusercontent.com/milvus-io/milvus/v2.5.10/configs/milvus.yaml
将文件中的authorizationEnabled改为true。默认的用户名和密码是root:Milvus。如果要改密码,密码和还必须8-64个字符之间,且必须包含大写字母、小写字母、数字和特殊字符中的三种(怎么这么多事)。
▼text复制代码security: authorizationEnabled: true
最后就是将我们配置文件挂载进去。修改docker-compose文件,在valumes中添加一条
▼text复制代码standalone: volumes: - ${DOCKER_VOLUME_DIRECTORY:-.}/volumes/milvus:/var/lib/milvus - /you local path/milvus.yaml:/milvus/configs/milvus.yaml # 更换为你的yaml文件地址
可视化
可以安装可视化工具attu。
作者为了以后使用方便,安装的windows版本



milvus的使用
增查操作我们基本都由eino完成
创建Collection
创建Collection,在 Milvus 中,Collection(集合) 是数据组织的核心单元,类似于传统数据库中的“表”,但专门针对向量数据优化。可以存储向量和和非向量数据(下方的案例也可以体现),如文本、元数据(json格式)
作者选择直接通过attu创建。这里展示我们创建的collection的结构,需要注意的是,collection一旦创建,结构便不可修改,所以创建的时候需要谨慎。

其中collection最重要的字段便是vector,该字段专门用于存储向量数据,其设计目的是支持高效的相似性搜索。存储该数据要求,数据的长度必须为8的倍数。
eino添加数据
我们使用eino来collection中添加数据。注意这里,作者花了两天天才解决的坑,真是坑死作者了 (;´д`)ゞ 作者创建时vector选择的二进制类型,与eino官方的案例一致。阿里云的embedded模型选择的维度也是128,正常只要模型的维度和存储数据维度一致即可了。但是作者使用的时候一直报错,提示维度对不齐。
▼text复制代码[Indexer.Store] failed to insert rows: the num_rows (1024) of field (vector) is not equal to passed num_rows (2): invalid parameter[expected=1024][actual=2]
作者绞尽脑汁,用尽各种办法,换了多个模型,把64-2048所有维度都测试了一遍,始终不行。没办法了,只能深入底层,为此只能去看eino源码了。
在"github.com/cloudwego/eino-ext/components/indexer/milvus"源码中有个store方法,就是将数据存储到collection。其中在debug时,作者发现我们的文本数据被模型转换为向量后的维度是正常的128位数,理应来讲我们没问题的,但是!,作者注意模型处理完的向量是[][]float6464位数的,我们要存储为binary类型的,于是作者往下继续找,发现了一个eino源码中的方法,将[]float64转换为[]byte类型
▼text复制代码func vector2Bytes(vector []float64) []byte { float32Arr := make([]float32, len(vector)) for i, v := range vector { float32Arr[i] = float32(v) } bytes := make([]byte, len(float32Arr)*4) for i, v := range float32Arr { binary.LittleEndian.PutUint32(bytes[i*4:], math.Float32bits(v)) } return bytes }
然后作者突然受到启发!自己测试后发现,我们的float64的类型的数据转换为byte后切片的长度不是原来的128位了!测试代码如下,测试后发现embeddings[0]也就是模型给我们返回的维度是128位。没问题,而bytes,也就是被vector2Bytes转化为byte类型的数据后,维度竟然是512位,怪不得维度不一致。
▼text复制代码emb, err := openai.NewEmbedder(ctx, &openai.EmbeddingConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "text-embedding-v3", Dimensions: &[]int{128}[0], }) embeddings, err := emb.EmbedStrings(ctx, []string{"我是中国人"}) // 调用向量模型 将我是中国人变为向量数据 if err != nil { log.Fatalf("Failed to embed: %v", err) return } fmt.Println(len(embeddings[0])) // 与Dimensions位数相同 var vectors [][]byte for _, v := range embeddings { bytes := vector2Bytes(v) fmt.Println("len bytes :", len(bytes)) // 打印[]float64转变为[]byte后的长度 vectors = append(vectors, bytes) }
于是作者自信满满的将collection的vector维度改为512,这下不得直接拿捏bug,结果还是维度对不齐!。作者人傻了。但是离答案更进一步了。报错信息里由原来的1024位对2,变成了8对1,于是作者继续让vector的维度翻倍,就这样不断尝试下,最后锁定在了vector维度为4096时,成功存储了数据。 于是作者将模型维度变为1024,继续尝试,最后vector维度锁定在了32768时成功存入数据。
不难发现,模型维度和vector维度之间是32倍的关系,这个关系的原因是因为数据位数的问题,我们的float64是64位,而byte是8位,两者是8倍的关系。但是上方的vector2Bytes方法,将float64转换为了float32,在转化为byte类型,就是4倍的关系,然后vector中的二进制数的维度好像是bit数,一个byte占8bit,最后的关系就是 float64 -> float32 -> byte -> bit 最后维度变化就是 原始维度 -> *4 -> *8 -> vector的维度
eion中使用milvus
▼text复制代码package main import ( "context" "log" "github.com/cloudwego/eino-ext/components/embedding/openai" "github.com/milvus-io/milvus-sdk-go/v2/entity" "github.com/milvus-io/milvus-sdk-go/v2/client" "github.com/cloudwego/eino-ext/components/indexer/milvus" "github.com/cloudwego/eino/schema" ) func main() { // 连接Milvus ctx := context.Background() cli, err := client.NewClient(ctx, client.Config{ Address: "192.168.161.130:19530", Username: "root", Password: "Milvus", DBName: "Eino", }) if err != nil { log.Fatalf("Failed to create client: %v", err) return } defer cli.Close() if err != nil { log.Fatalf("Failed to create embedding: %v", err) return } // 创建collection字段结构 必须以我们的collection中的字段结构完全一致才可以 filed := []*entity.Field{ entity.NewField(). WithName("id"). WithIsPrimaryKey(true). WithDataType(entity.FieldTypeVarChar). WithMaxLength(255), entity.NewField(). WithName("vector"). WithIsPrimaryKey(false). WithDataType(entity.FieldTypeBinaryVector). WithDim(128 * 32), entity.NewField(). WithName("content"). WithIsPrimaryKey(false). WithDataType(entity.FieldTypeVarChar). WithMaxLength(1024), entity.NewField(). WithName("metadata"). WithIsPrimaryKey(false). WithDataType(entity.FieldTypeJSON), } //创建embedding客户端 这里注意 这里使用的是embedding/openai 要使用的类型是openai.Embedder emb, err := openai.NewEmbedder(ctx, &openai.EmbeddingConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "text-embedding-v3", Dimensions: &[]int{128}[0], }) // 如果不指定collection 会默认生成一个 eino_collection 默认的vector为81920太大了 所以作者建了一个新的collection indexer, err := milvus.NewIndexer(ctx, &milvus.IndexerConfig{ Client: cli, Embedding: emb, Collection: "eino_colletion_test", // 指定milvus的collection名称 Fields: filed, // 指定你的collection中的字段 如果你的collection不是默认的 就必须填写 }) if err != nil { panic(err) return } log.Printf("Indexer created success") // 要存储的 documents docs := []*schema.Document{ { ID: "milvus-1", Content: "国外", MetaData: map[string]any{ "h1": "milvus", }, }, { ID: "milvus-2", Content: "国内", MetaData: map[string]any{ "h1": "milvus", }, }, } ids, err := indexer.Store(ctx, docs) if err != nil { log.Fatalf("Failed to store: %v", err) return } log.Printf("Store success, ids: %v", ids) }
向量存储原理
在向量数据库中,查询与传统关系型数据库有所不同。向量库执行的是相似性搜索,而非精确匹配,具体流程我们在上一节教程中有了解,可以再复习下。
- 嵌入转换:当文档被添加到向量存储时,eino会使用嵌入模型(例如上方我们已经使用的百炼的text-embedding-v3)将文本转换为向量。
- 相似度计算:查询时,查询文本同样被转换为向量,然后系统计算此向量与存储中所有向量的相似度。
- 相似度度量:常用的相似度计算方法包括:
- 余弦相似度:计算两个向量的夹角余弦值,范围在-1到1之间
- 欧氏距离:计算两个向量间的直线距离
- 点积:两个向量的点积值
- 过滤与排序:根据相似度阈值过滤结果,并按相似度排序返回最相关的文档
文档过滤和检索
多查询扩展
多查询扩展,简单来说就用AI扩展用户的回答。这种方法说实话,太抽象。作者直接忽略该做法。实现也很简单,多并行调用几次ai,接收到问题后作为正式调用ai的输入即可。
eino中实现该功能不复杂,使用ReactAgent即可实现。底层使用的Graph。
查询重写和翻译
查询重写器,还是调用了ai,先让ai把我们的问题变得更加简洁明了。 其实可以用图来实现,开始调用ai优化问题,然后将问题再次喂给ai。而且调用问题的时候还可以使用比较廉价的模型。 需要注意的时,我们需要创建一个中间节点来将上一个模型的输出,填充到下一个模型的输入。为了方便演示,作者只调用了一个模型。
▼text复制代码package main import ( "context" "fmt" openaiModel "github.com/cloudwego/eino-ext/components/model/openai" "github.com/cloudwego/eino/components/model" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" "github.com/cloudwego/eino/compose" ) func main() { const ( ChatTemplate = "ChatTemplate" ChatModel = "ChatModel" QuestionTemplate = "ChatTemplateQuestion" QuestionModel = "ChatModelQuestion" Lambda = "Lambda" ) ctx := context.Background() g := compose.NewGraph[map[string]any, *schema.Message]() chatModel, err := newChatModel(ctx) if err != nil { panic(err) } _ = g.AddChatModelNode(ChatModel, chatModel) questionModel, err := newChatModel(ctx) if err != nil { panic(err) } _ = g.AddChatModelNode(QuestionModel, questionModel) chatTemplateKeyOfChatTemplate, err := newChatTemplate(ctx) if err != nil { panic(err) return } _ = g.AddChatTemplateNode(ChatTemplate, chatTemplateKeyOfChatTemplate) questionTemplate, err := newQuestionChatTemplate(ctx) if err != nil { panic(err) } _ = g.AddChatTemplateNode(QuestionTemplate, questionTemplate) _ = g.AddLambdaNode( Lambda, compose.InvokableLambda[*schema.Message, map[string]any](func(ctx context.Context, input *schema.Message) (map[string]any, error) { fmt.Println("question:", input.Content) return map[string]any{ "question": input.Content, }, nil }), ) _ = g.AddEdge(compose.START, QuestionTemplate) _ = g.AddEdge(QuestionTemplate, QuestionModel) _ = g.AddEdge(QuestionModel, Lambda) _ = g.AddEdge(Lambda, ChatTemplate) _ = g.AddEdge(ChatTemplate, ChatModel) _ = g.AddEdge(ChatModel, compose.END) compile, err := g.Compile(ctx) if err != nil { panic(err) } output, err := compile.Invoke(ctx, map[string]any{ "question": "如果我在出国旅行的时候很不幸遇到了一些让我感到很困难的事情的时候该怎么办呢?", }) if err != nil { panic(err) } fmt.Println(output) } func newQuestionChatTemplate(ctx context.Context) (ctp prompt.ChatTemplate, err error) { pt := prompt.FromMessages( schema.FString, schema.UserMessage("{question},把问题变得精炼短小"), ) return pt, nil } // newChatTemplate 聊天模板 func newChatTemplate(ctx context.Context) (ctp prompt.ChatTemplate, err error) { pt := prompt.FromMessages( schema.FString, schema.UserMessage("{question}?"), ) return pt, nil } // newChatModel 创建模型 func newChatModel(ctx context.Context) (model.ToolCallingChatModel, error) { // model.ToolCallingChatModel 也可以替换为model.BaseChatModel chatModel, err := openaiModel.NewChatModel(ctx, &openaiModel.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "you key", // todo 替换为你的key Model: "qwen-turbo", }) if err != nil { return nil, err } return chatModel, nil }
翻译就不实现了,其实翻译也很简单,在最后调用模型的时候,在模板里面添加让模型翻译就行了。
检索器配置 (milvus查询案例)
检索器配置其实已经在创建检索器的时候设置好了。例如下方milvus.RetrieverConfig各种属性。详情属性见官方 至于文档过滤器,我们在上方介绍过markdown过滤器,详情其他过滤器。
▼text复制代码package main import ( "context" "log" "github.com/milvus-io/milvus-sdk-go/v2/entity" "github.com/cloudwego/eino-ext/components/embedding/openai" "github.com/cloudwego/eino-ext/components/retriever/milvus" "github.com/milvus-io/milvus-sdk-go/v2/client" ) func main() { ctx := context.Background() cli, err := client.NewClient(ctx, client.Config{ Address: "192.168.161.130:19530", Username: "root", Password: "Milvus", DBName: "Eino", }) if err != nil { log.Fatalf("Failed to create client: %v", err) return } defer cli.Close() embeddingClient, err := openai.NewEmbedder(ctx, &openai.EmbeddingConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "text-embedding-v3", Dimensions: &[]int{128}[0], }) if err != nil { log.Fatalf("Failed to create NewEmbedder: %v", err) return } retriever, err := milvus.NewRetriever(ctx, &milvus.RetrieverConfig{ Client: cli, Collection: "eino_colletion_test", // 指定milvus的collection名称 Embedding: embeddingClient, // 指定你的embedding客户端 VectorField: "vector", // 集合中的向量字段名称 TopK: 3, // 搜索数量 ScoreThreshold: 0, // 搜索结果阈值 MetricType: entity.JACCARD, // 向量的度量类型 默认HAMMING 必须与实际一致 OutputFields: []string{ // 搜索返回的字段 "id", "content", "metadata", }, }) docs, err := retriever.Retrieve(ctx, "查阅国内旅行推荐手册,国内旅行需要注意哪些问题,有推荐课程链接吗?") if err != nil { log.Fatalf("Failed to retrieve: %v", err) return } log.Printf("Retrieve success, docs: %v", docs) }
查询增强和关联
错误处理机制
在实际应用中,可能出现很多种情况,如果找不到相关文档,相似度低,查询超时等等,良好的错误处理机制会提升用户体验。
异常处理包括:
- 允许空上下文查询
- 提供友好的错误提示
- 引导用户提供必要信息
eino中并没有错误处理,因为本身go的错误处理机制问题,我们可以直接在正常逻辑中返回即可。例如空上下文问题,当我们的模板没有内容的时候,本身就会返回错误。以及检索其没有找到索引向量的时候也会报错。
▼text复制代码chatTemplateKeyOfChatTemplate, err := newChatTemplate(ctx) if err != nil { panic(err) return } // 没有数据产生本身就会报错 format, err := chatTemplateKeyOfChatTemplate.Format(ctx, map[string]any{}) if err != nil { panic(err) return }
工具调用
AI应用框架本身自带大量的工具可以使用。例如Eino Eino自带了多种查询工具,例如Google、维基百科、DuckDuckGo安全隐私搜索。http请求工具,sequentialthinking结构化思考工具,MCP工具等。
利用这些工具AI才能真正变成一个智能体。
Eino工具开发
Eino自定义工具
Eino中自定义工具并不复杂,其实在上文AI应用框架的tool标题已经介绍了工具的定义使用。
工具生态
eino工具生态不是很好,目前也是只有官方工具以供使用。想要拓展的话,可以使用mcp,因为eino的tool本身就封装了mcp,所以也是支持使用该方法。
eino工具进阶
而且eino本身设计挺复杂的,作者实在难以搞懂底层运行逻辑。最后也只是了解一点点
其实eino的底层基于图论。 组件是最小的运行节点,基于 Graph 模型 (node + edge) 的,以组件为原子节点的,以上下游类型对齐为基础的编排。
工具执行是在chatModel中,我们把工具绑定给模型,他不会产生任何工具的调用,当模型返回的信息中带有工具涉及的tool token,框架就会去找到对应的tool,执行tool。
MCP协议
使用MCP
kratos就在不久前刚刚更新了mcp包,我们直接更新新版的kratos即可。同时eino本身也支持mcp。两者都是用的开源sdkmark3labs/mcp-go
客户端SSE使用MCP
mcp官方 MCP-Go 支持stdio SSE两种连接方式。
- stdio:standard input output,标准输入输出,适用于mcp client和mcp server部署在同一个机器上,不涉及跨机器。
- sse:Server-Sent Events,服务器发送事件,是一个基于HTTP的协议。适用于mcp client与mcp server不在同一个机器中的场景,涉及网络传输。也是接下来我们要用的协议。
案例:创建高德mcp客户端
▼text复制代码package main import ( "context" "fmt" "log" "github.com/mark3labs/mcp-go/mcp" "github.com/mark3labs/mcp-go/client/transport" "github.com/mark3labs/mcp-go/client" ) func main() { ctx := context.Background() // 高德地图key mapKey := "your key" // 因为我们要使用第三方的mcp 不在同一个服务器上 使用SSE连接 sse, err := transport.NewSSE(fmt.Sprintf("https://mcp.amap.com/sse?key=%s", mapKey)) if err != nil { panic(err) } // 创建一个mcp客户端 c := client.NewClient(sse) // 启动客户端 err = c.Start(ctx) if err != nil { panic(err) } // 设置通知处理函数 c.OnNotification(func(notification mcp.JSONRPCNotification) { fmt.Printf("Received notification: %s\n", notification.Method) }) // 初始化客户端 initRequest := mcp.InitializeRequest{} initRequest.Params.ProtocolVersion = mcp.LATEST_PROTOCOL_VERSION initRequest.Params.ClientInfo = mcp.Implementation{ Name: "MCP-Go Simple Client Example", Version: "1.0.0", } initRequest.Params.Capabilities = mcp.ClientCapabilities{} serverInfo, err := c.Initialize(ctx, initRequest) if err != nil { log.Fatalf("Failed to initialize: %v", err) } // 显示服务器信息 fmt.Printf("Connected to server: %s (version %s)\n", serverInfo.ServerInfo.Name, serverInfo.ServerInfo.Version) fmt.Printf("Server capabilities: %+v\n", serverInfo.Capabilities) // 列出服务器支持的工具 if serverInfo.Capabilities.Tools != nil { fmt.Println("Fetching available tools...") toolsRequest := mcp.ListToolsRequest{} toolsResult, err := c.ListTools(ctx, toolsRequest) if err != nil { log.Printf("Failed to list tools: %v", err) } else { fmt.Printf("Server has %d tools available\n", len(toolsResult.Tools)) for i, tool := range toolsResult.Tools { fmt.Printf(" %d. %s - %s\n", i+1, tool.Name, tool.Description) } } } // 列出服务器获取可用资源 if serverInfo.Capabilities.Resources != nil { fmt.Println("Fetching available resources...") resourcesRequest := mcp.ListResourcesRequest{} resourcesResult, err := c.ListResources(ctx, resourcesRequest) if err != nil { log.Printf("Failed to list resources: %v", err) } else { fmt.Printf("Server has %d resources available\n", len(resourcesResult.Resources)) for i, resource := range resourcesResult.Resources { fmt.Printf(" %d. %s - %s\n", i+1, resource.URI, resource.Name) } } } fmt.Println("Client initialized successfully. Shutting down...") err = c.Close() if err != nil { return } }
输出结果:15个工具,和我们在百炼中的一样。
▼text复制代码Connected to server: amap-sse-server (version 1.0.0) Server capabilities: {Experimental:map[] Logging:0x158ffc0 Prompts:<nil> Resources:<nil> Tools:0xc00000b09a} Fetching available tools... Server has 15 tools available 1. maps_direction_bicycling - 骑行路径规划用于规划骑行通勤方案,规划时会考虑天桥、单行线、封路等情况。最大支持 500km 的骑行路线规划 2. maps_direction_driving - 驾车路径规划 API 可以根据用户起终点经纬度坐标规划以小客车、轿车通勤出行的方案,并且返回通勤方案的数据。 3. maps_direction_transit_integrated - 根据用户起终点经纬度坐标规划综合各类公共(火车、公交、地铁)交通方式的通勤方案,并且返回通勤方案的数据,跨城场景下必须传起点城市与终点城市 4. maps_direction_walking - 根据输入起点终点经纬度坐标规划100km 以内的步行通勤方案,并且返回通勤方案的数据 5. maps_distance - 测量两个经纬度坐标之间的距离,支持驾车、步行以及球面距离测量 6. maps_geo - 将详细的结构化地址转换为经纬度坐标。支持对地标性名胜景区、建筑物名称解析为经纬度坐标 7. maps_regeocode - 将一个高德经纬度坐标转换为行政区划地址信息 8. maps_ip_location - IP 定位根据用户输入的 IP 地址,定位 IP 的所在位置 9. maps_schema_personal_map - 用于行程规划结果在高德地图展示。将行程规划位置点按照行程顺序填入lineList,返回结果为高德地图打开的URI链接,该结果不需总结,直接返回! 10. maps_around_search - 周边搜,根据用户传入关键词以及坐标location,搜索出radius半径范围的POI 11. maps_search_detail - 查询关键词搜或者周边搜获取到的POI ID的详细信息 12. maps_text_search - 关键字搜索 API 根据用户输入的关键字进行 POI 搜索,并返回相关的信息 13. maps_schema_navi - Schema唤醒客户端-导航页面,用于根据用户输入终点信息,返回一个拼装好的客户端唤醒URI,用户点击该URI即可唤起对应的客户端APP。唤起客户端后,会自动跳转到导航页面。 14. maps_schema_take_taxi - 根据用户输入的起点和终点信息,返回一个拼装好的客户端唤醒URI,直接唤起高德地图进行打车。直接展示生成的链接,不需要总结 15. maps_weather - 根据城市名称或者标准adcode查询指定城市的天气 Client initialized successfully. Shutting down... Process finished with the exit code 0
将mcp提供的方法转换为工具提供给模型使用。不是很复杂,就是代码多一点。
▼text复制代码package main import ( "context" "fmt" "log" "github.com/cloudwego/eino/compose" openaiModel "github.com/cloudwego/eino-ext/components/model/openai" toolMcp "github.com/cloudwego/eino-ext/components/tool/mcp" "github.com/cloudwego/eino/components/model" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/schema" "github.com/mark3labs/mcp-go/mcp" "github.com/mark3labs/mcp-go/client/transport" "github.com/mark3labs/mcp-go/client" ) func main() { ctx := context.Background() // 高德地图key mapKey := "your key" // 因为我们要使用第三方的mcp 不在同一个服务器上 使用SSE连接 sse, err := transport.NewSSE(fmt.Sprintf("https://mcp.amap.com/sse?key=%s", mapKey)) if err != nil { panic(err) } // 创建一个mcp客户端 c := client.NewClient(sse) // 启动客户端 err = c.Start(ctx) if err != nil { panic(err) } // 设置通知处理函数 实际使用可忽略 c.OnNotification(func(notification mcp.JSONRPCNotification) { fmt.Printf("Received notification: %s\n", notification.Method) }) // 初始化客户端 initRequest := mcp.InitializeRequest{} initRequest.Params.ProtocolVersion = mcp.LATEST_PROTOCOL_VERSION initRequest.Params.ClientInfo = mcp.Implementation{ Name: "MCP-Go Simple Client Example", Version: "1.0.0", } initRequest.Params.Capabilities = mcp.ClientCapabilities{} serverInfo, err := c.Initialize(ctx, initRequest) if err != nil { log.Fatalf("Failed to initialize: %v", err) } // 显示服务器信息 fmt.Printf("Connected to server: %s (version %s)\n", serverInfo.ServerInfo.Name, serverInfo.ServerInfo.Version) fmt.Printf("Server capabilities: %+v\n", serverInfo.Capabilities) // 调用eino官方工具包 将mcp工具转化为eino工具 获取工具列表 tools, err := toolMcp.GetTools(ctx, &toolMcp.Config{ Cli: c, }) if err != nil { panic(err) } // 创建模板 ctp, err := newChatTemplateMcp(ctx) if err != nil { panic(err) } // 创建模型 modelMcp, err := newChatModelMcp(ctx) if err != nil { panic(err) } // 获取工具信息 var toolsInfo []*schema.ToolInfo for _, t := range tools { info, err := t.Info(ctx) if err != nil { fmt.Errorf("GetToolInfo failed, err=%v", err) continue } toolsInfo = append(toolsInfo, info) } // 为模型绑定工具 withTools, err := modelMcp.WithTools(toolsInfo) if err != nil { panic(err) } const ( chatModel = "chatModel" template = "template" modelTool = "modelTool" ) // 创建图 g := compose.NewGraph[map[string]any, []*schema.Message]() _ = g.AddChatModelNode(chatModel, withTools) _ = g.AddChatTemplateNode(template, ctp) // 创建工具节点 toolsNode, err := compose.NewToolNode(ctx, &compose.ToolsNodeConfig{ Tools: tools, }) _ = g.AddToolsNode(modelTool, toolsNode) _ = g.AddEdge(compose.START, template) _ = g.AddEdge(template, chatModel) _ = g.AddEdge(chatModel, modelTool) _ = g.AddEdge(modelTool, compose.END) // 编译图 compile, err := g.Compile(ctx) if err != nil { panic(err) } // 跑图 output, err := compile.Invoke(ctx, map[string]any{ "question": "上海东方明珠附近有什么宾馆和饭店吗?", }) if err != nil { panic(err) } fmt.Println("invoke result: ", output) fmt.Println("Client initialized successfully. Shutting down...") err = c.Close() if err != nil { return } } // newChatTemplate 聊天模板 func newChatTemplateMcp(ctx context.Context) (ctp prompt.ChatTemplate, err error) { pt := prompt.FromMessages( schema.FString, schema.UserMessage("{question}"), ) return pt, nil } // newChatModelMcp 创建模型 func newChatModelMcp(ctx context.Context) (model.ToolCallingChatModel, error) { // model.ToolCallingChatModel 也可以替换为model.BaseChatModel chatModel, err := openaiModel.NewChatModel(ctx, &openaiModel.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "your key", // todo 替换为你的key Model: "qwen-turbo", }) if err != nil { return nil, err } return chatModel, nil }
最后生成的结果
▼text复制代码Connected to server: amap-sse-server (version 1.0.0) Server capabilities: {Experimental:map[] Logging:0x192dfa0 Prompts:<nil> Resources:<nil> Tools:0xc00033c8aa} invoke result: [tool: {"content":[{"type":"text","text":"{\"pois\":[{\"id\":\"B001543C3B\",\"name\":\"锦佳饭店\",\"address\":\"大名路60号(北外滩地区、近苏州路)\",\"typecode\":\"100100\",\"photo\":\"http://store.is.autonavi.com/showpic/b2e17def05c113b01184e70af5606f32\"}]}"}]} tool_call_id: call_1cd1e270c15a467abda6ce] Client initialized successfully. Shutting down... Process finished with the exit code 0
结果中生成的图片

创建MCP Server
这里就需要kratos了,因为kratos添加了mcp-server模块,快速启动server非常方便。 例如下方案例,创建一个tool函数,创建服务即可。Health 是一个检验tool是否健康的中间件,此处可以忽略
▼text复制代码package main import ( "context" "errors" "fmt" "net/http" tm "github.com/go-kratos/kratos/contrib/transport/mcp/v2" "github.com/go-kratos/kratos/v2" "github.com/mark3labs/mcp-go/mcp" ) // helloHandler is a tool handler func helloHandler(ctx context.Context, request mcp.CallToolRequest) (*mcp.CallToolResult, error) { name, ok := request.Params.Arguments.(string) if !ok { return nil, errors.New("name must be a string") } return mcp.NewToolResultText(fmt.Sprintf("Hello, %s!", name)), nil } // Health is a middleware that handles health checks func Health(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if r.URL.Path == "/health/ready" { w.WriteHeader(http.StatusOK) return } next.ServeHTTP(w, r) }) } func main() { srv := tm.NewServer("kratos-mcp", "v1.0.0", tm.Address(":8000"), tm.Middleware(Health)) tool := mcp.NewTool("hello_world", mcp.WithDescription("Say hello to someone"), mcp.WithString("name", mcp.Required(), mcp.Description("Name of the person to greet"), ), ) // Add tool handler srv.AddTool(tool, helloHandler) // creates a kratos application app := kratos.New( kratos.Name("kratos-app"), kratos.Server(srv), ) if err := app.Run(); err != nil { panic(err) } }
我们启动另一个程序来测试是否创建成功,代码和上方高德地图案例一致。只是把url换成本地而已
▼text复制代码sse, err := transport.NewSSE("http://127.0.0.1:8000/sse")
输出结果:成功打印hello_world - Say hello to someone
▼text复制代码Connected to server: kratos-mcp (version v1.0.0) Server capabilities: {Experimental:map[] Logging:<nil> Prompts:<nil> Resources:<nil> Tools:0xc00000ad3a} Fetching available tools... Server has 1 tools available 1. hello_world - Say hello to someone Client initialized successfully. Shutting down... Process finished with the exit code 0
AI智能体构建
使用AI智能体
在对应的平台,例如阿里云百炼创建一个智能体,智能体可以添加知识库、应用、mcp服务等多种功能。
在程序中使用的话,go推荐使用sdk或api实现,因为智能体本身已经把能集成的功能集成好了,不需要我们在扩展什么。所以这部不是langchain和eino能涉及的,毕竟langchain和eino本身就是可以用来构建智能体的。
创建AI智能体
ReactAgent
ReactAgent 介绍
eino中自带ReactAgent。
react agent 底层使用 compose.Graph 作为编排方案,一般来说有 2 个节点: ChatModel、Tools,中间运行过程中的所有历史消息都会放入 state 中,在将所有历史消息传递给 ChatModel 之前,会 copy 消息交由 MessageModifier 进行处理,处理的结果再传递给 ChatModel。直到 ChatModel 返回的消息中不再有 tool call,则返回最终消息。

当 Tools 列表中至少有一个 Tool 配置了 ReturnDirectly 时,ReAct Agent 结构会更复杂:在 ToolsNode 之后会增加一个 Branch,判断是否调用了一个 ReturnDirectly 的 Tool,如果是,直接 END,否则照旧进入 ChatModel。
ReactAgent使用
经过作者多次测试后发现,生成结果不是很稳定,有时生成的很优秀,有时候会生成不出来。可以设置Agent步数多一些,这样就能大大提高生成结果的准确性。
▼text复制代码package main import ( "context" "fmt" "log" openaiModel "github.com/cloudwego/eino-ext/components/model/openai" toolMcp "github.com/cloudwego/eino-ext/components/tool/mcp" "github.com/cloudwego/eino/components/model" "github.com/cloudwego/eino/components/prompt" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/flow/agent/react" "github.com/cloudwego/eino/schema" "github.com/mark3labs/mcp-go/client" "github.com/mark3labs/mcp-go/client/transport" "github.com/mark3labs/mcp-go/mcp" ) func main() { ctx := context.Background() // 高德地图key mapKey := "bd0145ed6d98cce8c91a10fd9a2682e3" // SSE连接 sse, err := transport.NewSSE(fmt.Sprintf("https://mcp.amap.com/sse?key=%s", mapKey)) if err != nil { panic(err) } // 创建一个mcp客户端 c := client.NewClient(sse) // 启动客户端 err = c.Start(ctx) if err != nil { panic(err) } // 初始化客户端 initRequest := mcp.InitializeRequest{} initRequest.Params.ProtocolVersion = mcp.LATEST_PROTOCOL_VERSION initRequest.Params.ClientInfo = mcp.Implementation{ Name: "MCP-Go Simple Client Example", Version: "1.0.0", } initRequest.Params.Capabilities = mcp.ClientCapabilities{} _, err = c.Initialize(ctx, initRequest) if err != nil { log.Fatalf("Failed to initialize: %v", err) } // 调用eino官方工具包 将mcp工具转化为eino工具 获取工具列表 tools, err := toolMcp.GetTools(ctx, &toolMcp.Config{ Cli: c, }) // 创建模板 ctp, err := newChatTemplateAgent(ctx) if err != nil { panic(err) } // 创建模型 modelAgent, err := newChatModelAgent(ctx) if err != nil { panic(err) } //创建reactAgent代理 agent, err := react.NewAgent(ctx, &react.AgentConfig{ ToolCallingModel: modelAgent, ToolsConfig: compose.ToolsNodeConfig{ Tools: tools, }, MaxStep: 20, // 设置最大步数 防止死循环 agent基本2步一个循环 设置n个循环的时候 至少执行2*(n-1)步 }) format, err := ctp.Format(ctx, map[string]any{ "question": "推荐两家离济南大明湖比较近的酒店或宾馆,以及这两家店距离有多远呢?", }) if err != nil { panic(err) } generate, err := agent.Generate(ctx, format) if err != nil { panic(err) } fmt.Println("invoke result : ", generate) } // newChatTemplate 聊天模板 func newChatTemplateAgent(ctx context.Context) (ctp prompt.ChatTemplate, err error) { pt := prompt.FromMessages( schema.FString, &schema.Message{ Role: schema.User, Content: "{question}?", }, ) return pt, nil } // newChatModelAgent 创建模型 func newChatModelAgent(ctx context.Context) (model.ToolCallingChatModel, error) { // model.ToolCallingChatModel 也可以替换为model.BaseChatModel chatModel, err := openaiModel.NewChatModel(ctx, &openaiModel.ChatModelConfig{ BaseURL: "https://dashscope.aliyuncs.com/compatible-mode/v1", APIKey: "sk-253f0d5c02384aefbc27348f22412229", // todo 替换为你的key Model: "qwen-turbo", }) if err != nil { return nil, err } return chatModel, nil }
输出:
▼text复制代码invoke result : assistant: 我为您找到了两家靠近济南大明湖的酒店: 1. **银座精宿酒店(济南历山路山东大学店)**: - 地址:历山路辅路与利农庄路交叉口东北160米 - 距离:约500米 2. **银座精宿酒店(济南大明湖历山路店)**: - 地址:历山路48号 - 距离:约600米 请注意,这些酒店都位于大明湖附近,具体选择可以根据您的需求决定。 finish_reason: stop usage: &{4117 128 4245}
应用搭建
现在掌握了eino的基本能力后,便可以开发应用了。
整个项目业务只有两个模块
用户模块
经典的增删改查环节,只不过这次为了简便,作者删除了casbin。只留下了一个jwt校验环节。用户就不提了,千篇一律的设计
就不过多赘述了
chat 核心模块
表结构设计
作者将用户的对话内容保存到了mysql。 表结构设计比较简单,用户表,会话表,消息表。用户每次对话需要建立一个会话,会话中每次用户对话后都会将对话内容保存下来。消息按照角色存储,每条消息存储内容的同时,存储消息发送者。便于前端构建对话框。
▼text复制代码DROP TABLE IF EXISTS `conversations`; CREATE TABLE `conversations` ( `conversation_id` bigint(0) NOT NULL AUTO_INCREMENT, `user_id` bigint(0) NULL DEFAULT NULL, `title` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '对话标题/摘要', `created_at` timestamp(0) NULL DEFAULT CURRENT_TIMESTAMP(0), `updated_at` timestamp(0) NULL DEFAULT CURRENT_TIMESTAMP(0) ON UPDATE CURRENT_TIMESTAMP(0), `deleted_at` timestamp(0) NULL DEFAULT NULL, PRIMARY KEY (`conversation_id`) USING BTREE ) ENGINE = InnoDB AUTO_INCREMENT = 124072407076865 CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic; -- ---------------------------- -- Table structure for messages -- ---------------------------- DROP TABLE IF EXISTS `messages`; CREATE TABLE `messages` ( `message_id` bigint(0) NOT NULL AUTO_INCREMENT COMMENT '消息id', `conversation_id` bigint(0) NOT NULL COMMENT '会话id', `sender_type` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT 'user or system', `content` text CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL, `created_at` timestamp(0) NULL DEFAULT CURRENT_TIMESTAMP(0), `tokens` int(0) NULL DEFAULT NULL COMMENT '消息的token数量', PRIMARY KEY (`message_id`) USING BTREE ) ENGINE = InnoDB AUTO_INCREMENT = 1 CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic; -- ---------------------------- -- Table structure for sys_users -- ---------------------------- DROP TABLE IF EXISTS `sys_users`; CREATE TABLE `sys_users` ( `id` bigint(0) NOT NULL AUTO_INCREMENT COMMENT '主键id', `uid` bigint(0) NOT NULL COMMENT '用户id', `username` varchar(64) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '用户名(登入)', `password` varchar(128) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '密码', `phone` varchar(16) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL DEFAULT '' COMMENT '手机', `status` tinyint(0) NOT NULL DEFAULT 1 COMMENT '1=正常 2=异常', `post_ids` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL DEFAULT '' COMMENT '多岗位', `created_at` datetime(0) NULL DEFAULT NULL COMMENT '创建时间', `deleted_at` datetime(0) NULL DEFAULT NULL COMMENT '删除时间 gorm soft delete', PRIMARY KEY (`id`) USING BTREE, UNIQUE INDEX `username`(`username`) USING BTREE, UNIQUE INDEX `idx_uid`(`uid`) USING BTREE, INDEX `idx_deleted_at`(`deleted_at`) USING BTREE ) ENGINE = InnoDB AUTO_INCREMENT = 7 CHARACTER SET = utf8mb4 COLLATE = utf8mb4_general_ci ROW_FORMAT = Dynamic;
架构设计
kratos框架的经典架构。 为了复用模型与mcp连接和工具。作者把它们都放在在data层导出。同时在biz层创建了config结构体保存他们,这样xxxUseCase只需要引入config就可方便的使用它们了。
▼text复制代码type AiAgentConfig struct { Mcp *mcpClient.Client Model model.ToolCallingChatModel Tool []tool.BaseTool } func NewAiAgentConfig( mcp *mcpClient.Client, model model.ToolCallingChatModel, tool []tool.BaseTool, ) *AiAgentConfig { return &AiAgentConfig{ Mcp: mcp, Model: model, Tool: tool, } }
接口选型
AI对话设计有两种选型,sse和websockt。 如果是频繁对话或者长时间链接的话推荐websocket。 但是对于我们的项目没有那么频繁的对话,就是用轻量的sse。
利用gin框架创建sse连接去建立对话。 完整的业务逻辑就是。用户发起对话 -> 保存用户消息到mysql -> ReactAgent开始代理 -> 流式传输给用户 同时 后端保存结果 -> 开启协程异步调用AI为会话起标题 ->用户看到消息。
这里注意sse跨域问题。因为sse链接跨域时会预先发送一个OPTIONS请求,我们在跨域中间件中要注意主动去为options请求返回一个200的状态码,才能建立sse链接。
▼text复制代码// Cors gin的跨域设置 func Cors() gin.HandlerFunc { return func(c *gin.Context) { //if strings.Contains(c.Request.Header.Get("Accept"), "text/event-stream") { // c.Header("Content-Type", "text/event-stream") // c.Header("Cache-Control", "no-cache") // c.Header("Connection", "keep-alive") // c.Next() // return //} method := c.Request.Method c.Header("Access-Control-Allow-Origin", "*") //c.Header("Access-Control-Allow-Headers", "Content-Type,AccessToken,X-CSRF-Token, Authorization, Token") c.Header("Access-Control-Allow-Headers", "*") c.Header("Access-Control-Allow-Methods", "POST, GET, OPTIONS") //c.Header("Access-Control-Allow-Methods", "*") c.Header("Access-Control-Expose-Headers", "Content-Length, Access-Control-Allow-Origin, Access-Control-Allow-Headers, Content-Type") c.Header("Access-Control-Allow-Credentials", "true") if method == "OPTIONS" { c.AbortWithStatus(200) } c.Next() } }
核心接口代码
主要接口的代码。
▼text复制代码type ChatService struct { v1.UnimplementedAiServiceServer log *log.Helper ai *biz.AiAgentUseCase } func NewChatService(ai *biz.AiAgentUseCase, logger log.Logger) *ChatService { return &ChatService{ log: log.NewHelper(log.With(logger, "module", "service/chat")), ai: ai, } } func (chat *ChatService) ChatUseAgent(c *gin.Context) { // 设置SSE响应头 c.Header("Content-Type", "text/event-stream") c.Header("Cache-Control", "no-cache") c.Header("Connection", "keep-alive") // 绑定请求参数 var req v1.NewChatRequest if err := c.ShouldBindJSON(&req); err != nil { chat.log.Errorf("failed to bind request: %v", err) sendSSEError(c, http.StatusBadRequest, "invalid request format") return } if req.Question == "" { sendSSEError(c, http.StatusBadRequest, "对话内容不能为空") return } // 收集对话结果 var collectedContent strings.Builder // 创建agent agent, err := chat.ai.NewChatAgent(c) if err != nil { chat.log.Errorf("failed to create chat agent: %v", err) sendSSEError(c, http.StatusInternalServerError, constant.ErrParams.Error()) return } // 填充模板 prompt, err := chat.ai.NewChatPrompt(c, &req) if err != nil { chat.log.Errorf("failed to generate prompt: %v", err) sendSSEError(c, http.StatusInternalServerError, constant.ErrParams.Error()) return } // 保存用户对话 atoi, err := strconv.Atoi(req.ConversionId) if err != nil { chat.log.Errorf("failed to convert conversation id: %v", err) sendSSEError(c, http.StatusInternalServerError, constant.ErrParams.Error()) } go func() { err = chat.ai.SaveMessage(c, &v1.SaveMessage{ ConversationId: int64(atoi), Content: req.Question, SenderType: constant.Eino_User, }) if err != nil { chat.log.Errorf("failed to convert conversation id: %v", err) //sendSSEError(c, http.StatusInternalServerError, constant.ErrParams.Error()) } }() // 执行agent流式处理 stream, err := agent.Stream(c, prompt) if err != nil { chat.log.Errorf("failed to start streaming: %v", err) sendSSEError(c, http.StatusInternalServerError, "failed to start streaming") return } defer stream.Close() // 创建SSE流 flusher, ok := c.Writer.(http.Flusher) if !ok { chat.log.Errorf("streaming not supported") sendSSEError(c, http.StatusInternalServerError, constant.ErrSystem.Error()) return } // 保持连接打开并持续发送消息 for { msg, err := stream.Recv() if err != nil { if errors.Is(err, io.EOF) { // 获取到完整的输出信息 completeResult := collectedContent.String() go func() { // 保存系统输出信息 err := chat.ai.SaveMessage(c, &v1.SaveMessage{ ConversationId: int64(atoi), Content: completeResult, SenderType: constant.Eino_System, }) if err != nil { chat.log.Errorf("failed to save message: %v", err) } // 生成新标题 err = chat.ai.MessageUpdateTitle(c, completeResult, int64(atoi)) if err != nil { chat.log.Errorf("failed to save message: %v", err) } }() sendSSEEvent(c, "end", "StreamCompleted") break } chat.log.Errorf("failed to receive message: %v", err) sendSSEError(c, http.StatusInternalServerError, constant.ErrSystem.Error()) return } // 收集内容 collectedContent.WriteString(msg.Content) // 发送SSE格式的消息 sendSSEEvent(c, "message", msg.Content) flusher.Flush() // 立即发送到客户端 } } // 发送SSE事件 func sendSSEEvent(c *gin.Context, event, data string) { c.SSEvent(event, data) } // 发送SSE错误 func sendSSEError(c *gin.Context, code int, message string) { c.JSON(code, gin.H{"error": message}) }
前端
在cursor与deepseek的双重加持下。 也是快速搭建了一类似deepseek页面风格的前端。不过还是自己手动优化不少内容。 同时AI模型输出结果是markdown格式 也是使用了bytemd来展示markdown内容 具体逻辑就不解释了,项目结构在github中有介绍。axios实例在request包中,service中用于定义接口。 核心对话的业务逻辑在ChatInputComponent.vue文件中




