智能协同云图库项目 go后端 + kratos框架 + 前端elementplus

智能云图库go版本

欢迎各位大佬批评指出问题,作者只不过换了个架构实现了项目,并不是完美的,可能有些方案实现的不理想,bug难免有些没有排查出来,欢迎各位指导批评。

本文章只介绍技术方案选型,业务逻辑鱼皮大大已经介绍的非常清楚了。我们不在过多赘述。 github作者源码

如果你是一个go小白,没听说过casbin,kratos。只用过MVC架构的gin框架之类的。那么你应该好好学学该文章涉及到的框架和技术。

同时作为go语言的特色goroutine和channel是每个go学者几乎必备的技能,如果你不理解,那么看懂作者的websocket模块基本理解了。注意了websocket模块代码有点多,作者将代码折叠了,要注意

不过自学的话确实不容易,作者刚接触也是花了将近一个星期才弄懂了kratos,花了一天半才基本摸清楚了casbin。(真希望有个当初能有个大佬能带带作者 /(ㄒoㄒ)/~~)

go-kratos框架

kratos

kratos微服务框架,领域驱动设计,由bilibili开发(不过现在好像已经成为独立框架了),当然也完全适用单体架构。框架整体的结构,设计理念官网非常详细,就像教科书一般。 详情见kratos官网 kratos框架可以说是非常理想的框架,至今仍在更新维护。国内很多中小厂都采用了该框架。

框架特性

  • 遵循 RESTful API 设计规范 & 基于接口的编程规范
  • 基于 Kratos 2.x 框架(支持微服务架构).
  • 基于 Casbin 的 RBAC 访问控制模型
  • 基于 gorm/gen 的数据库存储,可自行生成model与curd方法
  • 基于 WIRE 的依赖注入 -- 依赖注入本身的作用是解决了各个模块间层级依赖繁琐的初始化过程
  • 基于 Zap & Context 实现了日志输出,通过结合 Context 实现了统一的 TraceID/UserID 等关键字段的输出(同时支持日志钩子写入到Gorm)
  • 基于 JWT 的用户认证 -- 基于 JWT 的黑名单验证机制
  • 基于 Swaggo 自动生成 Swagger 文档 -- 独立于接口的 mock 实现
  • 基于 go mod 的依赖管理(国内源可使用:https://goproxy.cn/)
  • 类DDD的领域驱动设计

框架结构

下面是该项目的架构,其实和官网介绍的架构一致。 详情见官网kratos项目架构

  1. api:接口层,放置了proto文件,定义了api的接口,而这些接口由internal的service实现。
  2. cmd:项目的启动层或者入口层,通常存放了main函数以及wire依赖
  3. configs:配置层,存放各种配置文件,比如yaml配置文件和casbin配置文件
  4. internal:业务逻辑层,该层尤为重要,还会单独讲接。该层名为internal,在go中是一个特殊包名,放在 internal 目录下的包,只能被当前模块(module)或其直接父级目录下的包导入,外部项目无法导入这些包。
  5. pkg:公共包,一般用于存放工具类和公共模块,比如枚举值
  6. sql: 该包存放了该项目的sql文件,个人创建,可以无视
  7. test:测试包
  8. third_part:第三方部分,一般存放了google的proto文件
  1. biz 作为中间层,类似于DDD的domain,是连接service和data一起沟通其他服务的桥梁
  2. conf配置层的实体,在这里创建了配置文件具体的结构体,用于将配置文件中的具体内容变为结构体的方式,方便各个模块使用。其中该实体类由proto文件生成
  3. data仓储层,定义了访问数据库的方式
  4. pkg公共模块,该模块是真正意义上存放业务的公共模块,一般存放了oss,ai,以及各种中间件
  5. server:http和grpc实例在这里生成和创建
  6. service:DDD架构的appliaciton层,实现了api定义的服务层

API设计

API设计采用google规范(作者其实也没有完全采用)。详情见googleAPI设计指南中文网

gorm/gen

使用gorm/gen框架,自动生成结构体与查询对象,更加简练简单的去实现CURD。底层为gorm。gen主张极度的安全,所有有时候稍微复杂的sql不好写。不过使用gen前提是有gorm所以我们完全可以gen和gorm组合使用 gorm官网 推荐一篇包含基础CURD中文博客

gen

帮助我们自动生成代码,包括快速增删改查的方法,以及库表映射结构体。但是生成代码需要我们主动去main函数中运行。就像下方的使用案例,当我需要重新刷新model和query(就是在重新生成一遍代码)时,只需要修改package name为main直接跑一边就行。(就想java中mybatis需要我们手动生成dao层mapper和数据库表结构映射类一样)

text
复制代码
package data //将此处data改为main即可执行 import ( "fmt" "gorm.io/driver/mysql" "gorm.io/gen" "gorm.io/gorm" ) const MysqlConfig = "root:123456@tcp(127.0.0.1:3306)/c_picture?charset=utf8mb4&parseTime=True&loc=Local" func main() { // 连接数据库 db, err := gorm.Open(mysql.Open(MysqlConfig)) if err != nil { panic(fmt.Errorf("cannot establish db connection: %w", err)) } // 生成实例 g := gen.NewGenerator(gen.Config{ // 相对执行`go run`时的路径, 会自动创建目录 // 如果使用go run会以当前目录为起点 如果编辑器运行会以项目工程目录(根目录)为起点 OutPath: "internal/data/dal/query", // 查询代码目录 ModelPkgPath: "/model", // 模型代码目录 模型代码路径自动添加了一个前缀 ../OutPath 使用时注意 // WithDefaultQuery 生成默认查询结构体(作为全局变量使用), 即`Q`结构体和其字段(各表模型) // WithoutContext 生成没有context调用限制的代码供查询 // WithQueryInterface 生成interface形式的查询代码(可导出), 如`Where()`方法返回的就是一个可导出的接口类型 Mode: gen.WithDefaultQuery | gen.WithQueryInterface, }) g.ApplyBasic() // 设置目标 db g.UseDB(db) // 生成所有表 g.GenerateAllTable() // 生成所有表的查询代码 g.ApplyBasic( g.GenerateAllTable()...) // 一般情况下此处无需使用 总用其他解决方案 // g.ApplyBasic(g.GenerateModel("picture", // gen.FieldType("review_at", "*time.Time"), // 让 review_at 支持插入null // )) g.Execute() }

生成后的代码如下 转存失败,建议直接上传图片文件

gorm有个问题,当我们更新字段时,如果结构体中有些值是默认值,如果我们没有传值也没有指定更新具体字段,直接去更新表结构,会导致这些字段被默认值替换(当然gorm也有忽略默认值字段的update方法),此时我们可以用指针类型来代替值类型,传入一个空指针的话,gorm会字段识别为null,不插入。具体需改方案如下,但是我们也可以无视此方案,完全可以和mybatis一样根据值是否为空或者是否为字段类型默认值去选择字段是否要更新

text
复制代码
g.ApplyBasic(g.GenerateModel("picture", gen.FieldType("review_at", "*time.Time"), // 让 review_at 支持插入null ))

gorm

gorm就和mybatis一样了,用起来都大差不差。 注意的是删除操作如果数据库中存在一条delete_at字段,那么默认执行逻辑删除,否则执行硬删除

这里给两个案例分页查询和修改

text
复制代码
func (s *sysUserRepo) ListPage(ctx context.Context, page, size int32, condition biz.UserListCondition) ([]*model.SysUser, error) { m := s.data.Query(ctx).SysUser q := m.WithContext(ctx) if condition.Status != 0 { q = q.Where(m.Status.Eq(condition.Status)) } if condition.UserName != "" { q = q.Where(m.Username.Like("%" + condition.UserName + "%")) } if condition.RoleId != 0 { q = q.Where(m.RoleID.Eq(int64(condition.RoleId))) } limit, offset := pagination.ConvertPageSize(page, size) return q.Limit(limit).Offset(offset).Find() } func (s *sysUserRepo) UpdateByID(ctx context.Context, id int64, user *model.SysUser) error { if id == 0 { return errors.New("user can not update without id") } q := s.data.Query(ctx).SysUser // updates 支持使用结构体去修改 update的话需要单个字段的修改 _, err := q.WithContext(ctx).Select(q.Avatar, q.NickName).Where(q.UID.Eq(id)).Updates(user) return err }

gorm和gen的一些方法注意事项

联级查询以及scan

注意join以及scan方法 如果涉及到复杂的查询,例如我们查询角色ID为1的用户信息,左联角色信息。因为没有对应的结构体字段,所以我们要自己定义,不过也不难,直接吧两个gen生成的两个表合并起来即可。最后一定要用scan把要获取的信息映射到里面。 这里面还要注意,scan函数底层调用的是反射来对传入的指针进行修改,这里我们直接var声明变量后传入指针即可,如果使用make会导致切片长度固定,无法通过反射来进行修改,如reflect.Value.SetLen,就会报错。gorm里join就是内联的意思,如果找不到inner用join即可

text
复制代码
type Result struct { model.SysRole model.SysUser } func (r *sysRoleRepo) SelectUserMessageAndRole(ctx context.Context) error { ro := r.data.Query(ctx).SysRole user := r.data.Query(ctx).SysUser var result []Result err := user.WithContext(ctx).Select(ro.ALL, user.ALL).LeftJoin(ro, ro.Rid.EqCol(user.RoleID)).Where(ro.Rid.Eq(1)).Scan(&result) for _, sysUser := range result { r.log.Warnf("%+v", sysUser) } return err }
修改updates

gorm提供了多个修改方法update根据设置的字段更新,updates根据给的结构体去更新,自动忽略结构体默认值,但是这个方法要注意他会自动的去添加where条件查询,条件是该表中的你设置的主键。有些时候这样会导致你的无法正常更新(例如当你的部分主键是可以发生变化的时候)。这时候需要UpdateColumns方法,他会只根据你的条件以及给的结构体内容去更新。就像我们正常自己写的sql一样。不过该方法有个问题就是不会忽略结构体默认值,建议改之前find一下,拿到旧值改完在update

gorm原生查询

上方代码都是gen生成的快速查询方案,该方案特点就是适用于全部的数据库,mysql,GuassDB,Oracle。因此封装的快速crud方案以及函数都是适用于公共查询的 。如果此时我们需要一些单独的查询方案,例如mysql中独特的YEARWEEK(date)是其他查询没有的。这样的话我们就用原生查询。具体情况看业务如何。gen和grom原生查询混在一起因为类型不同难以处理,可能会使代码变得复杂不优雅。建议要么全用gen要么全用gorm,不过还是得看业务逻辑

text
复制代码
s.data.db.Model(&model.Picture{}).WithContext(ctx).Select("YEARWEEK(create_at) as period,count(id) as count"). Group("period").Order("period desc").Scan(&spaceUser)
gen group by传参问题

gen想要使用DATEFORMAT函数来作为group参数的时候的时候你会发现无法正常传入参数值,会生成下面这种sql。不仅仅是DATE_FORMAT函数只要是传参都不行,那么为了解决这个问题可以通过起别名的问题来解决,详情见下个标题。(花了两个小时终于在GitHub官方issue中找到问题解决方案了)

text
复制代码
DATE_FORMAT(`test`.`created_at`,?)
gen创建别名字段(方便的起别名)

有时遇到稍微复杂的sql难以用gen生成。直接手动创建别名可以通过NewField一个新字段。

go
复制代码
period := field.NewField("", "period") // 第一个参数为表名 起别名的话不要填写 q = q.Select(m.CreatedAt.DateFormat("%Y-%m-%d").As("period"), m.ID.Count().As("count")). Group(period). Order(period.Asc())

Casbin

Casbin作为一个权限管理框架,支持多种语言,多种模型以及模型的扩展,生态好,功能强大。这里我只涉及一些介绍和扩展。详细使用请看官方文档。kratos集成请看kratos-casbin该库使用很方便,与kratos高度集成,一站式搭建流程。 Casbin中文官方 使用:

  1. 添加模型
  2. 数据库存放策略
  3. 制定角色表,用于指定角色具有的权限
  4. 制作成中间件管理即可。

其中要注意的是:

  1. 用了grpc框架后接口路径我们要用grpc方法来代替。(这个要求和我们配置的模型文件有关,想要存放具体接口路径修改模型文件即可)
  2. 只要为casbin添加了gorm数据库,casbin就能自己去创建对应的规则存储库。无需我们操作。
  3. casbin支持多种模型但我们使用市面上绝大多数的模型。RBAC。详细的RBAC介绍以及RBAC扩展见下文
  4. casbin继承后几乎不需要我们更改代码,需要扩展功能往往只需要修改casbin的模型配置文件,和规则存储库(存储规则的库表)即可。

模型

我们要创建一个conf文件放在配置文件中,作为模型文件供casbin使用

text
复制代码
# PERM模型 # sub 访问实体(想要访问资源的用户) # obj访问的资源(将被访问的资源) # act访问方法(用户对资源执行的操作) # eft:策略结果,一般为空,默认指定allow # 请求格式 例如下方 [request_definition] # 指定请求格式为 主体 对象 操作 r = sub, obj, act # 策略 在grpc框架中我们如果使用的话 p要保存为grpc定义的方法名 /package.Service/Method [policy_definition] p = sub, obj, act # policy生效范围 [policy_effect] # 结果是allow就放行 e = some(where (p.eft == allow)) # 匹配规则 支持 * 和 : 通配(如 /api/:id /api/*) [matchers] m = r.sub == p.sub && (keyMatch2(r.obj, p.obj) || keyMatch(r.obj, p.obj)) && ( r.act == p.act || p.act == '*')

Casbin RBAC权限介绍

RBAC(Role-Based Access Control,基于角色的访问控制)。一个用户一个角色,不同的角色有着不同的权限,casbin自建规则库要求字段是这样的(casbin自己建的哦) 转存失败,建议直接上传图片文件 接下来我们要了解这些字段的意义

  • ptype:策略类型,表示当前的策略类型,一般三种,当然还能往下扩展g3,g4等等
    • p:权限策略,定义 角色/资源操作权限
    • g:角色继承策略,为了我们方便扩展角色权限,我们可以让新的角色继承老的角色的全部权限
    • g2:扩展分组,就是在继承一层。例如g继承了p,g2继承了g...
  • v0:主题
  • v1:资源
  • v2:操作
  • v3往后均为扩展。
顺序发生变化,那么对应的v0 v1 v2也要发生变化。

那你要问了,怎么用呢?直接看一条数据。

上面是两种数据填写方式,意思是角色admin有着UploadFileService的全部接口的全部权限。第二行就是只有/api.....这个接口的POST权限。

角色继承

假如我们创建一个新角色,但是这个角色比普通用户只是多一点接口权限而已,重新建一个角色引入一便接口太过于麻烦,那么我们只需要让该角色继承普通用户,在扩展一点接口就可以了。继承方法见下方。让super_user继承了user,superdouble继承了user,再往下同理。 转存失败,建议直接上传图片文件

注意如果确实开启了继承功能要在配置文件中这样添加配置

text
复制代码
# PERM模型 # sub 访问实体(想要访问资源的用户) # obj访问的资源(将被访问的资源) # act访问方法(用户对资源执行的操作) # eft:策略结果,一般为空,默认指定allow # 请求格式 例如下方 [request_definition] # 指定请求格式为 (主体 对象 操作)这样复合grpc方法名 r = sub, obj, act # 策略 在grpc框架中我们如果使用的话 p要保存为grpc定义的方法名 /package.Service/Method [policy_definition] p = sub, obj, act # 角色定义 角色继承 如果只想要一种角色即可 注释此部分代码 [role_definition] g = _, _ # 用户-角色关系 g2 = _, _ # 角色-角色继承关系 # policy生效范围 [policy_effect] # 结果是allow就放行 e = some(where (p.eft == allow)) # 匹配规则 支持 * 和 : 通配(如 /api/:id /api/*) [matchers] m = r.sub == p.sub && (keyMatch2(r.obj, p.obj) || keyMatch(r.obj, p.obj)) && ( r.act == p.act || p.act == '*') && (g(r.sub, p.sub) || g2(r.sub, p.sub))

多租户

多租户是RBAC的扩展,定义四个模型:用户,域,角色,接口。 意思在这个域中的角色的接口权限。对于本项目有来说我们的权限如果不想自己手动校验,那么使用多租户是件很方便的事情,不同用户在不用域中的角色不同,可以自动帮我们校验好。 例如在后面的公共空间模块,不同的角色在不同的空间需要做不同的校验,如果使用多租户系统会很方便。但是如果结合我们的项目来看,其实单纯的继承+判断用户是否在该空间,即可完成。两种方案均可。 但是每次请求如何传递用户的域呢?我们不能保存到token中吧,每次请求的token不一致或者频繁的更换token肯定是不理想的方案。 我们可以放到请求头中,前端发送请求的时候就可以发送不同域。这样就能完成多租户功能了 当然在本项目中开启和关闭多租户只需要在中间件中注释掉支持多租户即可

text
复制代码
// Auth 认证中间件 func Auth(s *conf.Auth, repo biz.CasbinRuleRepo) middleware.Middleware { // selector 加载多个中间件 return selector.Server( // JWT 认证部分 jwt.Server源码中默认配置令牌格式为 Authorization: Bearer <token> jwt.Server( // 该函数返回配置中的 UserSecretKey func(token *jwtV5.Token) (interface{}, error) { return []byte(s.User.UserSecretKey), nil }, // 签名方法 jwt.WithSigningMethod(jwtV5.SigningMethodHS256), // 自定义 claims 类型为 TokenClaims jwt.WithClaims(func() jwtV5.Claims { return &authz.TokenClaims{} }), ), // casbin 权限校验 casbin.Server( // 获取casbin 配置模型 casbin.WithCasbinModel(repo.GetModel()), // 获取casbin 权限规则 casbin.WithCasbinPolicy(repo.GetAdapter()), // 创建安全用户 casbin.WithSecurityUserCreator(authz.NewSecurityUser), // 启用自动加载策略,每30秒刷新一次 casbin.WithAutoLoadPolicy(true, 30*time.Second), //casbin.WithDomainSupport(), // 启用多租户支持 ), // 添加白名单 ).Match(AuthWhiteListMatcher()).Build() }

获取到所有接口信息的方案

获取到所有接口信息我们就能入库保存,保存后就能配合casbin一起使用了

方案一 (不是很推荐)

grpc获取接口信息,但是缺点是无法获取到具体的请求方法,毕竟grpc没有请求方法。这个问题稍微严重一点。其中grpc就是*grpc.Server对象,来自"github.com/go-kratos/kratos/v2/transport/grpc"

text
复制代码
// 编辑解析server结构 for key, info := range a.grpc.Server.GetServiceInfo() { a.log.Infof("key:%s,info:%v\n", key, info) if strings.HasPrefix(key, "api") { // 匹配到api开头的接口数据 split := strings.Split(key, ".") for _, method := range info.Methods { sysApis = append(sysApis, &model.SysAPI{ ID: 0, APIGroup: split[len(split)-1], Path: fmt.Sprintf("/%s/%s", key, method.Name), Description: fmt.Sprintf("操作%s接口", split[len(split)-1]), UpdatedAt: now, CreatedAt: now, Method: "*", }) } } }
方案二(推荐)

我们利用proto生成的openapi.json文件,解析文件,从文件中能拿到大量的信息包括接口路径,操作ID,接口注释,请求方法,接收参数以及返回参数。极其好用。唯一小小缺点就是不会拿到操作全称,不过毫无影响,前缀来自包名,主要我们不动项目包结构就完全没影响(都动包结构了肯定得重新部署,所以可以说是毫无影响),我们手动拼接即可。 举例:某个操作全程 /api.go_web.v1.RolesService/RolesTest 但是通过解析文件拿到的:RolesService_RolesTest 首先导入库,这个库可以帮助我们解析yaml json yml等类型的配置文件,超级好用。

text
复制代码
go get github.com/getkin/kin-openapi/openapi3

使用方法

text
复制代码
package main import ( "context" "fmt" "net/http" "github.com/getkin/kin-openapi/openapi3" ) func main() { background := context.Background() // 创建对象 loader := openapi3.NewLoader() // 通过地址获取文件 也允许使用二进制等等 doc, err := loader.LoadFromFile("api/go_web/v1/openapi.yaml") if err != nil { panic(err) } // 校验文件是否复合规范 err = doc.Validate(background) if err != nil { panic(err) } // 遍历所有接口是否包含这个四个类型的信息 因为接口路径可能有重复的 但是接口方法不重复 所以我们都要遍历一遍看看 method := []string{ http.MethodPost, http.MethodGet, http.MethodDelete, http.MethodPut, } // 通过路径获取到接口内容 for _, path := range doc.Paths.InMatchingOrder() { // 获取接口内容对象 pathItem := doc.Paths.Value(path) // 获取对象操作 operations := pathItem.Operations() for _, s := range method { operation, ok := operations[s] if !ok { continue } fmt.Printf("%s ", path) // 路径 fmt.Printf("%s ", s) // 方法 fmt.Printf("%s ", operation.Description) // 描述 fmt.Println(operation.OperationID) // 操作id } } }

casbin缺点

任何框架都不是万能的,casbin也有缺点,可能这也是哪些中厂大佬们不愿意使用的原因吧。

  1. casbin的数据会存放到内存,适配器很容易全量规则load(把全部的权限都读取存放到内存)
  2. 如果规则达到10w条后用起来需要非常谨慎,多节点部署的话,如果监听规则变更断了还会导致多个节点之间的权限不统一。
  3. casbin在少角色的系统中用起来确实方便,但是涉及到动态角色变化,多租户等等复杂情况,表存储容量就会急剧上升,在微服务系统中不是很友好。

auth权限校验框架

使用go自带的jwt5封装,封装为kratos中间件使用。其中auth模块和casbin配套使用,进一步增强以及方便我们对于权限的管理。

中间件回调

在grpc自动注册http路由会遇到一些问题,就是某些数据无法通过grpc进行传递,例如form-data格式数据,此时我们无奈只能单开路由,单开路由但是还是会碰到一个问题:无法通过中间件,此时我们需要中间件回调。重新走一遍中间件。

text
复制代码
r := srv.Route("/") v1group := r.Group("/v1") v1group.POST("/file/upload", func(ctx http.Context) error { http.SetOperation(ctx, "/api.go_web.v1.FileLoad/UploadFile") var respon v1.FileLoadResponse // 中间件回调 会影响响应时长 很小的可能会造成context超时 h := ctx.Middleware(func(ctx context.Context, req interface{}) (interface{}, error) { return file.FileLoad(ctx, nil) }) out, err := h(ctx, &respon) if err != nil { return err } reply := out.(*v1.FileLoadResponse) return ctx.Result(200, reply) })

GRPC

go学者必须掌握的一项技术,go的接口甚至一些类型结构体全都由grpc定义。kratos脚手架可以帮助我们将proto文件生成为go文件,但是麻烦。 由此我们还得掌握另一个技术buf。 在项目根目录定义好buf.yaml以及buf.gen.yaml两个配置文件,使用buf命令去编辑proto文件。 此外还得掌握新的校验工具protovalidate,官方开始支持全面使用该校验器,(确实好用你别说),该工具支持在grpc中编写对参数的校验规则。一个proto文件走天下了 buf官方 protovalidate校验器校验规则

buf使用方法

首先在你的首页创建两个文件buf.gen.yaml buf.yaml。如下是已经配置好的文件,此后我们只需要在控制台输入buf generate就会自动翻新全部的proto文件。简单快速方便。而且可以将所有的接口统一到一个openapi.yaml文档中。buf脚手架可以去官方了解和下载

  • buf.gen.yaml
text
复制代码
version: v2 managed: enabled: true # override: # 强制使用以下定义的配置而不是proto文件原配置 # - file_option: go_package_prefix # value: api/go_web/v1 # inputs: # - directory: api/go_web/v1 # disable: # - file_option: go_package_prefix # module: google/api/annotations.proto plugins: - local: protoc-gen-openapi # swagger out: api/go_web/v1 opt: - naming=proto # 指定文件类型 - depth=3 # 限制递归解析嵌套消息的深度 - default_response=false # 自动添加默认错误响应 - enum_type=string # 控制枚举值的 OpenAPI 表示形式 - output_mode=source_relative # 输出为每个proto一个文件 与下方冲突 - output_mode=merge # 将所有grpc的swagger文档合并到一个中 - fq_schema_naming=true # Schema 命名风格 - title=Go Web API # 指定生成的yaml文件的title # - local: protoc-gen-go # go 不用添加 lang=go # out: api/go_web/v1 # opt: # - paths=source_relative - local: protoc-gen-go-grpc # grpc 不用添加 lang=go out: api/go_web/v1 opt: - paths=source_relative - local: protoc-gen-go-http # http 能生成具体文件的插件一定要指定 lang=go语言 out: api/go_web/v1 opt: - paths=source_relative - local: protoc-gen-go-errors # 错误 不用添加 lang=go out: api/go_web/v1 opt: - paths=source_relative - remote: buf.build/protocolbuffers/go # 已经取代下方字段校验 # - local: protoc-gen-validate # 已经取代下方字段校验 out: api/go_web/v1 opt: paths=source_relative # - local: protoc-gen-validate # 字段校验 # out: api/go_web/v1 # opt: # - paths=source_relative # - lang=go inputs: - directory: api/go_web/v1

buf.yaml

text
复制代码
# For details on buf.yaml configuration, visit https://buf.build/docs/configuration/v2/buf-yaml version: v2 modules: - path: api/go_web/v1 lint: use: - STANDARD # 代码检查规范 breaking: # 兼容性检查 use: - FILE # 基于文件历史变化检查兼容性 except: - ENUM_VALUE_PREFIX # 允许枚举值不使用统一前缀 - ENUM_PASCAL_CASE # 允许枚举值不使用帕斯卡命名 - PACKAGE_VERSION_SUFFIX # 允许包名不包含版本后缀 - PACKAGE_DIRECTORY_MATCH # 允许包名和目录结构不匹配 # service_suffix: Service # 要求服务名称以API结尾 如service UserServiceAPI deps: - buf.build/googleapis/googleapis # - buf.build/envoyproxy/protoc-gen-validate # 已被protovalidate替换 - buf.build/bufbuild/protovalidate

protovalidate

在这里介绍一下校验器的基本用法,有两个重要字段一个是参数校验规则,一个cel参数校验反馈。两者互相独立,推荐使用cel,灵活,两者搭配使用也不是不可。

text
复制代码
message PictureUploadByBatchRequest { string search_text = 1[json_name = "searchText", //JSON序列化时对应的JSON字段 (buf.validate.field).cel = { // cel校验规则 id: "search_text.length" // 唯一ID message: "关键词长度不符合要求" // 校验失败的错误信息 expression: "this.size() >= 1 && this.size() <= 30" // cel表达式 校验条件 }, (buf.validate.field).string.min_len = 1, (buf.validate.field).string.max_len = 30]; int32 count = 2[json_name = "count", (buf.validate.field).cel = { id: "count.value" message: "抓取的图片数目不符合要求" expression: "this >=1 && this <=30" }, (buf.validate.field).int32 = {gt: 1, lt: 30}]; }

如何完成参数校验以及错误返回呢? 我们只需要在service层导入包github.com/bufbuild/protovalidate-go。 使用以下代码轻松校验。例如下方图片编辑接口。其中error_convert.ErrorValidateConvertKratos(err),是我自行设置的将protocalidate的错误转换为kratos的错误用以适配kratos框架使用的。正常直接返回err未必不可。

text
复制代码
func (p *PictureService) PictureEdit(ctx context.Context, req *v1.PictureEditRequest) (*v1.PictureEditResponse, error) { if err := protovalidate.Validate(req); err != nil { return nil, error_convert.ErrorValidateConvertKratos(err) } err := p.picture.EditPicture(ctx, req) return &v1.PictureEditResponse{}, err }

pprof

pprof是go语言内置的性能分析工具,用于诊断程序的CPU,内存,协程等性能问题。他通过采样数据生成可视化报告,帮助开发者定位性能瓶颈。 使用非常简单,在http中引入即可

text
复制代码
srv := http.NewServer(opts...) srv.Handle("/debug/pprof/", pprof.NewHandler())

一般pprof主要分析allocs(内存分配),heap(堆内存),goroutine(协程),profile(CPU分析)。通过分析内存占用情况找出内存泄漏或者高内存占用的地方。

swagger-ui

在kratos中使用swagger-ui文档非常简单,安装官方插件go-kratos/swagger-api

安装插件

text
复制代码
go get -u github.com/go-kratos/swagger-api

在我们的http服务中注册路由,其中openapiv2.NewHandler支持多种配置,不过一般无需添加,作者添加一个无关紧要的配置作为案例。

text
复制代码
srv := http.NewServer(opts...) // swagger swaggerUi := openapiv2.NewHandler(openapiv2.WithGeneratorOptions( generator.UseJSONNamesForFields(true), )) srv.HandlePrefix("/q/", swaggerUi)

踩坑记录:srv.HandlePrefix。 先介绍一下srv.HandlePrefix

text
复制代码
func (s *Server) HandlePrefix(prefix string, h http.Handler)

它能够注册处理一个以prefix为开头的URL路径的处理器,它不会移除这个前缀,同时如果出现不匹配的前缀就会返回404 not found,

因为本项目涉及到websocket,作者又开通了gin路由,为了方便将两个路由放到一个端口使用,作者又注册了一个路由,将gin和grpc的http服务放到一块了,但是为了区分两者,作者添加了httpSrv.HandlePrefix用于区分两者,但是两个srv.HandlePrefix(上方swagger一个,grpc的http服务一个)导致我的路由无法匹配到一起了,怎么都无法访问swagger文档,不得不说这里出现问题要怪我画蛇添足了。 解决方案很简单,把newApp中的http服务的路径匹配/v1删除掉即可这样不影响我们接口访问

text
复制代码
func newApp(logger log.Logger, conf *conf.Server, hs *http.Server, gin *gin.Engine) *kratos.App { // 集成gin框架 httpSrv := http.NewServer( http.Address(conf.Http.Addr), http.Timeout(conf.Http.ReadTimeout.AsDuration()), //添加超时时间 ) httpSrv.HandlePrefix("/gin", gin) //httpSrv.HandlePrefix("/v1", hs) //原先的http服务 httpSrv.HandlePrefix("", hs) return kratos.New( kratos.ID(id), kratos.Name(Name), kratos.Version(Version), kratos.Metadata(map[string]string{}), // 元数据 默认为空 通过HTTP Header 进行传递 kratos.Logger(logger), kratos.Server( //gs, // grpc 服务 gs *grpc.Server, httpSrv, //hs, // http 服务 用于注册其他服务的grpc ), ) }

项目流程

用户模块

比较统一的模块,curd+jwt一套流程走起。略

图片模块

方案设计

cos 腾讯云存储

大量图片存储巨麻烦,我们一般选择使用第三方云存储。保存到自己本地服务器过于累赘,即影响内存又妨碍以后迁移问题。 因为模块自带腾讯云存储,所以我们直接选择腾讯云存储

文件上传下载问题

一般推荐使用SDK工具包,不推荐使用API,API不是很方便 我们自己也已经封装好了完善的上传方法

文件上传

可以使用最简单的SDK再配合我们的文件校验,完成图片上传。我们只不过是需要考虑,图片在数据库中的保存逻辑问题。

文件下载
  1. 后端下载文件,生成流,和前端一直交互,一直看,但是开销大
  2. 获取到文件下载输入流,返回给前端用户。
  3. 通过URL路径直接下载,适合单一被随意允许用户公开访问的资源

对于安全性要求高的场景推荐先后端服务校验权限,在从COS下载文件到服务器,在返回给前端。 或者用户登录后给用户一个密钥,拿着密钥直接从对象存储下载,不用经过后端,这样做性能高。 COS go的SDK提供了两种下载方式 直接下载到本地和下载为字节流 均可使用 但是由于项目本身图片都是公开的,我们直接使用第三种方式下载,凭借URL链接访问即可

图片上传

图片解析代码实现

可以使用使用go自带的image.Decode库手动解析,但是没法获取图片颜色等属性。 cos的sdk可以方便我们进行解析。使用时调用CI方法即可。

在提供一个第三方库Imagemeta,用于解析元数据,而不是原始图像信息Imagemeta的github链接 注意要多引入一个包

text
复制代码
// UploadPictureAndParse 上传图像并解析图像 k /文件夹/文件名 _ "image/jpeg" func (t *tencentClient) UploadPictureAndParse(ctx context.Context, file *multipart.FileHeader, key string) (image.Image, string, string, error) { // 获取图片文件 fileOpen, err := file.Open() if err != nil { return nil, "", "", err } // 关闭文件流 defer func(fileOpen multipart.File) { err := fileOpen.Close() if err != nil { return } }(fileOpen) // 获取图片信息与图片格式 decode, format, err := image.Decode(fileOpen) if err != nil { t.log.Error(err) return nil, "", "", constant.ErrFilePictureParse } // 上传图片 并获取图片路径 filepath, err := t.UploadFile(ctx, file, key) if err != nil { return nil, "", "", err } return decode, format, filepath, nil }

使用CI调用时,可以获取多种数据。

text
复制代码
res, _, err := t.client.CI.Put(ctx, path, fileOpen, opt) res.OriginalInfo.ImageInfo

url地址上传文件

自己手撸吧,go没有很多方便使用的库只能手撸了

text
复制代码
// UploadPictureByUrl 通过路径上传文件 func (t *tencentClient) UploadPictureByUrl(ctx context.Context, fileUrl string) (*biz.PictureInfo, error) { path, format, fileName, size, err := t.validFileUrl(fileUrl) if err != nil { return nil, err } fileOpen, err := DownloadToMultipartFile(fileUrl) if err != nil { t.log.Error(err) return nil, constant.ErrFileDownload } // 关闭文件流 defer func(fileOpen multipart.File) { err := fileOpen.Close() if err != nil { return } }(fileOpen) // 获取图片信息与图片格式 decode, format, err := image.Decode(fileOpen) if err != nil { t.log.Error(err) return nil, constant.ErrFilePictureParse } _, err = fileOpen.Seek(0, io.SeekStart) // 重置文件指针 因为 image.Decode(fileOpen)以将文件读取完 if err != nil { return nil, constant.ErrFilePictureParse } // 上传图片 并获取图片路径 _, err = t.client.Object.Put(ctx, path, fileOpen, nil) if err != nil { t.log.Error(err) return nil, err } return &biz.PictureInfo{ Url: filePath + path, PicName: fileName, PicSize: size, PicWidth: int64(int32(decode.Bounds().Dx())), PicHeight: int64(int32(decode.Bounds().Dy())), PicScale: float32(decode.Bounds().Dx() * 1.0 / decode.Bounds().Dy()), PicFormat: format, }, nil } // DownloadToMultipartFile 下载文件并返回 multipart.File 接口对象 func DownloadToMultipartFile(url string) (multipart.File, error) { // 发送 HTTP 请求 resp, err := http.Get(url) if err != nil { return nil, err } defer func(Body io.ReadCloser) { err := Body.Close() if err != nil { } }(resp.Body) if resp.StatusCode != http.StatusOK { return nil, err } // 创建临时文件 tmpFile, err := os.CreateTemp("", "downloaded-*") if err != nil { return nil, err } // 将内容写入临时文件 _, err = io.Copy(tmpFile, resp.Body) if err != nil { err := tmpFile.Close() if err != nil { return nil, err } err = os.Remove(tmpFile.Name()) if err != nil { return nil, err } return nil, err } // 重置文件指针到开头(便于后续读取) _, err = tmpFile.Seek(0, io.SeekStart) if err != nil { err := tmpFile.Close() if err != nil { return nil, err } err = os.Remove(tmpFile.Name()) if err != nil { return nil, err } return nil, err } return tmpFile, nil }

批量抓取图片

java有jsoup库,我go有colly!而且这个库还在更新! colly爬虫库自带html解析简单好用能上手! colly官方

然后为了减少人工审核速度,我们直接在请求中填写好要完善的名称种类标签以及简介。

后端设计
text
复制代码
// BatchPicture 爬虫爬取图片 func (p *PictureUseCase) BatchPicture(ctx context.Context, req *v1.PictureUploadByBatchRequest) (int32, error) { //抓取内容 这个网站不容易被封 fetchUrl := fmt.Sprintf("https://cn.bing.com/images/async?q=%s&mmasync=1", req.SearchText) //解析tag var tags string if len(req.Tags) != 0 { bytes, err := json.Marshal(req.Tags) tags = string(bytes) if err != nil { return 0, constant.ErrParams } } //解析内容 c := colly.NewCollector() // 限制并发和延迟(避免被封) err := c.Limit(&colly.LimitRule{ DomainGlob: "*bing.com*", Parallelism: 2, Delay: 1 * time.Second, }) if err != nil { return 0, err } // 3. 图片计数器 var downloadedCount int maxDownload := int(req.Count) // 从请求参数获取最大下载数量 // 4. 解析图片元素 对于我们遇到的每一个goquerySelector元素 都会进入一次循环 直到全部元素遍历完毕 c.OnHTML("img.mimg", func(e *colly.HTMLElement) { // 全部下载完毕 停止解析 if downloadedCount >= maxDownload { return } // src获取到的一般是缩略图 源地址高清图请将对于我们遇到的每一个goquerySelector元素设置为a.iusc 从murl中获取源地址高清图 //imgM := e.Attr("m") //var data map[string]interface{} //err := json.Unmarshal([]byte(imgM), &data) //murl := data["murl"] //p.log.Infof("murl:%s", murl) // 对路径进行处理 防止转义或者与对象存储产生冲突 以及确保获取完整图片内容 imgUrl := e.Attr("src") if imgUrl == "" { imgUrl = e.Attr("data-src") } p.log.Infof("imgUrl Original:%s", imgUrl) // 去除链接参数 markIndex := strings.Index(imgUrl, "?") if markIndex > 0 { // 防止下标越界 imgUrl = imgUrl[:markIndex] } p.log.Infof("imgUrl:%s", imgUrl) // 处理相对路径 if !strings.HasPrefix(imgUrl, "http") { imgUrl = "https:" + imgUrl } // 流式传输 //// 下载图片 //resp, err := http.Get(imgUrl) //if err != nil { // p.log.Errorf("下载图片失败: %v, URL: %s", err, imgUrl) // return //} //defer resp.Body.Close() // //// 获取图片内容 //imgData, err := io.ReadAll(resp.Body) //if err != nil { // p.log.Errorf("读取图片数据失败: %v", err) // return //} // 5. 获取文件名 fileName := filepath.Base(imgUrl) if !strings.HasSuffix(strings.ToLower(fileName), ".jpg") && !strings.HasSuffix(strings.ToLower(fileName), ".jpeg") && !strings.HasSuffix(strings.ToLower(fileName), ".png") { fileName = fmt.Sprintf("%d_%s.jpg", time.Now().Unix(), req.SearchText) } //// 6. 上传cos //pictureInfo, err := p.oss.UploadFileBytes(ctx, imgData) pictureInfo, err := p.oss.UploadPictureAndParseByUrl(ctx, imgUrl) if err != nil { p.log.Errorf("上传图片失败: %v", err) return } // 7. 保存图片数据到数据库 可以优化为批量插入 但是因为此方法仅限管理员使用 没有任何并发量 故可以等待 _, err = p.SaveBatchPicture(ctx, &model.Picture{ URL: pictureInfo.Url, PicName: fmt.Sprintf("%s%d", req.SearchText, downloadedCount+1), PicSize: pictureInfo.PicSize, PicWidth: int32(pictureInfo.PicWidth), PicHeight: int32(pictureInfo.PicHeight), PicScale: float64(pictureInfo.PicScale), PicFormat: pictureInfo.PicFormat, Category: req.Category, Tags: tags, Introduction: req.Introduction, }) if err != nil { p.log.Errorf("保存图片失败: %v", err) return } downloadedCount++ p.log.Infof("成功下载并上传第 %d/%d 张图片: %s", downloadedCount, maxDownload, pictureInfo.Url) }) // 8. 错误处理 c.OnError(func(r *colly.Response, err error) { p.log.Errorf("请求失败: %v, URL: %s", err, r.Request.URL) }) // 9. 开始爬取 err = c.Visit(fetchUrl) if err != nil { return 0, constant.ErrUrlNotFound } // 10. 等待所有异步请求完成 c.Wait() //返回成功下载图片数量 return int32(downloadedCount), nil }

图片优化

图片查询优化

一提到查询优化就想到了redis...,如果你是这么想的,没毛病,你想对了。 每次从数据获取数还是比较慢的。

多级缓存

多级缓存兼具各个缓存的优点,本地缓存的高性能,分布式缓存的数据一致性。 这里单体架构推荐go-cache库。分布式还是用redis就好。 github

go-cache使用时建议将cache对象集成到biz层需要的服务中去,使用极其方便。

text
复制代码
// ListPageCache 批量获取文件 缓存 func (p *PictureUseCase) ListPageCache(ctx context.Context, req *v1.PictureListRequest) (pictures []*model.Picture, total int32, err error) { req.ReviewStatus = constant.Pass // 强制只查询已审核图片 // 1. 构造缓存Key MD5哈希请求参数可以减少key的长度 marshal, err := json.Marshal(req) if err != nil { return nil, 0, constant.ErrParams } redisKey := fmt.Sprintf("%s%x", constant.RedisPictureList, md5.Sum(marshal)) total, err = p.pictureRepo.Count(ctx, req) if err != nil { return nil, 0, err } // 2. 尝试从一级缓存(本地)获取 // 五分钟的缓存 10分钟清理一过期数据 if pictures, found := p.cache.Get(redisKey); found { return pictures.([]*model.Picture), total, nil } // 3. 一级缓存未命中,尝试从二级缓存(分布式缓存)获取 if pictures, err := p.pictureRepo.ListPageCache(ctx, redisKey); err == nil && pictures != nil { // 二级缓存命中,返回结果 并且更新一级缓存 p.cache.Set(redisKey, pictures, cache.DefaultExpiration) return pictures, total, err } // 4. 多级缓存均未命中,查询数据库 pictures, err = p.pictureRepo.ListPage(ctx, req) if err != nil { return nil, 0, fmt.Errorf("%w: 查询失败", constant.ErrParams) } // 4. 异步更新多级缓存(忽略错误) go func() { if picturesJson, err := json.Marshal(pictures); err == nil { p.cache.Set(redisKey, pictures, cache.DefaultExpiration) _ = p.pictureRepo.SetCache(ctx, redisKey, string(picturesJson)) } }() // 5. 返回结果 return pictures, total, err }

图片上传优化

压缩方案

压缩方案其实交给第三方完成即可,我们可以让腾讯云直接压缩,简单轻松。官方GO的图片压缩为WebP等格式的SDK见下方。 GO-SKD图片高级压缩 同时我们在做一个缩略图(就是将原先的图片大小不变,但是图像等比例缩小),三种不同的图片:原图,压缩图,缩略图。同时上传至cos,只需要配置好Rule即可。其中put方法上传时本身就是上传了原图,我们只需要为Rule提供缩略图以及webp图即可 案例

text
复制代码
func (t *tencentClient) UploadPictureAndParse(ctx context.Context, file *multipart.FileHeader) (*biz.PictureInfo, error) { path, err := t.validFile(file) if err != nil { return nil, err } // 获取图片文件 fileOpen, err := file.Open() if err != nil { return nil, err } // 关闭文件流 defer func(fileOpen multipart.File) { err := fileOpen.Close() if err != nil { return } }(fileOpen) // 处理图片上传路径 split := strings.Split(path, ".") pathWebp := split[0] + ".webp" pathThumbnail := split[0] + "_thumbnail" + "." + split[1] // 定制规格同时保存原图与缩略图以及压缩图 开始上传压缩图 pic := &cos.PicOperations{ IsPicInfo: 0, // 是否返回原图信息 0默认不 1返回 Rules: []cos.PicOperationsRules{ { FileId: pathWebp, // 保存路径 Rule: "imageMogr2/format/webp", // 转换格式 webp }, { FileId: pathThumbnail, // 保存路径 Rule: fmt.Sprintf("imageMogr2/thumbnail/%dx%d>", 128, 128), // 转换格式缩略图 }, }, } opt := &cos.ObjectPutOptions{ ACLHeaderOptions: nil, ObjectPutHeaderOptions: &cos.ObjectPutHeaderOptions{ XOptionHeader: &http.Header{}, }, } opt.XOptionHeader.Add("Pic-Operations", cos.EncodePicOperations(pic)) // 同时上传三种图片 Put函数本身就上传了一次原图 _, err = t.client.Object.Put(ctx, path, fileOpen, opt) if err != nil { // 压缩图上传失败 仍然使用原图 t.log.Error(err) return nil, err } _, err = fileOpen.Seek(0, io.SeekStart) // 重置文件指针 因为 image.Decode(fileOpen)以将文件读取完 if err != nil { return nil, constant.ErrFilePictureParse } // 获取图片信息与图片格式 decode, format, err := image.Decode(fileOpen) if err != nil { t.log.Error(err) return nil, constant.ErrFilePictureParse } return &biz.PictureInfo{ Url: filePath + path, PicName: strings.Split(path, "/")[2], PicSize: file.Size, PicWidth: int64(int32(decode.Bounds().Dx())), PicHeight: int64(int32(decode.Bounds().Dy())), PicScale: float32(decode.Bounds().Dx() * 1.0 / decode.Bounds().Dy()), PicFormat: format, Webp: filePath + pathWebp, ThumbnailUrl: filePath + pathThumbnail, }, nil }

其他可扩展选项:断点续传与分片上传。这两者功能均可以由cos自身实现,自己开发难度较大算了吧...

图片存储优化

深度沉降:可以控制桶的声明周期功能,将长时间未访问的数据进行沉降,大幅降低数据存储成本。 转存失败,建议直接上传图片文件 转存失败,建议直接上传图片文件

数据删除:当图片被删除或者图片被覆盖更新的时候,我们应该把cos中的图片同步删除掉。

空间模块

私有空间的权限控制

私有空间的权限和公共图库是不同的,我们需要对所有的图片操作都添加和空间有关的权限校验逻辑。 这样的话要对大面积的图片逻辑进行更改。 创建图片的校验不需要过于臃肿,我们明知图片和空间严格挂钩,只要上传图片时让前端必须给我提供空间id即可。为了保证前端一定传给我们,前端可以设置传值的时候默认传0,也就是公共空间的值。当然,go的结构体默认也会初始化int类型为0,也就是说即使不传值,也会默认使用公共空间。 增删改

空间级别和限额控制

1. 上传图片时校验和更新额度。

现在我们的上传图片代码已经不叫复杂了。此时如果在进行校验在添加逻辑比较头疼。我们不妨对业务逻辑进行优化。 单张图片最大2M,即使空间满了我们也允许上传。如果用户在限流之前大量上传图片,那么我们可以通过限流+定时任务检测来进行限制。

图片功能扩展

图片搜索 - 基础属性搜索

主要是让前端实现模糊查询的各个字段。后端改动很少

以图搜图-数据抓取

最简单的就是使用第三方API 百度以图搜图 bing图片

复杂方案:抓取第三方数据,还是爬虫。 整体需要调用三个API,第一个API是向百度发送post请求,获取给定的url页面地址。 第二步获取到对应地址的HTML页面内容,提取其中包含的firstUrl的js脚本,返回对应的图片列表的页面地址。 第三步:对图片列表进行解析,将json数据序列化为我们的结构体数据。 其中百度在请求头中添加了一个asc-token,我们必须添加这个请求头才可以正常访问,但是目前可能有bug,这个请求头的值可以是任意的。嗯......

颜色搜图

通过欧几里得(欧式)距离算法,就是每种颜色的平方差的和求根号。贴一张鱼皮的图片 转存失败,建议直接上传图片文件 写个计算方法,数据库查数据后用算法为每一个数据计算出结果后进行排序即可

text
复制代码
func (p *PictureUseCase) ListPageByColor(ctx context.Context, req *PictureListCondition) ([]*model.Picture, int32, error) { if err := p.space.CheckAuthSpace(ctx, req.SpaceId, false); err != nil { return nil, 0, err } pictures, err := p.pictureRepo.FindAll(ctx, req) type colorSimilarity struct { Picture *model.Picture Similarity float64 } var colorSimilarities []colorSimilarity for _, picture := range pictures { if picture.Ave == "" { //无颜色的图片 不参与排列 continue } colorSum, err := util.CalculateSimilarityHex(req.Ave, picture.Ave) if err != nil { // 颜色计算出现问题 跳过该数据 continue } colorSimilarities = append(colorSimilarities, colorSimilarity{ Picture: picture, Similarity: colorSum, }) } // 按照相似度降序排序 sort.Slice(colorSimilarities, func(i, j int) bool { return colorSimilarities[i].Similarity > colorSimilarities[j].Similarity }) // 只去前12条数据 limit := 12 if len(colorSimilarities) < limit { limit = len(colorSimilarities) } var result []*model.Picture for i := 0; i < limit; i++ { result = append(result, colorSimilarities[i].Picture) } return result, int32(len(pictures)), err }

批量操作

简答的批量操作只靠sql即可完成,因为gorm不支持updatebatch操作,可以用 gorm-bulk-insert库,或者for循环操作。当然如果同步让所有要改的数据改为同一个值,仅仅靠id in (...)即可完成。 但是如果是超大批量操作,例如操作数据超过千条,需要分批操作并且利用通道和协程辅助完成。粘贴一个实例

text
复制代码
func (p *PictureUseCase) ListEditBatch(ctx context.Context, ids []int64, req *PictureListCondition) error { if err := p.space.CheckAuthSpace(ctx, req.SpaceId, false); err != nil { return err } marshal, _ := json.Marshal(req.Tags) tags := string(marshal) batchSize := 100 //每批最多处理100条数据 workers := (len(ids) + batchSize - 1) / batchSize // 创建工作通道 每个通道只发送100条数据 idChan := make(chan []int64, workers) // 用于接收每个协程的错误 errChan := make(chan error, workers) var wg sync.WaitGroup for _ = range workers { wg.Add(1) go func() { defer wg.Done() for id := range idChan { select { case <-ctx.Done(): //主动终止 errChan <- ctx.Err() return default: err := p.pictureRepo.EditList(ctx, id, &model.Picture{ Category: req.Category, Tags: tags, SpaceID: req.SpaceId, PicName: req.PicName, }) errChan <- err } } }() } // 分批发送数据 for i := 0; i < workers; i++ { // 计算当前批处理的起始和结束索引 防止下标越界 start := i * batchSize end := min(start+batchSize, len(ids)) select { case idChan <- ids[start:end]: case <-ctx.Done(): break } } // 发送完毕立刻关闭通道释放资源 close(idChan) // 收集错误 var errs []error go func() { // 等待任务结束 结束后立刻关闭通道 wg.Wait() close(errChan) }() for err := range errChan { if err != nil { errs = append(errs, err) } } if len(errs) > 0 { p.log.Errorf("批量更新失败: %v", errs) return errs[0] } return nil }

AI编辑模块

AI模块接入

使用阿里云百炼(鱼皮都推荐了,就不找其它的了)

ai调用有个通病,就是慢,一直阻塞程序容易超时,体验很差。我们采用异步调用,每次请求时返回任务的id,任务完成后我们可以根据任务的id去查询结果。整个流程就是,前端发送请求给后端,后端接收后调用API,调用完成后获取到任务id发送给前端并保存到数据库中,前端根据任务id轮询发送请求到后端尝试获取任务结果。这种调用方式让前后端解耦,阻止了后端因为一个任务长时间的停滞阻塞问题。 同步流程 转存失败,建议直接上传图片文件

异步流程 转存失败,建议直接上传图片文件

总之代码不是很复杂。其中百炼返回的响应结构体我用proto文件编辑,但是因为最终生成的结构体中包含两个只有proto文件生成的特殊字段。导致返回的响应无法映射到结构体上。不影响最终的结构体映射因为有protojson库可以帮助完成普通json转换为proto结构体。 但是,我最后本来想通过对结构体类型,进行switch判断,来判断出最后失败、成功还是轮询。但是不行了,只能通过拿出json字段中的具体内容来一步步判断是否成功了。不过也不复杂。 还有一个问题就是这个百炼扩图后绝大多是的图片大小都远超2mb,多的从原图300kb涨到了8mb,导致我们上传到自己服务器无法通过校验,于是我直接将原图url改为了生成图片的url反正不消耗我们自己的资源,无非是用户打开详情页多加载一会...

text
复制代码
func (a *aliAiRepo) GetImageOutPaintingTask(ctx context.Context, id string) (*v1.ImageOutPaintingSuccessTaskResponse, error) { // 创建请求 client := http.Client{ Timeout: 5 * time.Second, } req, err := http.NewRequest(http.MethodGet, strings.ReplaceAll(constant.GET_OUT_PAINITING_TASK_URL, "%s", id), nil) if err != nil { a.log.Error(err) return nil, err } req.Header.Set("Authorization", "Bearer "+a.conf.AliKey) resp, err := client.Do(req) defer func(Body io.ReadCloser) { if err := Body.Close(); err != nil { } }(resp.Body) // 检查请求状态 if resp.StatusCode != http.StatusOK { return nil, constant.ErrApi } body, err := io.ReadAll(resp.Body) if err != nil { return nil, err } var result map[string]interface{} // 未知类型的响应 保存到RawMessage中 if err := json.Unmarshal(body, &result); err != nil { // 如果是 JSON 则用 json.Unmarshal return nil, err } put := result["output"].(map[string]interface{}) switch put["task_status"].(string) { case "SUCCEEDED": // 返回成功逻辑 var response v1.ImageOutPaintingSuccessTaskResponse if err := protojson.Unmarshal(body, &response); err != nil { return nil, err } return &response, nil case "FAILED": // 返回错误逻辑 return nil, errors2.New(put["message"].(string)) case "RUNNING": // 返回进行中逻辑(例如轮询或等待) return nil, nil default: return nil, constant.ErrApi } }

空间分析

总体来说这一节比较轻松,但是gorm/gen让我浪费了一天时间。好在最后解决了。

踩坑记录:gorm/gen起别名问题。解决方案已在gorm中

团队空间

权限控制

采用经典的RBAC权限控制模式,(我们使用的casbin就是采用的RBAC,当然casbin不仅仅只支持这一种模型)。 casbin已经介绍了最优雅最安全的多租户功能。但是不得不说的是,多租户确实麻烦。对于权限校验不麻烦,但是每一次更换校验模型,我们就要更改大量代码。

其次就是简单的角色继承,我们让空间角色继承不同的角色。用户原本的角色不需要改动。当用户进入对应空间时,我们主动搜索用户在对应空间的角色,然后将角色通过请求头传递进去X-Domain,这样我们可以将用户域或用户的角色传进去。为了保证安全可以将信息加密。 前端想要更新也无需重新登陆,添加一个刷新个人信息按钮,重新获取空间角色信息即可。

最后最简单的方案就是自定义校验规则:我们直接在不同接口中获取到用户在空间的权限,如果权限不足就报错。这种方法其实也是用的最多的。作者咨询了多位大佬,大佬们表示使用框架做可以但是用起来太麻烦,不如直接代码式权限校验。

分库分表

go没有那么好用的库支持分库分表。

不过只要能做到对于一个id进行一个计算后,最后总能得到一个固定的分片值,那么对条数据无论是任何操作都可以进行了。

感谢站内KingYen大佬的建议(大佬在评论区),gorm官方有一个分页插件gorm.io/sharding

分表后以及该库有不少缺点和注意事项罗列一下:

  1. 该库只支持gorm,不支持gen,gen的话主要为了安全而设计,不是很灵活。
  2. 该库无法动态扩容,无法自动创建表,需要手动创建。
  3. 这个库比较旧了,可能已经没有需要完善的地方了,这个库安装后可能会让gorm部分包版本倒退,出现报错,安装完这个库后建议在重新更新下gorm库
  4. 分表后无法在进行模糊查询,like ,in 等操作都被禁止,条件查询只能使用 =,而且条件中必须包含分页主键,这就意味着分页查询也无法进行。解决方案作者只想到两个:一、搜索全部的数据,手动模糊查询(有些复杂而且很麻烦,数据量大了很占内存)。二、使用ES中间件完成搜索功能(这可能是最理想的方案了)
  5. 分表后如果想查询全部数据需要自己手动通过拼接所有的表名来实现查询,下方有案例
  6. 分表下表从0开始,例如我设置picture分为4个表,那么表名就是从picture_0picture_3

总之如果考虑好要分表了,那么就要考虑这个项目的复杂性,一般情况下还是不要分表为好。

下载包

text
复制代码
go get -u gorm.io/sharding

使用方法:在生成gorm实例后注册一个中间件即可

text
复制代码
db, err := gorm.Open(mysql.Open(conf.Database.Source), &gorm.Config{ DisableForeignKeyConstraintWhenMigrating: true, Logger: gormLogger.Default.LogMode(toGormLogLevel(conf.Database.LogLevel)), }) if err != nil { logs.Fatalf("failed opening connection to mysql: %v", err) } // 注册分表中间件 err = db.Use(sharding.Register(sharding.Config{ ShardingKey: "pid", //指定分表字段 NumberOfShards: 2, // 指定分表数量 PrimaryKeyGenerator: sharding.PKSnowflake, // 指定主键生成器 }, "picture")) // 指定表名 if err != nil { logs.Fatalf("failed register sharding middleware: %v", err) }

其中,中间件由多个配置属性。前缀和后缀如果添加了,就是相当于在表名的基础上在拼接上对应的字符,例如picture表名添加一个space_前缀,最后生成的表名就是space_picture

参数类型说明
ShardingKeystring用于路由分表的字段名
NumberOfShardsint分表总数
PrimaryKeyGeneratorPKGenerator主键生成器(支持雪花算法,uuid,以及自定义算法)
ShardingAlgorithmfunc(value interface{}) (suffix string, err error)自定义分表算法(可选)
TablePrefixstring表前缀(可选)
TableSuffixstring表后缀(可选)

注册后我们只需要正常使用gorm即可,操作时会自动帮我们路由到对应的分表,不过分表后的查询也会相应的变成全量查询。推荐还是使用gen帮助我们生成对应的表映射结构体。

创建图片案例

text
复制代码
func (s *pictureRepo) Save(ctx context.Context, picture *model.Picture) (*model.Picture, error) { err := s.data.db.WithContext(ctx).Model(&model.Picture{}).Save(picture).Error return picture, err }

查询所有数据

text
复制代码
func (s *pictureRepo) ListPage(ctx context.Context, condition *biz.PictureListCondition) ([]*model.Picture, error) { var pictures []*model.Picture var pictureSharding []*model.Picture for i := range 4 { err := s.data.db.WithContext(ctx).Model(&model.Picture{}).Table(fmt.Sprintf("picture_%d", i)).Scan(&pictureSharding).Error if err != nil { return pictures, err } pictures = append(pictures, pictureSharding...) } return pictures, nil }

图片协同

关于websocket作者也是一知半解。对于框架选择提供两种:第一种,轻量简单好用。第二种是某位大佬推荐(看实际项目运行来看,确实第二种更复杂但是可维护性与可用性更高),但是第二种已经归档,不再更新,这也是最遗憾的地方,只能盼望有后来者更新了。不过目前仍然可以使用,未必就不行了。

  1. gorilla/websocket
  2. socket.io

第二种本质是实现了Socket.IO,这也是大佬为什么使用的原因,它自己本身搭载了redis适配器,支持广播的发送接收。

作者选择第一种,简单快速上手就够了,我们的项目用不到第二种,无需太复杂。

websocket

全双工通信协议,让客户端和服务端之间保持实时且持续的连接,常用于聊天通话直播等需要长连接的地方。

go搭建websocket

虽然我们找到了gorilla/websocket库方便我们开发,但是还是要自己处理很多地方的。 grpc因为不支持websocket,所以我们要引入gin框架,因为只需要gin框架连接websocket即可,所以我们不需要在添加其他校验中间件,只引入跨域中间件即可。

websocket整体设计(自上而下的设计方便代码阅读)
go
复制代码
package service import ( "context" v1 "dream_cloud_chart/api/go_web/v1" "dream_cloud_chart/internal/biz" "dream_cloud_chart/internal/conf" "dream_cloud_chart/internal/data/dal/model" "dream_cloud_chart/internal/pkg/authz" "dream_cloud_chart/pkg/common/constant" "fmt" "net/http" "strconv" "sync" "time" "github.com/golang-jwt/jwt/v5" "go.uber.org/zap" "google.golang.org/protobuf/encoding/protojson" "github.com/gin-gonic/gin" "github.com/go-kratos/kratos/v2/log" "github.com/gorilla/websocket" ) type WebSocketService struct { log *log.Helper auth *conf.Auth picture *biz.PictureUseCase user *biz.SysUserUseCase rwLocker sync.RWMutex // 读写锁 提高并发效率 pictureEdit sync.Map pictureSession map[int64]map[int64]*Node } func NewWebSocketService(picture *biz.PictureUseCase, logger log.Logger, user *biz.SysUserUseCase, auth *conf.Auth) *WebSocketService { return &WebSocketService{ log: log.NewHelper(log.With(logger, "module", "service/picture")), picture: picture, user: user, auth: auth, rwLocker: sync.RWMutex{}, pictureSession: map[int64]map[int64]*Node{}, pictureEdit: sync.Map{}, } } // Node 本核心在于形成userid和Node的映射关系 type Node struct { Conn *websocket.Conn //并行转串行 Conn是IO型数据 大量IO并行的话容易乱 串行数据就能保证数据顺序了 DataQueue chan []byte // 该节点的用户信息 User *model.SysUser PictureId int64 // 并发控制 ctx context.Context // 控制协程生命周期 cancel context.CancelFunc // 终止函数 } // NewNode 创建新Node(工厂方法) func NewNode(conn *websocket.Conn, user *model.SysUser, pictureId int64) *Node { ctx, cancel := context.WithCancel(context.Background()) return &Node{ Conn: conn, User: user, PictureId: pictureId, DataQueue: make(chan []byte), //无缓冲 减少内存占用 ctx: ctx, cancel: cancel, } } func (w *WebSocketService) Chat(c *gin.Context, writer http.ResponseWriter, request *http.Request) { w.log.Info("正在连接websocket") // 获取路径上的参数 query := request.URL.Query() //role := query.Get("role") token := query.Get("token") pictureIds := query.Get("picture") // 校验用户权限 claims, err := w.TokenValid(c, token, constant.SpaceRoleAdmin) if err != nil { return } pictureId, err := strconv.Atoi(pictureIds) if err != nil { return } // 创建并升级连接 conn, err := (&websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true }, }).Upgrade(writer, request, nil) if err != nil { w.log.Error("websocket upgrade failed", zap.Error(err)) return } // 创建Node节点 node := NewNode(conn, &model.SysUser{ UID: claims.UserID, Username: claims.Nickname, NickName: claims.Nickname, }, int64(pictureId)) // 设置连接关闭处理器 conn.SetCloseHandler(func(code int, text string) error { // 退出编辑模式 w.handleExitEditMessage(nil, node) // 清理会话 w.rwLocker.Lock() _, ok := w.pictureEdit.Load(node.PictureId) if ok { w.pictureEdit.Delete(node.PictureId) } if m, ok := w.pictureSession[node.PictureId]; ok { delete(m, node.User.UID) } w.rwLocker.Unlock() // 终止接收和发送协程 node.cancel() w.log.Infof("%s %d用户会话已被清理", node.User.NickName, node.User.UID) return nil }) // 添加节点 // 读写锁防止并发问题 w.rwLocker.Lock() nodes, ok := w.pictureSession[int64(pictureId)] if ok { nodes[claims.UserID] = node w.pictureSession[int64(pictureId)] = nodes } else { w.pictureSession[int64(pictureId)] = map[int64]*Node{ claims.UserID: node, } } w.rwLocker.Unlock() // 启动发送和接收消息的协程 提高并发性能 go w.SendProc(node) go w.RecvProc(node) // 为其他在编辑的用户添加广播 w.BroadCastToPicture(int64(pictureId), &v1.PictureTeamEditResponse{ Message: fmt.Sprintf("%s加入编辑", claims.Nickname), Type: constant.PICTURE_TEAM_MESSAGE_INFO, User: &v1.UserData{ Nickname: claims.Nickname, Uid: strconv.FormatInt(claims.UserID, 10), }, }) } // BroadCastToPicture 对操作指定图片的所有用户进行广播 func (w *WebSocketService) BroadCastToPicture(pictureId int64, request *v1.PictureTeamEditResponse) { nodes := w.pictureSession[pictureId] marshal, _ := protojson.Marshal(request) // 串行发送消息 //for _, node := range nodes { // node.DataQueue <- marshal //} // 并行发送消息 var wg sync.WaitGroup for _, node := range nodes { wg.Add(1) go func(n *Node) { defer wg.Done() select { case n.DataQueue <- marshal: case <-time.After(100 * time.Millisecond): // 超时保护 w.log.Warn("broadcast timeout") } }(node) } wg.Wait() } // TokenValid token校验 func (w *WebSocketService) TokenValid(c *gin.Context, token, role string) (*authz.TokenClaims, error) { // 校验权限 if role != constant.SpaceRoleEditor && role != constant.SpaceRoleAdmin { return nil, constant.ErrWebSocketNotAuth } // 校验token parse, err := jwt.ParseWithClaims(token, &authz.TokenClaims{}, func(token *jwt.Token) (interface{}, error) { return []byte(w.auth.User.UserSecretKey), nil }) if err != nil { return nil, constant.ErrUserAccountHasException } return parse.Claims.(*authz.TokenClaims), nil } // SendProc 发送协程 func (w *WebSocketService) SendProc(node *Node) { for { select { case data := <-node.DataQueue: if err := node.Conn.WriteMessage(websocket.TextMessage, data); err != nil { w.log.Error("write message failed", zap.Error(err)) node.cancel() return } //监听推出信号 case <-node.ctx.Done(): return } } } // RecvProc 接收的协程 接收到数据后调度发送给需要接收的人 func (w *WebSocketService) RecvProc(node *Node) { for { _, data, err := node.Conn.ReadMessage() if err != nil { w.log.Error("read message failed", err) return } //defer node.Conn.Close() w.dispatch(data, node) } } // dispatch 后端调度处理方便我们处理不同的类型数据 func (w *WebSocketService) dispatch(data []byte, node *Node) { // 解析data为PictureTeamEditResponse msg := new(v1.PictureTeamEditRequest) err := protojson.Unmarshal(data, msg) w.log.Info(string(data)) w.log.Infof("msg=%v", msg) if err != nil { w.log.Info(data) w.log.Error("unmarshal json failed", err.Error()) return } // 根据CMD 对逻辑进行处理 switch msg.Type { case constant.PICTURE_TEAM_MESSAGE_ENTER_EDIT: // 进入编辑 w.handleEnterEditMessage(msg, node) case constant.PICTURE_TEAM_MESSAGE_EXTI_EDIT: // 推出编辑 w.handleExitEditMessage(msg, node) case constant.PICTURE_TEAM_MESSAGE_ACTION: // 执行操作 w.handleEditActionMessage(msg, node) default: w.handleErrorMessage(msg, node) } } // 加入编辑 func (w *WebSocketService) handleEnterEditMessage(req *v1.PictureTeamEditRequest, node *Node) { _, ok := w.pictureEdit.Load(node.PictureId) if !ok { //没有用户编辑当前图片 则可以进入编辑 w.pictureEdit.Store(node.PictureId, node.User.UID) PictureTeamEditResponse := &v1.PictureTeamEditResponse{ Type: constant.PICTURE_TEAM_MESSAGE_ENTER_EDIT, Message: fmt.Sprintf("%s进入编辑", node.User.NickName), User: &v1.UserData{ Uid: strconv.FormatInt(node.User.UID, 10), Nickname: node.User.NickName, }, } w.BroadCastToPicture(node.PictureId, PictureTeamEditResponse) } } // 执行操作 func (w *WebSocketService) handleEditActionMessage(req *v1.PictureTeamEditRequest, node *Node) { find, ok := w.pictureEdit.Load(node.PictureId) if ok && find == node.User.UID { PictureTeamEditResponse := &v1.PictureTeamEditResponse{ Type: constant.PICTURE_TEAM_MESSAGE_ACTION, Message: fmt.Sprintf("%s执行%s", node.User.NickName, req.EditAction), User: &v1.UserData{ Uid: strconv.FormatInt(node.User.UID, 10), Nickname: node.User.NickName, }, } // 广播 w.BroadCastToPicture(node.PictureId, PictureTeamEditResponse) } } // 退出编辑 func (w *WebSocketService) handleExitEditMessage(req *v1.PictureTeamEditRequest, node *Node) { find, ok := w.pictureEdit.Load(node.PictureId) if ok && find == node.User.UID { // 确认确实是当前用户 w.pictureEdit.Delete(node.PictureId) PictureTeamEditResponse := &v1.PictureTeamEditResponse{ Type: constant.PICTURE_TEAM_MESSAGE_EXTI_EDIT, Message: fmt.Sprintf("%s退出编辑", node.User.NickName), User: &v1.UserData{ Uid: strconv.FormatInt(node.User.UID, 10), Nickname: node.User.NickName, }, } // 广播 w.BroadCastToPicture(node.PictureId, PictureTeamEditResponse) } } // 错误 func (w *WebSocketService) handleErrorMessage(req *v1.PictureTeamEditRequest, node *Node) { find, ok := w.pictureEdit.Load(node.PictureId) if ok && find == node.User.UID { // 确认确实是当前用户 w.pictureEdit.Delete(node.PictureId) PictureTeamEditResponse := &v1.PictureTeamEditResponse{ Type: constant.PICTURE_TEAM_MESSAGE_ERROR, Message: fmt.Sprintf("出现错误"), User: &v1.UserData{ Uid: strconv.FormatInt(node.User.UID, 10), Nickname: node.User.NickName, }, } // 广播 w.BroadCastToPicture(node.PictureId, PictureTeamEditResponse) } }

至此我们可以说是项目结束。整体本身就是类DDD架构,无需再改为ddd架构

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP