api/rag-pipeline/RAG-Pipeline开发规范.md
所有流水线应在 /home/liuzimu/code/KM/service/rag_job_engine.go 文件的 registerDefaultPipelines 函数中注册。
func registerDefaultPipelines() {
// 注册流水线
pipelines.RegisterPipeline("流水线类型", func() pipelines.Pipeline {
return pipelines.New流水线名WithDB(model.DB)
})
}
// 注册RechunkAndReindex流水线
pipelines.RegisterPipeline("rechunk_and_reindex", func() pipelines.Pipeline {
return pipelines.NewRechunkAndReindexPipelineWithDB(model.DB)
})
WithDB 后缀)使用 /home/liuzimu/code/KM/service/rag_job_engine.go 文件中的 CreateJob 函数创建任务:
func CreateJob(ctx context.Context, eid int64, jobType string, startParameters string) (*model.RagJob, error)
ctx: 上下文eid: 企业IDjobType: 流水线类型(与注册时使用的类型一致)startParameters: 流水线参数的JSON字符串// 准备参数
params := RechunkAndReindexParameters{
Eid: eid,
FileID: fileID,
UserID: userID,
}
paramsJSON, _ := json.Marshal(params)
// 创建任务
job, err := CreateJob(ctx, eid, "rechunk_and_reindex", string(paramsJSON))
if err != nil {
return err
}
每个流水线应定义自己的参数结构体,例如:
type RechunkAndReindexParameters struct {
Eid int64 `json:"eid"`
FileID int64 `json:"file_id"`
UserID int64 `json:"user_id"`
}
type ReindexParameters struct {
Eid int64 `json:"eid"`
FileID int64 `json:"file_id"`
UserID int64 `json:"user_id"`
RunAIIndexTask bool `json:"run_ai_index_task"` // 是否运行AI索引任务,默认为false
}
type PipelineParameters struct {
Eid int64 `json:"eid"`
FileID int64 `json:"file_id"`
UserID int64 `json:"user_id"`
RunAIIndexTask bool `json:"run_ai_index_task"`
}
流水线文件应放在 /home/liuzimu/code/KM/rag-pipeline/pipelines/ 目录下。
文件名应使用下划线分隔的小写字母,例如:rechunk_and_reindex.go
type 流水线名Pipeline struct {
*BasePipeline
}
必须提供两个构造函数:
// 不带数据库连接的构造函数
func New流水线名Pipeline() Pipeline {
pipeline := &流水线名Pipeline{
BasePipeline: NewBasePipeline("流水线类型"),
}
// 注册流水线
RegisterPipeline("流水线类型", func() Pipeline {
return New流水线名Pipeline()
})
return pipeline
}
// 带数据库连接的构造函数
func New流水线名PipelineWithDB(db *gorm.DB) Pipeline {
pipeline := &流水线名Pipeline{
BasePipeline: NewBasePipelineWithDB("流水线类型", db),
}
return pipeline
}
在 Initialize 方法中添加流水线所需的步骤:
func (p *流水线名Pipeline) Initialize() error {
// 添加步骤
step := steps.New步骤名Step(p.DB)
if err := p.AddStep("步骤名", step); err != nil {
return err
}
// 添加更多步骤...
return nil
}
根据步骤顺序准备对应的参数:
func (p *流水线名Pipeline) PrepareStepParameters(order int) interface{} {
// 从上下文中获取参数
var eid, fileID, userID int64
// 获取参数逻辑...
// 根据步骤顺序返回对应的参数
switch order {
case 1:
return steps.步骤名Parameters{
Eid: eid,
FileID: fileID,
UserID: userID,
}
// 更多步骤...
default:
return nil
}
}
实现流水线的执行逻辑:
func (p *流水线名Pipeline) Execute(job *model.RagJob) error {
// 解析参数
var params 流水线名Parameters
if job.StartParameters != "" {
if err := json.Unmarshal([]byte(job.StartParameters), ¶ms); err != nil {
return fmt.Errorf("failed to unmarshal parameters: %v", err)
}
}
// 将参数添加到上下文中
p.Context["eid"] = params.Eid
p.Context["file_id"] = params.FileID
p.Context["user_id"] = params.UserID
// 如果尚未初始化,则初始化流水线
if len(p.Steps) == 0 {
if err := p.Initialize(); err != nil {
job.Status = model.RagJobStatusFailed
job.FailureReason = err.Error()
return err
}
}
// 使用自身作为执行器调用基础执行方法
return p.BasePipeline.ExecuteWithExecutor(job, p)
}
步骤文件应放在 /home/liuzimu/code/KM/rag-pipeline/steps/ 目录下。
文件名应使用下划线分隔的小写字母,例如:cleanup.go
type 步骤名Step struct {
BaseStep
DB *gorm.DB
}
type 步骤名Parameters struct {
Eid int64 `json:"eid"`
FileID int64 `json:"file_id"`
UserID int64 `json:"user_id"`
// 其他参数...
}
type 步骤名Result struct {
// 结果字段...
Success bool `json:"success"`
}
func New步骤名Step(db *gorm.DB) *步骤名Step {
return &步骤名Step{
DB: db,
}
}
实现步骤的执行逻辑:
func (s *步骤名Step) Execute(parameters any) error {
// 开始处理
s.Step.StartProcessing(parameters)
// 类型断言获取参数
params, ok := parameters.(步骤名Parameters)
if !ok {
err := fmt.Errorf("invalid parameters type, expected 步骤名Parameters")
s.Step.CompleteWithError(err.Error())
return err
}
// 执行步骤逻辑...
// 创建结果
result := 步骤名Result{
// 设置结果字段...
Success: true,
}
// 完成步骤并返回结果
s.Step.CompleteSuccessfully(result)
return nil
}
在步骤执行过程中,如果发生错误,应该:
s.Step.CompleteWithError(errMsg) 标记步骤失败if err != nil {
errMsg := fmt.Sprintf("操作失败: %v", err)
s.Step.CompleteWithError(errMsg)
return fmt.Errorf(errMsg)
}
在关键操作处添加适当的日志记录,便于调试和监控。
对于涉及多个数据库操作的场景,应使用事务确保数据一致性。
使用流水线的 Context 字段在步骤之间传递数据。
如果需要并行执行某些步骤,可以使用 ParallelGroups 字段定义并行组。
参考 rechunk_and_reindex.go 文件,这是一个完整的流水线实现示例。
参考 cleanup.go 文件,这是一个完整的步骤实现示例。
这份规范文档基于 rag-pipeline 的现有代码实现,涵盖了流水线开发的主要方面。遵循这些规范可以确保代码的一致性和可维护性。