Eino的基础组件(六)Lambda

深入理解 Eino 框架中最灵活的组件 Lambda,掌握四种交互模式、构造方法、内置工具与最佳实践,让自定义逻辑无缝融入编排流程

🎪 写在前面

在构建 LLM 应用时,我们已经接触了很多 Eino 的标准组件:ChatModel 负责和模型交互,Retriever 负责检索文档,Loader 负责加载数据,Transformer 负责转换文档……这些组件覆盖了大部分常见场景。但总有一些业务逻辑是标准组件无法直接满足的:

  • 从模型输出中提取特定字段(比如从 JSON 字符串解析出结构化数据)
  • 过滤检索结果(比如只保留得分高于阈值的文档)
  • 聚合多个上游的输出(比如合并多路召回的结果)
  • 格式转换(比如把单个元素包装成列表)
  • 自定义的业务规则判断(比如根据用户等级选择不同的处理路径)

这些逻辑虽然简单,但不属于任何标准组件的职责。如果没有一种机制让你嵌入自定义代码,整个编排系统就会变得僵化——要么勉强用标准组件拼凑,要么就得自己实现一个新组件。

Lambda 就是为解决这个问题而生的。它是 Eino 中最灵活、最基础的组件类型,允许你把任意 Go 函数包装成组件,无缝集成到 Graph 或 Chain 中。Lambda 是编排系统的"胶水代码"——它连接各个标准组件,填补逻辑空白,让数据按你的意图流转。

Eino 的 Lambda 设计非常优雅:它不是简单的"函数节点",而是根据输入输出是否为流,定义了四种交互模式(Invoke、Stream、Collect、Transform),并且框架会自动在这些模式之间转换。你可以用最简单的方式定义 Lambda(一个普通函数),也可以精细控制每种模式的行为(实现多个函数)。

这篇文章会从 Lambda 的本质讲起,理解四种交互模式的设计思路,然后通过实战代码掌握各种构造方法,最后探讨内置 Lambda 和最佳实践。如果你想让自己的业务逻辑优雅地融入 Eino 编排,这篇文章是必读。


🧮 Lambda 的本质

Lambda 的核心是一个简单的理念:把 Go 函数包装成 Eino 组件。但"包装"并不是简单的适配器模式,Eino 对 Lambda 做了精心的抽象。

🎮 四种交互模式

Eino 中的所有组件(包括 Lambda)都实现了 Runnable 接口,该接口定义了四种调用方法:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
type Runnable[I, O any] interface {
    // Invoke:非流式输入,非流式输出
    Invoke(ctx context.Context, input I, opts ...Option) (O, error)
    
    // Stream:非流式输入,流式输出
    Stream(ctx context.Context, input I, opts ...Option) (*StreamReader[O], error)
    
    // Collect:流式输入,非流式输出
    Collect(ctx context.Context, input *StreamReader[I], opts ...Option) (O, error)
    
    // Transform:流式输入,流式输出
    Transform(ctx context.Context, input *StreamReader[I], opts ...Option) (*StreamReader[O], error)
}

这四种方法对应了输入输出是否为流的四种组合:

方法输入输出典型场景
Invoke完整数据完整数据数据转换、过滤、聚合
Stream完整数据流式数据生成流式内容(不常用)
Collect流式数据完整数据拼接流式输入(如拼接 LLM 输出)
Transform流式数据流式数据流式转换(如逐 token 处理)

为什么需要四种模式?

在 LLM 应用中,流式处理非常常见。模型可能以流式方式输出 token,你的 Lambda 可能需要:

  • 拼接流:把流式 token 拼成完整文本(Collect)
  • 转换流:逐 token 处理,比如过滤敏感词(Transform)
  • 混合使用:上游是流,但你的 Lambda 只需要完整数据(框架自动 Collect)

Eino 的设计哲学是:组件只需要实现它真正需要的模式,其他模式框架自动转换

🔬 模式详解

Lambda 提供了四种对应的函数类型定义:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
// Invoke 函数:输入输出都是完整数据
type Invoke[I, O, TOption any] func(ctx context.Context, input I, opts ...TOption) (output O, err error)

// Stream 函数:输入完整,输出流式
type Stream[I, O, TOption any] func(ctx context.Context, input I, opts ...TOption) (output *schema.StreamReader[O], err error)

// Collect 函数:输入流式,输出完整
type Collect[I, O, TOption any] func(ctx context.Context, input *schema.StreamReader[I], opts ...TOption) (output O, err error)

// Transform 函数:输入输出都是流式
type Transform[I, O, TOption any] func(ctx context.Context, input *schema.StreamReader[I], opts ...TOption) (output *schema.StreamReader[O], err error)

类型参数说明

  • I:输入类型
  • O:输出类型
  • TOption:自定义选项类型(可选)

自动转换规则

Eino 框架会根据你实现的函数类型,自动生成其他模式。比如:

  • 你只实现了 Invoke,框架会自动提供 Collect(先拼接输入流,再调用 Invoke)
  • 你实现了 InvokeStream,框架会自动提供 Transform(先 Collect 输入,再 Stream 输出)

这意味着大部分情况下,你只需要实现 Invoke 函数,其他模式框架帮你处理。


🏭 构造 Lambda

Eino 提供了多种构造 Lambda 的方法,从最简单的单函数到复杂的多模式实现。

🎲 InvokableLambda

最常用的构造方式,适合 99% 的场景:

1
2
3
4
lambda := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
    // 自定义逻辑
    return strings.ToUpper(input), nil
})

特点

  • 输入输出都是完整数据(非流式)
  • 类型可以是任意自定义类型
  • 框架自动提供其他三种模式

实战示例:提取消息内容

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
extractContent := compose.InvokableLambda(func(ctx context.Context, msg *schema.Message) (string, error) {
    if msg == nil {
        return "", fmt.Errorf("message is nil")
    }
    return msg.Content, nil
})

// 在 Chain 中使用
chain := compose.NewChain[*schema.Message, string]()
chain.
    AppendChatModel(chatModel).  // 输出 *schema.Message
    AppendLambda(extractContent)  // 输入 *schema.Message,输出 string

🌊 StreamableLambda

输入完整数据,输出流式数据:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
lambda := compose.StreamableLambda(func(ctx context.Context, input string) (*schema.StreamReader[string], error) {
    sr, sw := schema.Pipe[string](10)
    
    go func() {
        defer sw.Close()
        // 逐字符发送
        for _, char := range input {
            sw.Send(string(char), nil)
        }
    }()
    
    return sr, nil
})

适用场景

  • 模拟流式输出(测试用)
  • 把批量数据转换成流式处理(不常用)

大部分情况下,你不需要手动实现 StreamableLambda,因为框架会根据 InvokableLambda 自动生成流式输出(把完整结果包装成单次发送的流)。

🧲 CollectableLambda

输入流式数据,输出完整数据:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
lambda := compose.CollectableLambda(func(ctx context.Context, input *schema.StreamReader[string]) (string, error) {
    var result strings.Builder
    
    for {
        chunk, err := input.Recv()
        if err != nil {
            if err == io.EOF {
                break
            }
            return "", err
        }
        result.WriteString(chunk)
    }
    
    return result.String(), nil
})

适用场景

  • 拼接流式输出(如拼接 LLM 的 token)
  • 对流式数据做聚合操作

但实际上,框架已经内置了常见类型的流拼接逻辑(string*schema.Message[]*schema.Message),大部分情况下不需要手动实现 CollectableLambda。

🔀 TransformableLambda

输入输出都是流式数据:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
lambda := compose.TransformableLambda(func(ctx context.Context, input *schema.StreamReader[string]) (*schema.StreamReader[string], error) {
    sr, sw := schema.Pipe[string](10)
    
    go func() {
        defer sw.Close()
        for {
            chunk, err := input.Recv()
            if err != nil {
                if err == io.EOF {
                    break
                }
                sw.Send("", err)
                break
            }
            
            // 转换逻辑:转大写
            sw.Send(strings.ToUpper(chunk), nil)
        }
    }()
    
    return sr, nil
})

适用场景

  • 流式数据转换(如过滤、映射)
  • 需要保持流式特性的中间处理

TransformableLambda 是四种模式中最复杂的,因为你需要手动管理输入输出的两个 StreamReader。

🎚️ 自定义 Option

如果你的 Lambda 需要接受运行时配置,可以定义自定义 Option:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
type MyOptions struct {
    Prefix string
    Suffix string
}

type MyOption func(*MyOptions)

func WithPrefix(prefix string) MyOption {
    return func(o *MyOptions) {
        o.Prefix = prefix
    }
}

// 创建带 Option 的 Lambda
lambda := compose.InvokableLambdaWithOption(
    func(ctx context.Context, input string, opts ...MyOption) (string, error) {
        // 处理 options
        options := &MyOptions{Prefix: "[", Suffix: "]"}
        for _, opt := range opts {
            opt(options)
        }
        
        return options.Prefix + input + options.Suffix, nil
    },
)

// 使用时传入 Option
result, _ := lambda.Invoke(ctx, "hello", WithPrefix("<"), WithSuffix(">"))
// 输出: <hello>

适用场景

  • Lambda 需要根据不同场景调整行为
  • 配置参数较多,不适合写死在函数内部

🎨 AnyLambda

如果你需要精细控制每种模式的行为,可以用 AnyLambda 同时实现多个函数:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
lambda, err := compose.AnyLambda(
    // Invoke 函数
    func(ctx context.Context, input string, opts ...MyOption) (string, error) {
        return "Invoke: " + input, nil
    },
    // Stream 函数
    func(ctx context.Context, input string, opts ...MyOption) (*schema.StreamReader[string], error) {
        sr, sw := schema.Pipe[string](1)
        go func() {
            defer sw.Close()
            sw.Send("Stream: "+input, nil)
        }()
        return sr, nil
    },
    // Collect 函数
    func(ctx context.Context, input *schema.StreamReader[string], opts ...MyOption) (string, error) {
        var result string
        for {
            chunk, err := input.Recv()
            if err != nil {
                break
            }
            result += chunk
        }
        return "Collect: " + result, nil
    },
    // Transform 函数
    func(ctx context.Context, input *schema.StreamReader[string], opts ...MyOption) (*schema.StreamReader[string], error) {
        sr, sw := schema.Pipe[string](10)
        go func() {
            defer sw.Close()
            for {
                chunk, err := input.Recv()
                if err != nil {
                    break
                }
                sw.Send("Transform: "+chunk, nil)
            }
        }()
        return sr, nil
    },
)

适用场景

  • 四种模式的逻辑差异较大,无法通过自动转换实现
  • 需要针对流式和非流式做不同的优化

大部分情况下,InvokableLambda 足够用了。AnyLambda 是高级特性,只在需要极致性能或特殊逻辑时才使用。


💻 实战示例

🛠️ 数据转换

场景:从 Retriever 检索到的文档列表中,只保留得分高于阈值的文档。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
filterByScore := compose.InvokableLambda(func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
    const threshold = 0.7
    
    var filtered []*schema.Document
    for _, doc := range docs {
        score, ok := doc.MetaData["score"].(float64)
        if ok && score >= threshold {
            filtered = append(filtered, doc)
        }
    }
    
    return filtered, nil
})

// 在 Chain 中使用
chain := compose.NewChain[string, []*schema.Document]()
chain.
    AppendRetriever(retriever).   // 输出 []*schema.Document
    AppendLambda(filterByScore)   // 过滤低分文档

💭 消息处理

场景:从 ChatModel 的输出中提取工具调用信息。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
extractToolCalls := compose.InvokableLambda(func(ctx context.Context, msg *schema.Message) ([]string, error) {
    if msg == nil || len(msg.ToolCalls) == 0 {
        return nil, nil
    }
    
    var toolNames []string
    for _, tc := range msg.ToolCalls {
        toolNames = append(toolNames, tc.Function.Name)
    }
    
    return toolNames, nil
})

// 使用示例
chain := compose.NewChain[*schema.Message, []string]()
chain.AppendLambda(extractToolCalls)

toolNames, _ := chain.Compile(ctx).Invoke(ctx, message)
fmt.Println("Model wants to call tools:", toolNames)

🌀 流式处理

场景:从 ChatModel 的流式输出中,过滤掉空字符串 chunk。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
filterEmptyChunks := compose.TransformableLambda(func(ctx context.Context, input *schema.StreamReader[*schema.Message]) (*schema.StreamReader[*schema.Message], error) {
    sr, sw := schema.Pipe[*schema.Message](10)
    
    go func() {
        defer sw.Close()
        for {
            chunk, err := input.Recv()
            if err != nil {
                if err != io.EOF {
                    sw.Send(nil, err)
                }
                break
            }
            
            // 过滤空内容
            if chunk != nil && chunk.Content != "" {
                sw.Send(chunk, nil)
            }
        }
    }()
    
    return sr, nil
})

// 在 Graph 中使用
g := compose.NewGraph[string, *schema.Message]()
_ = g.AddChatModelNode("model", chatModel)
_ = g.AddLambdaNode("filter", filterEmptyChunks)
_ = g.AddEdge("model", "filter")

🔗 编排集成

场景:多路召回后,合并结果并去重。

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
mergeAndDeduplicate := compose.InvokableLambda(func(ctx context.Context, input map[string]any) ([]*schema.Document, error) {
    // input 来自 Parallel 节点,包含多个 Retriever 的结果
    var allDocs []*schema.Document
    
    for key, val := range input {
        docs, ok := val.([]*schema.Document)
        if ok {
            allDocs = append(allDocs, docs...)
        }
    }
    
    // 去重(按 ID)
    seen := make(map[string]bool)
    var unique []*schema.Document
    for _, doc := range allDocs {
        if !seen[doc.ID] {
            seen[doc.ID] = true
            unique = append(unique, doc)
        }
    }
    
    return unique, nil
})

// 在 Chain 中使用
parallel := compose.NewParallel()
parallel.
    AddRetriever("vector", vectorRetriever).
    AddRetriever("keyword", keywordRetriever).
    AddRetriever("rule", ruleRetriever)

chain := compose.NewChain[string, []*schema.Document]()
chain.
    AppendParallel(parallel).       // 多路召回
    AppendLambda(mergeAndDeduplicate)  // 合并去重

🎁 内置 Lambda

Eino 提供了两个开箱即用的 Lambda 工具。

📦 ToList 转换器

作用:把单个元素转换成包含该元素的列表。

1
2
3
4
5
6
7
8
9
// 创建 ToList Lambda
toList := compose.ToList[*schema.Message]()

// 使用场景:ChatModel 返回单个 Message,但下游需要 Message 列表
chain := compose.NewChain[[]*schema.Message, []*schema.Message]()
chain.
    AppendChatModel(chatModel).  // 输出 *schema.Message
    AppendLambda(toList).        // 转换成 []*schema.Message
    AppendDocumentTransformer(transformer)  // 输入 []*schema.Message

类型定义

1
func ToList[T any]() *Lambda

泛型参数 T 是元素类型,返回的 Lambda 输入是 T,输出是 []T

🔎 MessageParser 解析器

作用:从 ChatModel 返回的 Message 中解析 JSON 内容,转换成结构化数据。

使用场景 1:从 Message.Content 解析

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
// 定义目标结构体
type IntentResult struct {
    Intent string `json:"intent"`
    Confidence float64 `json:"confidence"`
}

// 创建解析器
parser := schema.NewMessageJSONParser[*IntentResult](&schema.MessageJSONParseConfig{
    ParseFrom:    schema.MessageParseFromContent,  // 从 Content 字段解析
    ParseKeyPath: "",  // 解析整个 Content
})

// 包装成 Lambda
parserLambda := compose.MessageParser(parser)

// 在 Chain 中使用
chain := compose.NewChain[*schema.Message, *IntentResult]()
chain.AppendLambda(parserLambda)

// 执行
message := &schema.Message{
    Content: `{"intent": "book_flight", "confidence": 0.95}`,
}
result, _ := chain.Compile(ctx).Invoke(ctx, message)
fmt.Println("Intent:", result.Intent)  // 输出: book_flight

使用场景 2:从 ToolCall 结果解析

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
// 创建解析器(从 ToolCall 解析)
parser := schema.NewMessageJSONParser[*UserInfo](&schema.MessageJSONParseConfig{
    ParseFrom: schema.MessageParseFromToolCall,  // 从 ToolCall 结果解析
})

parserLambda := compose.MessageParser(parser)

// 使用示例
message := &schema.Message{
    ToolCalls: []*schema.ToolCall{
        {
            Function: schema.FunctionCall{
                Name:      "get_user_info",
                Arguments: `{"name": "Alice", "age": 30}`,
            },
        },
    },
}

userInfo, _ := parserLambda.Invoke(ctx, message)

配置说明

  • ParseFrom:指定从哪里解析
    • MessageParseFromContent:从 Message.Content 解析
    • MessageParseFromToolCall:从 Message.ToolCalls[0].Function.Arguments 解析
  • ParseKeyPath:如果只需要解析 JSON 中的某个字段,用点号分隔的路径(如 "data.result"

MessageParser 在意图识别、结构化输出等场景非常有用,避免手动解析 JSON 的重复代码。


✨ 最佳实践

📍 使用场景

何时使用 Lambda?

  1. 数据格式转换

    • 标准组件之间类型不匹配时(如单个 Message 转 Message 列表)
    • 提取组件输出的特定字段
  2. 过滤与筛选

    • 根据业务规则过滤检索结果
    • 去除无效数据
  3. 聚合与合并

    • 多路召回后合并结果
    • 拼接多个组件的输出
  4. 自定义业务逻辑

    • 计算得分
    • 规则判断
    • 日志记录

何时不用 Lambda?

  • 如果逻辑可以通过标准组件实现,优先用标准组件(如文档分割用 Transformer,不要用 Lambda)
  • 如果逻辑非常复杂且可复用,考虑封装成独立的自定义组件

⛔ 注意事项

1. 避免阻塞操作

Lambda 函数在 Graph 的执行流程中是同步调用的,不要在 Lambda 中做长时间的阻塞操作(如网络请求、大量计算)。如果必须做,考虑:

  • 在 Lambda 外部做异步处理
  • 或者把耗时操作封装成独立的组件

2. 错误处理

Lambda 的错误会中断整个编排流程,务必做好错误处理:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
safeLambda := compose.InvokableLambda(func(ctx context.Context, input string) (string, error) {
    defer func() {
        if r := recover(); r != nil {
            log.Printf("Lambda panic: %v", r)
        }
    }()
    
    // 业务逻辑
    return process(input)
})

3. 类型安全

Lambda 的输入输出类型必须与上下游节点对齐,编译期就会检查。但如果使用 any 类型,要在运行时做类型断言:

1
2
3
4
5
6
7
lambda := compose.InvokableLambda(func(ctx context.Context, input any) (string, error) {
    str, ok := input.(string)
    if !ok {
        return "", fmt.Errorf("expected string, got %T", input)
    }
    return str, nil
})

4. 流式处理的 goroutine 管理

在 TransformableLambda 或 StreamableLambda 中,记得在 goroutine 中 defer sw.Close(),确保流正确关闭。

🚄 性能优化

1. 优先使用 InvokableLambda

除非你真的需要流式处理,否则用 InvokableLambda。框架的自动转换已经足够高效,手动实现四种模式反而容易出错。

2. 避免不必要的拷贝

Lambda 的输入输出都是引用传递(指针、切片、map),不要无谓地拷贝大对象:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
// ❌ 不好:拷贝整个切片
badLambda := compose.InvokableLambda(func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
    result := make([]*schema.Document, len(docs))
    copy(result, docs)
    return result, nil
})

// ✅ 好:直接返回
goodLambda := compose.InvokableLambda(func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
    return docs, nil
})

3. 批量处理优于逐个处理

如果可能,尽量一次处理多个元素,而不是逐个调用 Lambda:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
// ✅ 好:批量处理
batchLambda := compose.InvokableLambda(func(ctx context.Context, docs []*schema.Document) ([]*schema.Document, error) {
    // 一次过滤整个列表
    return filterDocs(docs), nil
})

// ❌ 不好:在外层循环中多次调用 Invoke
for _, doc := range docs {
    singleLambda.Invoke(ctx, doc)  // 每次调用都有开销
}

📝 总结

Lambda 是 Eino 编排系统中最灵活的组件,它让你能够把任意自定义逻辑无缝集成到 Graph 或 Chain 中。核心要点:

  1. 四种交互模式:Invoke、Stream、Collect、Transform,对应输入输出是否为流的四种组合
  2. 自动转换:大部分情况只需要实现 InvokableLambda,框架自动提供其他模式
  3. 灵活构造:从简单的单函数到复杂的多模式实现,满足不同需求
  4. 内置工具:ToList 和 MessageParser 解决常见场景,开箱即用
  5. 最佳实践:优先用标准组件,Lambda 作为"胶水"填补空白;注意错误处理和性能优化

掌握 Lambda 的使用,你就能让 Eino 的编排系统真正为你的业务服务——不再受限于标准组件,而是可以自由定义数据流转的每一个环节。无论是简单的格式转换,还是复杂的业务规则,Lambda 都能优雅地完成任务。

下一篇文章,我们会探讨 Eino 的流式处理机制和 Callback 系统,看看如何在编排中实现实时反馈和监控。


参考资料

最后更新于 2026-09-02 17:26 UTC
그 경기 끝나고 좀 멍하기 있었는데 여러분 이제 살면서 여러가
使用 Hugo 构建
主题 StackJimmy 设计