Kratos 云盘项目:DDD 分层、依赖注入与服务治理

2025-01-15T10:30:00+08:00 | 27分钟阅读 | 更新于 2025-01-15T10:30:00+08:00

@

学习目标

学完这篇,你应该能:

  1. 说清楚 为什么 一个 Go 微服务要用 service → biz → data 的 DDD 四层,并用"依赖倒置"解释 biz 为什么只碰接口、不碰 GORM/Redis。
  2. 照着代码画出 依赖注入装配图,讲明白 Google Wire 的 ProviderSet / wire.NewSet / wire_gen.go 是怎么在编译期把对象拼起来的,以及为什么不用运行时 reflect
  3. 讲清楚 Kratos App 的生命周期NewData 返回 cleanup 函数,进程退出时资源按什么顺序释放。
  4. 解释一个 .proto 文件怎么靠 google.api.http 注解,同时暴露 HTTP(8000) 和 gRPC(9000) 两套协议。
  5. etcd 服务注册发现NewEtcdRegistrar / RegisterTTL(15s) / Namespace 讲成面试故事,包括"为什么要心跳 TTL"。

前置知识(不用很熟,但得知道):

  • Go 基础:interface、构造函数、context.Contextdefer
  • 一点点依赖注入概念(知道"控制反转"这个词就行);
  • gRPC / Protobuf 大概长什么样(不用写过);
  • 听说过 MySQL、Redis、etcd 是干嘛的即可。

本章动手 3 件

  • 打开 internal/biz/repo.go,把 UserRepo 接口抄一遍,体会"业务只定义要什么,不关心怎么存";
  • 跑一遍 wire 生成:go run github.com/google/wire/cmd/wire,对比 wire.gowire_gen.go 的差异;
  • configs/config.yaml 里一个值用环境变量覆盖(比如 CORS_ALLOWED_ORIGINS),理解 envsource 的用途。

下面正式进入"面试问答"环节。每个知识点我先给一个生活类比建立直觉,再上本项目的真实代码。


Q1. 这个项目为什么用 DDD 四层(service → biz → data)?依赖倒置到底倒置了什么?

答:

先给个类比。想象你开一家餐厅:

  • 服务员(service):负责接待客人、点单、上菜,把客人的话翻译成后厨能懂的指令。
  • 厨师长(biz,业务逻辑层):只关心"做菜的流程和配方",他才不管土豆是从哪个菜市场买的、用哪口锅炒。
  • 采购/仓库(data,数据层):真正去菜市场买菜、把菜存进冷库。

关键点:厨师长不会亲自跑去菜市场。他只会在墙上贴一张"我要土豆、要牛肉"的采购单(接口),采购部照着单子去办。哪天换个供应商(MySQL 换成 Postgres,或者本地存储换成 OSS),厨师长完全不用改配方——这就是"依赖倒置(Dependency Inversion)"。

在本项目里,这条链路在 README.md 写得很清楚:

分层依赖:service → biz → data,禁止反向。biz 只依赖接口,实现在 data

注意这句话的"禁止反向"四个字。如果 biz 里出现了 gorm.DB 或者 redis.Client,那就破功了。我们把箭头画出来:

graph TD
    Client["客户端
浏览器 / gRPC 调用方"] --> SVC["service 层
proto ↔ biz 转换
CtxUserID 取用户"] SVC --> BIZ["biz 层
领域对象 + 用例
只依赖 Repo / Cache / Locker 接口"] BIZ -->|定义接口 实现在 data| DATA["data 层
GORM / Redis / Kafka / MinIO
实现 biz 定义的接口"] DATA --> DB[("MySQL")] DATA --> RDB[("Redis")] DATA --> MQ[("Kafka")] DATA --> OS[("本地 / MinIO 存储")]

图里那个 只依赖接口 的箭头是核心:biz 指向 data,但不是"依赖 data 的具体类",而是"data 去实现 biz 定义的接口"。方向被倒过来了——高层(biz)定义它需要什么(抽象),低层(data)去满足它。

真实代码佐证 —— internal/biz/repo.goUserRepo 就是那张"采购单":

// UserRepo 定义用户数据持久化的接口
// biz 层仅依赖此接口,不依赖任何具体实现
type UserRepo interface {
	Create(ctx context.Context, user *User) (*User, error)
	FindByID(ctx context.Context, id uint64) (*User, error)
	FindByUsername(ctx context.Context, username string) (*User, error)
	// ... 省略若干方法
	UpdatePassword(ctx context.Context, id uint64, hashedPassword string) error
	GetUserStorage(ctx context.Context, userID uint64) (total int64, used int64, err error)
}

biz 只声明"我要能 Create / FindByID / UpdatePassword",至于这些方法背后是 GORM 还是别的什么,biz 一个字都不关心。具体实现在 internal/data/user.go

// 编译期确保 userRepo 实现了 biz.UserRepo 接口。
var _ biz.UserRepo = (*userRepo)(nil)

type userRepo struct {
	db *gorm.DB
}

func NewUserRepo(db *gorm.DB) biz.UserRepo {
	return &userRepo{db: db}
}

⚠️ 新手坑:依赖倒置不是"把接口写在某层"就够了,必须保证 编译期 真正有人实现了它。本项目用 var _ biz.UserRepo = (*userRepo)(nil) 在编译期做了硬校验——下一题细讲。


Q2. biz 层"不依赖任何具体存储/框架,只依赖接口"具体是怎么解耦的?除了 Repo 还有哪些接口?

答:

光有一个 UserRepo 不够。一个真实的业务用例,往往还要:

  • 查缓存(Redis);
  • 加分布式锁(防并发改存储配额);
  • 把"文件删除"这类事件丢到消息队列(Kafka);
  • 把文件存到本地磁盘或对象存储(MinIO/OSS)。

如果 biz 直接 import "cloud-disk/internal/data/cache"import ".../data/lock",那 biz 就又和具体实现绑死了。本项目的解法很统一:凡是 biz 需要的"外部能力",都在 biz 包里自己定义成一个小接口,让 data / event 去实现。

来看这几个接口的真实定义。

① 缓存接口 Cache —— internal/biz/auth.go

// Cache 是 biz 层使用的缓存层接口
// 避免直接依赖 cache 包
type Cache interface {
	Get(ctx context.Context, key string) (string, error)
	Set(ctx context.Context, key string, value string, ttl time.Duration) error
	Delete(ctx context.Context, keys ...string) error
	Exists(ctx context.Context, key string) (bool, error)
	DeleteByPattern(ctx context.Context, pattern string) error
}

为什么这样设计?因为"Redis 缓存"只是 Cache 的一种实现。哪天你想换成本地内存缓存做单测,直接写一个 fakeCache 实现这个接口塞进去即可,biz 一行都不用动。这就是"面向接口编程"的红利。

② 分布式锁接口 Locker —— internal/biz/storage.go

// Locker 封装外部锁接口,用于并发安全的存储更新
type Locker interface {
	Lock(ctx context.Context, key string) error
	Unlock(ctx context.Context, key string) error
}

// NewLockAdapter 将兼容 lock.Lock 的实现封装为 biz.Locker
func NewLockAdapter(lockFn func(ctx context.Context, key string) (bool, error), unlockFn func(ctx context.Context, key string) error) Locker {
	return &lockFuncs{lock: lockFn, unlock: unlockFn}
}

注意 NewLockAdapter:data 层的 lock.Lock 签名是 Lock(ctx, key) (bool, error),而 biz 的 LockerLock(ctx, key) error。两者签名不完全一致,于是用一个"适配器(adapter)“把 bool 返回值吃掉。这就是经典的 适配器模式——接口不兼容时套一层,biz 永远只看自己想要的签名。

③ 存储接口 Storage —— internal/biz/storage.go:本地存储、MinIO、OSS 都去实现它(var _ biz.Storage = (*localStorage)(nil)(*minioStorage)(nil)(*ossStorage)(nil)),所以切存储驱动时 biz 无感。

④ 事件发布接口 EventPublisher —— internal/biz/event.go

// EventPublisher 定义业务层发布异步事件的接口,由 event 包实现。
// 通过接口注入避免 biz 与 event 包之间的循环依赖。
type EventPublisher interface {
	Publish(ctx context.Context, eventType string, payload interface{}) error
}

这里还有个额外好处:biz 和 event 包之间避免了循环依赖。因为谁都不直接 import 谁的实现,只依赖接口。event 包实现 EventPublisher,在 wire 装配时注入给 biz,依赖方向是单向往下的。

解耦的本质一句话:biz 是"需求方”,它把对外部世界的所有依赖都抽象成"我需要你能做 X/Y/Z"的小接口合同;具体谁来做、怎么做,由最外层的装配代码(wire)决定。这就是控制反转——“要什么"由 biz 说了算,“给什么"由装配器说了算。


Q3. 那个 var _ Repo = (*x)(nil) 到底在干嘛?为什么要在编译期校验接口实现?

答:

先说结论:这行代码不是运行时逻辑,它只是一个编译期断言,意思是"如果 (*userRepo) 没完全实现 biz.UserRepo 接口,编译直接报错”。

看真实代码(internal/data/user.go):

// 编译期确保 userRepo 实现了 biz.UserRepo 接口。
var _ biz.UserRepo = (*userRepo)(nil)

拆开看:

  • (*userRepo)(nil):声明一个值为 nil、类型为 *userRepo 的指针。不会真的分配内存,纯粹是类型层面的东西。
  • var _ biz.UserRepo = ...:把上面这个指针赋给 _(空白标识符,表示"我不在乎这个值,只为类型检查”),类型要求是 biz.UserRepo

如果 userRepo 漏实现了 UserRepo 里的任何一个方法(比如你重构时把 UpdatePassword 删了),Go 编译器会发现"类型不匹配",编译失败,而不是等到线上运行时才 panic。

为什么这么重要?回到 Q1 的依赖倒置:biz 定义接口,data 实现接口。一旦接口和实现分处两个包,Go 的接口是隐式实现(没有 implements 关键字),编译器不会强制你在 data 里"声明我实现了某个接口"。这就埋了个雷——某天有人改了 UserRepo 接口(加了个方法),却忘了去 data 补实现,编译居然能过,直到运行时调那个方法才炸。

var _ biz.UserRepo = (*userRepo)(nil) 就是人工补上这层"强约束"。本项目里到处都是这种断言,例如:

// internal/data/recycle.go
var _ biz.RecycleRepo = (*recycleRepo)(nil)
// internal/data/share.go
var _ biz.ShareRepo = (*shareRepo)(nil)
// internal/data/storage/local.go
var _ biz.Storage = (*localStorage)(nil)
// internal/event/handlers.go
var _ IdempotencyStore = (*RedisIdempotencyStore)(nil)

⚠️ 面试加分点:你可以补一句——Go 接口是"鸭子类型"的隐式实现,var _ Iface = (*Impl)(nil) 是把隐式变显式的最佳实践,任何"跨包实现接口"的地方都该加。它零运行时开销,纯编译期保险丝。


Q4. Google Wire 是什么?ProviderSet / wire.NewSet / wire_gen.go 是怎么协作的?为什么用编译期 DI 而不是运行时 reflect?

答:

生活类比:假设你要组装一台电脑。两种办法——

  • 运行时(reflect):拿到一堆零件,运行时用"反射"去猜每个插槽该插什么,插错了当场冒烟。Go 里典型的 uber/digjava spring 早期都偏这种,“类型安全弱、报错晚、性能差”。
  • 编译期(Wire):你先写一份"装机清单"(ProviderSet),Wire 在编译前就根据清单把零件和插槽对好,直接生成一份"已经插好线"的 wire_gen.go。运行时零反射、零猜疑,插错线编译都过不了。

Wire 的三个核心概念在本项目里全部能找到:

wire.NewSet 定义"我能提供哪些构造函数" —— 各层都有自己的 ProviderSet:

internal/biz/biz.go

// ProviderSet is biz providers.
var ProviderSet = wire.NewSet(NewTokenManager, NewUserUsecase, NewFileUsecase, NewRecycleUsecase, NewShareUsecase)

internal/data/data.go

// ProviderSet is data providers.
var ProviderSet = wire.NewSet(
	NewData,
	wire.FieldsOf(new(*Data), "DB", "RDB", "KafkaClient", "KafkaWriter"),
	wire.FieldsOf(new(*conf.Data), "Database", "Redis", "Kafka", "Scheduler", "Lock"),
	wire.FieldsOf(new(*conf.Storage), "Local"),
	NewUserRepo,
	NewGormFileRepo,
	NewGormFolderRepo,
	NewGormUploadRepo,
	NewFileAccessLogRepo,
	NewRecycleRepo,
	NewShareRepo,
)

internal/server/server.go

// ProviderSet is server providers.
var ProviderSet = wire.NewSet(NewGRPCServer, NewHTTPServer, NewEtcdClient, NewEtcdRegistrar)

internal/service/service.go

// ProviderSet is service providers.
var ProviderSet = wire.NewSet(NewUserService, NewFileService, NewRecycleService, NewShareService)

wire.Build 在入口声明"我要装配一个 App" —— cmd/server/wire.go

//go:build wireinject
// +build wireinject

func wireApp(*conf.Server, *conf.Data, *conf.Auth, *conf.Storage, *conf.Etcd, *slog.Logger) (*kratos.App, func(), error) {
	panic(wire.Build(server.ProviderSet, data.ProviderSet, cache.ProviderSet, lock.ProviderSet, taskProviderSet, storage.ProviderSet, event.ProviderSet, biz.ProviderSet, service.ProviderSet, newServers, newApp, newBizCache, newBizLocker))
}

注意 //go:build wireinject 这个 build tag:这个文件只给 Wire 工具看,正常编译会被排除,所以里面的 panic(wire.Build(...)) 永远不会真的执行。

wire_gen.go 是 Wire 生成的"已插好线"的代码 —— cmd/server/wire_gen.go 节选:

// Code generated by Wire. DO NOT EDIT.

func wireApp(confServer *conf.Server, confData *conf.Data, auth *conf.Auth, confStorage *conf.Storage, etcd *conf.Etcd, logger *slog.Logger) (*kratos.App, func(), error) {
	tokenManager := biz.NewTokenManager(auth)
	dataData, cleanup, err := data.NewData(confData)
	if err != nil {
		return nil, nil, err
	}
	db := dataData.DB
	userRepo := data.NewUserRepo(db)
	// ... 中间省略一众 repo / usecase / service 的构造
	userUsecase := biz.NewUserUsecase(userRepo, tokenManager, bizCache, locker)
	userService := service.NewUserService(userUsecase)
	// ... 继续向上装配
	grpcServer := server.NewGRPCServer(confServer, tokenManager, userService, fileService, recycleService, shareService)
	httpServer := server.NewHTTPServer(confServer, tokenManager, userService, fileService, recycleService, shareService, db, client)
	// ... etcd、app
	app := newApp(logger, v2, registrar)
	return app, func() {
		cleanup2()
		cleanup()
	}, nil
}

直觉上这就是 Wire 帮你把"先 new 谁、再传给谁"的顺序手写了一遍。Wire 做的事其实很"笨"但很靠谱:它读你所有的 ProviderSet,做依赖图拓扑排序,然后吐出普通 Go 代码。

为什么用编译期 DI 而非运行时 reflect? 面试标准答案:

  1. 类型安全:依赖对不上,编译就挂,而不是运行时 panic;
  2. 零运行时开销:生成的是普通函数调用,没有反射、没有 map 查找,性能等同手写;
  3. 可读/可调试wire_gen.go 是普通 Go 代码,go to definition 能跳进去,断点能打;
  4. 报错早、报错准:缺 provider、循环依赖,Wire 直接告诉你哪里缺,而不是线上炸。

⚠️ 新手坑wire_gen.go 头部写着 DO NOT EDIT!它必须靠 go run github.com/google/wire/cmd/wire 重新生成。手改它等于和工具作对——下次生成就覆盖了。改依赖要去对应的 ProviderSetwire.go 改,再重新 wire

新增一个业务模块要改哪几个 ProviderSet? 这是面试官最爱追问的。以"新增一个 Notify 通知模块"为例:

  1. internal/biz/biz.goProviderSetNewNotifyUsecase
  2. internal/service/service.goProviderSetNewNotifyService
  3. 如果 data 里要新写 repo,在 internal/data/data.goProviderSetNewNotifyRepo
  4. cmd/server/wire.gowire.Build(...) 通常不用改(因为各层 ProviderSet 都已经在列表里了,新增的构造函数在对应 ProviderSet 里就会被自动纳入拓扑);
  5. 重新跑 wire 生成 wire_gen.go

一句话:改的是各层自己的 ProviderSet,入口 wire.go 几乎不动


Q5. Kratos 的 App 生命周期是怎样的?NewData 返回的 cleanup 函数干嘛用?资源按什么顺序释放?

答:

生活类比:开餐厅前你得"办执照、租房子、接水电"(初始化);餐厅营业(Run);关门时你得"断水断电、退租、注销执照"(释放)。cleanup 就是那张"关门 Checklist"。Go 里靠 defer cleanup() 保证不管怎么退出都执行。

在本项目里,装配入口 cmd/server/main.go 的关键几行:

app, cleanup, err := wireApp(bc.Server, bc.Data, bc.Auth, bc.Storage, bc.Etcd, logger)
if err != nil {
	panic(err)
}
defer cleanup()     // 进程退出前,按顺序释放所有资源

if err := app.Run(); err != nil {   // 阻塞运行,直到收到退出信号
	panic(err)
}

wireApp 返回的 cleanup 是从哪里来的?看 internal/data/data.goNewData

// NewData 创建一个新的 Data 实例,将所有数据库连接组装在一起。
// 返回 Data 实例、一个清理函数以及可能遇到的错误。
func NewData(c *conf.Data) (*Data, func(), error) {
	db, err := NewGormDB(c.Database)
	if err != nil {
		return nil, nil, err
	}
	rdb, err := NewRedisClient(c.Redis)
	if err != nil {
		return nil, nil, err
	}
	kc, kw, err := NewKafkaClient(c.Kafka)
	if err != nil {
		return nil, nil, err
	}

	// 对所有领域模型执行自动迁移
	if db != nil {
		if err := model.AutoMigrate(db); err != nil {
			return nil, nil, err
		}
	}

	cleanup := func() {
		log.Info("closing the data resources")

		if db != nil {
			sqlDB, err := db.DB()
			if err == nil {
				if err := sqlDB.Close(); err != nil {
					log.Error("failed to close MySQL", "err", err)
				}
			}
		}
		if err := rdb.Close(); err != nil {
			log.Error("failed to close Redis", "err", err)
		}
		if err := kw.Close(); err != nil {
			log.Error("failed to close Kafka writer", "err", err)
		}
	}
	return &Data{DB: db, RDB: rdb, KafkaClient: kc, KafkaWriter: kw}, cleanup, nil
}

几个要点:

  • NewData 同时返回三样东西*Data(装好连接的结构体)、cleanup(关门清单)、error
  • model.AutoMigrate(db) 在连接建立后自动建表——这就是为什么 README 说"GORM AutoMigrate 会自动建表",你不用手建表。
  • cleanup 的释放顺序:MySQL → Redis → Kafka Writer。注意它和"创建顺序"(先 NewGormDB、再 NewRedisClient、再 NewKafkaClient)是逆序的,这是资源释放的常见好习惯——后申请的先释放。

wire_gen.go 里最终的 cleanup 是怎么组合的?看它返回的那段:

return app, func() {
	cleanup2()   // 先关 etcd 客户端
	cleanup()    // 再关 data(MySQL → Redis → Kafka)
}, nil

也就是说整体释放顺序是:etcd 客户端 → MySQL → Redis → Kafka Writer。原因也合理:etcd 注册器是"对外宣告我还活着"的组件,先把它摘掉(停止心跳),再慢慢关内部资源,避免"资源都快没了还对外宣称健康"的尴尬窗口。

⚠️ 新手坑cleanupfunc() 不是 func() error。所以内部每个 Close 自己吞错误、只打日志。如果你的资源释放必须感知错误(比如优雅刷盘),得在 cleanup 里自己处理,不能靠返回值往外抛。


Q6. 同一个 proto 文件,怎么同时暴露 HTTP(8000) 和 gRPC(9000) 两套协议?

答:

生活类比:你开一家店,既在街边有实体门面(HTTP,浏览器/前端直接访问),又有电话下单专线(gRPC,给其他微服务内部调用)。你只需要一份菜单(.proto),只是在菜单上额外标注"这道菜既能堂食也能外送"——google.api.http 注解就是那个标注。

本项目的 api/file/v1/file.proto 开头就 import "google/api/annotations.proto",然后每个 RPC 后面挂一段 option (google.api.http)

service FileService {
  // ListFiles lists files and folders under a parent directory with cursor-based pagination.
  rpc ListFiles(ListFilesRequest) returns (ListFilesReply) {
    option (google.api.http) = { get: "/api/v1/files" };
  }

  rpc CreateFolder(CreateFolderRequest) returns (CreateFolderReply) {
    option (google.api.http) = { post: "/api/v1/folders"; body: "*" };
  }

  rpc Download(DownloadRequest) returns (stream DownloadReply) {
    option (google.api.http) = { get: "/api/v1/file/download" };
  }

  rpc InitUpload(InitUploadRequest) returns (InitUploadReply) {
    option (google.api.http) = { post: "/api/v1/upload/init"; body: "*" };
  }
}

{ get: "/api/v1/files" } 意思是:这条 gRPC 方法同时映射成 GET /api/v1/files 这个 HTTP 路由。body: "*" 表示把整个 HTTP JSON body 反序列化成请求消息。

那 HTTP server 怎么知道这些路由的? 靠 protoc 生成时吐出的 RegisterXxxHTTPServer。看 internal/server/http.go

srv := khttp.NewServer(opts...)
// 注册服务路由
userv1.RegisterUserServiceHTTPServer(srv, us)
filev1.RegisterFileServiceHTTPServer(srv, fs)
recyclev1.RegisterRecycleServiceHTTPServer(srv, rs)
sharev1.RegisterShareServiceHTTPServer(srv, ss)

RegisterFileServiceHTTPServer 就是把 google.api.http 里的那些路由,连同"proto 消息 ↔ JSON"的转换逻辑一起,注册到 Kratos 的 HTTP server 上。这就是 grpc-gateway 思路在 Kratos 里的内建实现:一份 proto,两套传输,不用写两套 handler。

但本项目还有个"混合模式"值得讲:一部分路由走 proto 自动生成(如 /api/v1/files),另一部分因为要流式返回文件内容 / 手动控制 JWT,用了 srv.Route("/").Handle(...) 手动挂:

// 手动路由:流式返回文件预览内容
srv.Route("/").Handle("GET", "/api/v1/file/{id}/preview", getFilePreviewHandler(fs, tm))
// 健康检查端点
srv.Route("/").Handle("GET", "/healthz", healthLivenessHandler())
srv.Route("/").Handle("GET", "/readyz", healthReadinessHandler(db, rdb))

这正好引出下面这张双协议 + 混合路由的图:

graph LR
    FE["前端 React
:5173 代理 /api → :8000"] CLI["gRPC 调用方
:9000"] FE -->|HTTP 8000| HTTP["Kratos HTTP Server
RegisterXxxHTTPServer
+ 手动 Route"] CLI -->|gRPC 9000| GRPC["Kratos gRPC Server
RegisterXxxServer"] HTTP --> A["service 层
proto↔biz 转换"] GRPC --> A A --> B["biz 层
业务用例"] HTTP -.->|流式/特殊路由| H["手动 Handler
preview / healthz / readyz"] subgraph PROTO["api/*/v1/*.proto"] ANN["google.api.http 注解
get/post + body"] end ANN -.生成.-> HTTP

一句话总结:proto 是"单一事实来源(single source of truth)",HTTP 路由和 gRPC 方法都从它生成;特殊需求再用 srv.Route 手动补。

⚠️ 新手坑body: "*"body: "field" 不一样。前者整包映射,后者只把指定字段放 body、其余走 query/path。写错会导致前端 POST 的 JSON 字段对不上。另外流式接口(stream)在 HTTP 侧是 chunked 响应,网关/反代要允许长连接。


Q7. etcd 服务注册与发现:NewEtcdRegistrar、RegisterTTL(15s)、Namespace 是什么?为什么需要心跳 TTL?

答:

生活类比:微服务就像商场里一家家店铺。顾客(调用方)不可能记住每家店的具体门牌,于是商场入口有个电子导览牌(etcd),每家店开门时在牌上写"我在 3 楼 A12",关门就擦掉。但万一某店突然倒闭没人来擦牌(进程崩溃),导览牌还写着"它在",顾客就会白跑一趟。解决办法:要求每家店每 15 秒来"戳一下"导览牌(心跳 TTL),超过 15 秒没戳,牌自动熄灭。

本项目的注册逻辑在 internal/server/registry.go

// NewEtcdClient 根据配置创建 etcd 客户端。
func NewEtcdClient(c *conf.Etcd) (*clientv3.Client, func(), error) {
	if c == nil || len(c.Endpoints) == 0 {
		slog.Warn("etcd 配置为空,跳过 etcd 客户端创建")
		return nil, func() {}, nil
	}
	cfg := clientv3.Config{ Endpoints: c.Endpoints }
	// ... 用户名/密码/超时等按需填充
	client, err := clientv3.New(cfg)
	if err != nil {
		return nil, nil, err
	}
	return client, func() { _ = client.Close() }, nil
}

// NewEtcdRegistrar 创建 etcd 服务注册器。
func NewEtcdRegistrar(client *clientv3.Client) registry.Registrar {
	if client == nil {
		return nil
	}
	return etcd.New(
		client,
		etcd.Namespace("/microservices"),         // 注册路径前缀,做环境/租户隔离
		etcd.RegisterTTL(15 * time.Second),        // 心跳 TTL:15 秒不续约就自动摘除
	)
}

逐项解释:

  • etcd.Namespace("/microservices"):所有服务实例都注册在 /microservices/... 这个前缀下。相当于把"服务目录"和 etcd 里别的 key 隔开,也方便一套 etcd 多套环境共用时按前缀区分。
  • etcd.RegisterTTL(15 * time.Second):这是租约(lease)时长。Kratos 的 etcd 注册器会后台定时续租这个租约(心跳)。只要进程活着,key 就一直有效;一旦进程挂了不再续租,etcd 在 15 秒后自动删除这个 key,调用方就发现不了它了。
  • NewEtcdRegistrar 返回 nil 的优雅降级:当 etcd 没配置(c == nil || len(c.Endpoints) == 0),返回 nil。再看 main.gonewApp
func newApp(logger *slog.Logger, servers []transport.Server, registrar registry.Registrar) *kratos.App {
	opts := []kratos.Option{
		kratos.ID(id),
		kratos.Name(Name),
		// ...
		kratos.Server(servers...),
	}
	if registrar != nil {
		opts = append(opts, kratos.Registrar(registrar))  // 没配 etcd 就不注册
	}
	return kratos.New(opts...)
}

也就是说etcd 是可选能力:本地开发不配 etcd,应用照样跑,只是不在注册中心登记而已。这是个很务实的取舍。

为什么需要心跳 TTL(而不是直接写死一个 key)? 因为分布式系统里"优雅下线"是小概率,“进程被 kill -9 / 机器断电 / 网络分区"才是常态。没有 TTL,死掉的实例会永远躺在注册中心,调用方持续把流量打过去,全部超时。TTL 让"失效实例"能自动过期,是服务发现可用性的基石。

把整条"注册 + 心跳 + 发现"的时序画出来:

sequenceDiagram
    participant App as 云盘 App 实例
    participant Reg as NewEtcdRegistrar
    participant ETCD as etcd 集群
    participant Caller as 调用方(其他服务)

    App->>Reg: 启动后注册实例
/microservices/cloud-disk/实例ID Reg->>ETCD: put key + 15s 租约(lease) ETCD-->>Reg: 租约 ID loop 每 15 秒 心跳续租 Reg->>ETCD: keepAlive(续租) ETCD-->>Reg: 租约续期成功 end Caller->>ETCD: 查询 /microservices/cloud-disk ETCD-->>Caller: 返回健康实例列表 Note over App,ETCD: 进程崩溃 / 网络断开 App--xReg: 心跳停止 ETCD->>ETCD: 15s 后租约过期 ETCD-->>Caller: 列表里不再含该实例(自动摘除)

⚠️ 新手坑RegisterTTL 是"最大失联容忍时间”,不是"心跳间隔"。Kratos 内部会用一个比 TTL 短的间隔去续租。把 TTL 调太大(比如 5 分钟),故障时流量会错误路由很久;调太小(比如 1 秒),网络抖动就会误摘实例。15 秒是常见平衡点。


Q8. 配置中心:conf.pb.go 怎么来的?config.yaml 从哪来?envsource 环境变量覆盖有什么用?

答:

生活类比:程序启动时得先读"说明书"(配置)。这本说明书有两层——模板(conf.proto) 告诉你有哪些可配置项、各自什么类型;填好的纸(config.yaml) 才是真正的值。而 envsource 就像"橡皮擦+笔":允许你用环境变量临时改写纸上的某些值,最适合容器化部署(不同环境只改环境变量,不重打包)。

① conf.pb.go 由 conf.proto 生成:本项目 internal/conf 包(使用 conf.NewEnvExpandSource 等可见是 Kratos 的 config proto 体系)。conf.pb.go 是 protoc 把 conf.proto 编译出的 Go 结构体(Bootstrap 包含 Server / Data / Auth / Storage / Etcd 等)。wireApp 的参数类型 (*conf.Server, *conf.Data, *conf.Auth, *conf.Storage, *conf.Etcd, *slog.Logger) 就来自这里。

② config.yaml 是值来源README.md 给出关键片段:

data:
  database:
    source: "root:admin123@tcp(127.0.0.1:3306)/cloud_disk?parseTime=True&loc=Local"
auth:
  jwt_secret: "cloud-disk-jwt-secret-key-2024"
storage:
  driver: "local"          # 或 "oss"
  local:
    upload_dir: "./data/uploads"
recycle:
  auto_clean_days: 7

cmd/server/main.go 的加载入口:

func init() {
	flag.StringVar(&flagconf, "conf", "../../configs", "config path, eg: -conf config.yaml")
}

func main() {
	// ...
	c := config.New(
		config.WithSource(
			conf.NewEnvExpandSource(file.NewSource(flagconf)),  // 关键:envsource 包装
		),
	)
	defer c.Close()

	if err := c.Load(); err != nil {
		panic(err)
	}
	var bc conf.Bootstrap
	if err := c.Scan(&bc); err != nil {   // 反序列化到 conf.Bootstrap
		panic(err)
	}
	// ... wireApp(bc.Server, bc.Data, bc.Auth, bc.Storage, bc.Etcd, logger)
}

③ envsource 的作用conf.NewEnvExpandSource(file.NewSource(flagconf)) 这行的含义是——先读 flagconf 指向的 yaml 文件,再把内容里 ${ENV_VAR} 形式的占位符用环境变量展开。例如 yaml 里写:

auth:
  jwt_secret: "${JWT_SECRET}"

部署到容器时只注入环境变量 JWT_SECRET=xxx,不用把密钥写进镜像。这解决了"配置随环境变、密钥不能进代码仓库"的工程痛点。

graph TD
    PROTO["conf.proto
配置 schema 定义"] -->|protoc 生成| PB["conf.pb.go
Go 结构体 Bootstrap"] YAML["configs/config.yaml
各环境实际值"] --> SRC["file.NewSource(flagconf)"] SRC --> ENV["conf.NewEnvExpandSource
展开 ${ENV} 占位符"] ENV --> CFG["config.New + Load + Scan"] CFG --> BC["conf.Bootstrap
bc.Server/Data/Auth/Storage/Etcd"] BC --> WIRE["wireApp(...) 注入到各层"] PB -.结构体类型.-> BC

⚠️ 新手坑envsource 是"展开环境变量",不是"用环境变量覆盖整个 yaml 文件"。它只替换 yaml 里显式写了 ${XXX} 的地方。如果你想"整段配置用环境变量提供",得自己写 source 或者换配置中心(如 Nacos/Apollo)。另外 flagconf 默认是 ../../configs(目录),本地 go run ./cmd/server 时相对路径别搞错。


Q9. 面试必考题:要在这个项目里"加一个新功能模块",完整流程是怎样的?

答:

这是把前面所有知识点串起来的"综合题"。假设我们要加一个 “标签(Tag)模块”:给文件打标签、按标签检索。严格按本项目的分层,步骤是:

第 1 步:写 proto(单一事实来源)api/tag/v1/tag.proto 定义 TagService,每个 RPC 加 google.api.http 注解:

service TagService {
  rpc CreateTag(CreateTagRequest) returns (CreateTagReply) {
    option (google.api.http) = { post: "/api/v1/tags"; body: "*" };
  }
  rpc ListByTag(ListByTagRequest) returns (ListByTagReply) {
    option (google.api.http) = { get: "/api/v1/tags/{tag}/files" };
  }
}

make proto 生成 tag.pb.gotag_grpc.pb.go / tag_http.pb.go

第 2 步:biz 层——定义领域对象 + 用例 + repo 接口 internal/biz/tag.go

type Tag struct {
	ID   uint64
	Name string
}

// TagRepo 是 biz 定义的"标签持久化"接口(只声明要什么)
type TagRepo interface {
	Create(ctx context.Context, tag *Tag) (*Tag, error)
	ListFileIDsByTag(ctx context.Context, userID uint64, tag string) ([]uint64, error)
}

func NewTagUsecase(repo TagRepo, cache Cache) *TagUsecase {
	return &TagUsecase{repo: repo, cache: cache}
}

注意 TagRepoCache 都是接口——biz 不碰 GORM。

第 3 步:data 层——实现 repo 接口 internal/data/tag.go

var _ biz.TagRepo = (*tagRepo)(nil)   // 编译期校验

type tagRepo struct{ db *gorm.DB }

func NewTagRepo(db *gorm.DB) biz.TagRepo {
	return &tagRepo{db: db}
}
// ... 实现 Create / ListFileIDsByTag,内部用 GORM

第 4 步:service 层——proto ↔ biz 转换 internal/service/tag.go 实现生成的 TagServiceServer 接口,把 v1.CreateTagRequest 转成 biz.Tag 调 usecase,再把结果转回 v1.CreateTagReply。还要把 NewTagService 加进 service.ProviderSet

第 5 步:server 层——注册路由internal/server/http.goNewHTTPServer 里加 tagv1.RegisterTagServiceHTTPServer(srv, ts);gRPC 侧同理(本项目 gRPC 目前还是占位,但结构一致)。

第 6 步:wire 串起来

  • biz.ProviderSetNewTagUsecase
  • data.ProviderSetNewTagRepo
  • service.ProviderSetNewTagService
  • cmd/server/wire.gowire.Build(...) 列表一般不用改(各层 ProviderSet 已纳入拓扑),只需重跑 wire 生成 wire_gen.go

第 7 步:建表 model.AutoMigrate 里加上 Tag 模型,启动时自动建表。

把全流程画成一张"加模块数据流":

graph TD
    P["1. api/tag/v1/tag.proto
+ google.api.http"] -->|make proto| GEN["生成 pb.go / http.pb.go"] GEN --> SVC["4. service 层
实现 TagServiceServer
proto↔biz 转换"] P --> BIZ["2. biz 层
Tag 领域对象 + TagRepo 接口
+ NewTagUsecase"] BIZ -->|接口| DATA["3. data 层
tagRepo 实现 TagRepo
var _ biz.TagRepo"] DATA -->|NewTagRepo| W["6. wire 装配
各层 ProviderSet 加入
wire.go 几乎不动"] SVC --> W W --> APP["wireApp 生成
HTTP/gRPC 注册 + etcd"] DATA -->|model.AutoMigrate| DB[("MySQL 自动建表")]

一句话面试总结:“自外而内定义契约(proto → biz 接口),自内而外实现(data 实现 → service 适配 → wire 装配),依赖方向始终向下,接口是分层的接缝。”


Q10. 这个项目的分层有什么"违反"或取舍?service 层 CtxUserID 用裸字符串 key 的小坑,biz 用自定义 ctxKey 类型更安全——展开讲讲。

答:

没有架构是完美的,能讲清"我这个项目哪里不优雅、为什么",比背八股更显功力。本项目有两个典型取舍:

取舍一:service 层 CtxUserID 用裸字符串 key 取 context 值

internal/service/service.go

// CtxUserID 从 JWT 中间件设置的上下文中提取用户 ID。
func CtxUserID(ctx context.Context) uint64 {
	if id, ok := ctx.Value("user_id").(uint64); ok {   // 注意:裸字符串 "user_id"
		return id
	}
	if id, ok := ctx.Value("user_id").(string); ok {
		// ... 把字符串数字解析成 uint64
	}
	return 0
}

问题在于 Go 的 context.WithValue 的 key 是 interface{},用裸字符串 "user_id" 当 key,有两个隐患:

  1. 容易撞键:任何别的包只要也用 "user_id" 这个字符串往 ctx 里塞值,就会互相覆盖/误读。字符串是全局共享的,没有命名空间保护。
  2. 类型不安全ctx.Value("user_id") 返回 interface{},你得自己 .(uint64) / .(string) 断言,而且本项目 service 层还不得不兼容两种类型(uint64 和 string),这是历史包袱导致的别扭。

对比 biz / server 层的正确做法internal/server/middleware.go 用了自定义类型当 key:

// ctxKey 是 context 键的自定义类型,避免使用裸字符串导致与其他包冲突。
type ctxKey string

const (
	ctxKeyUserID   ctxKey = "user_id"
	ctxKeyUsername ctxKey = "username"
	ctxKeyToken    ctxKey = "token"
)

// JWT 中间件注入:
ctx = context.WithValue(ctx, ctxKeyUserID, claims.UserID)

这里 key 是 ctxKey 类型(底层是 string,但类型不同)。即使别处有另一个 type otherKey string; otherKey("user_id"),因为类型不一样,ctx.Value 也取不到彼此的值——类型系统天然做了命名空间隔离。这才是正确的 context key 写法。

顺带一提,server.CtxUserID(middleware.go 里那个)就是用的 ctxKeyUserID,而 service.CtxUserID 还在用裸 "user_id"。两者能"巧合地"互通,是因为当前实现里 JWT 中间件注入的就是裸字符串 "user_id"(middleware.go 第 138 行 context.WithValue(ctx, ctxKeyUserID, ...)——注意这里其实是用 ctxKeyUserID 注入的,而 service 层用裸字符串取,理论上取不到!)。这恰恰暴露了项目里存在两套取用户 ID 的方式、且 key 类型不一致的隐患,是面试时可以主动点出的"可改进点":统一成 server.CtxUserID 那种基于自定义类型 key 的版本,删掉 service 层那个裸字符串版本。

取舍二:gRPC server 目前还是"占位"

internal/server/grpc.go

func NewGRPCServer(c *conf.Server, tm *biz.TokenManager, us *service.UserService, fs *service.FileService, rs *service.RecycleService, ss *service.ShareService) *grpc.Server {
	var opts = []grpc.ServerOption{
		grpc.Middleware(recovery.Recovery()),
	}
	// ...
	_ = tm // gRPC 的 JWT 中间件将在后续迭代中添加
	_ = us // gRPC 服务注册将在后续迭代中添加
	return srv
}

us / fs 等参数传进来却用 _ 丢弃,说明 gRPC 侧的服务注册和 JWT 中间件还没接上。这是一个明确的"演进中"取舍:先把 HTTP 跑通,gRPC 留作内部服务间通信的后续能力。面试时诚实承认"这是 TODO",比假装完美更可信。

取舍三:etcd 可选 + 本地可无注册中心

前面 Q7 讲过,NewEtcdRegistrar 在没配置时返回 nilnewApp 据此跳过注册。这是"降低本地开发心智负担"的务实取舍,代价是生产环境必须显式配 etcd 才能真正享受服务发现。

把"正确做法 vs 本项目小坑"放一张对照图收尾:

graph TD
    BAD["service.CtxUserID
ctx.Value(\"user_id\") 裸字符串
风险: 撞键 / 类型不安全"] -->|改进| GOOD["server.CtxUserID
ctx.Value(ctxKeyUserID)
自定义类型 key, 命名空间隔离"] NOTE["现状: JWT 中间件注入用 ctxKeyUserID
service 层用裸字符串取 → 类型不一致隐患
建议统一到自定义类型 key"] BAD -.关联.-> NOTE GOOD -.关联.-> NOTE

⚠️ 面试加分点:能说出"context key 必须用自定义类型而非裸字符串"已经是及格线;能进一步指出"我们项目里 service 层还残留裸字符串版本、和 server 层不一致,是待清理的技术债",就是超出预期的回答。


自测题与动手练习

自测题(合上文章试试能不能答出来):

  1. DDD 四层里,箭头方向是 service → biz → data,为什么说这是"依赖倒置"?biz 定义的接口到底由谁来实现?
  2. var _ biz.UserRepo = (*userRepo)(nil) 这行代码运行时会有什么副作用?它防的是哪类 bug?
  3. Google Wire 的 ProviderSetwire.NewSetwire.Buildwire_gen.go 各自扮演什么角色?为什么说它是"编译期 DI"?
  4. NewData 返回的 cleanupwire_gen.go 里最终被组合成什么顺序释放资源?etcd 客户端和 MySQL/Redis/Kafka 谁先关?
  5. google.api.http 注解 + RegisterXxxHTTPServer 实现了什么能力?本项目的 HTTP server 为什么还有一部分路由是 srv.Route("/").Handle(...) 手动注册的?

动手练习:

  1. 画依赖图:不参考文章,凭记忆画出 UserRepobiz 定义到 data 实现、再到被 NewUserUsecase 注入的完整调用链,并在图上标出 var _ 校验的位置。
  2. 加一个最小的 ProviderSet 改动:假装要新增一个 Audit(审计)模块,写出你需要改动的 4 个文件(biz.go / data.go / service.go / 各自的 impl 文件),并说明 wire.go 要不要动、为什么。
  3. 修一个 context key 隐患:把 internal/service/service.go 里的 CtxUserID 改造成使用自定义 type ctxKey string 类型(参考 internal/server/middleware.go),并保证和 JWT 中间件注入的 key 类型一致,然后写一句注释说明为什么这样更安全。

本章小结

  • 分层与依赖倒置service → biz → data 单向依赖,biz 只定义接口(Repo / Cache / Locker / Storage / EventPublisher),data 实现,靠 var _ Iface = (*Impl)(nil) 在编译期锁死实现。
  • 依赖注入:Google Wire 用各层 ProviderSet + 入口 wire.Build 在编译期生成 wire_gen.go,零反射、类型安全、可调试;新增模块只改各层 ProviderSet,入口几乎不动。
  • 生命周期NewData 返回 cleanupmaindefer cleanup() 保证退出时按 etcd → MySQL → Redis → Kafka 顺序释放;model.AutoMigrate 自动建表。
  • 双协议:一份 .proto + google.api.http 注解同时生成 HTTP(8000) 与 gRPC(9000),特殊流式/探活路由用 srv.Route 手动补。
  • 服务治理NewEtcdRegistrar 配合 RegisterTTL(15s) + Namespace("/microservices") 实现带心跳 TTL 的注册发现,etcd 缺省时优雅降级;配置经 envsource 支持 ${ENV} 环境变量覆盖,适合容器化。
  • 取舍与改进:service 层 CtxUserID 用裸字符串 key 是与 server 层自定义类型 key 不一致的历史隐患,统一到自定义类型 key 更安全;gRPC 侧 JWT 与注册仍在 TODO,是明确的演进中状态。
复习提示:
  • DDD 分层的本质biz 层只定义接口(契约),data 层提供实现——这样 biz 不依赖任何具体技术栈,可以独立测试和演进。
  • Wire 生成代码的价值:不是"减少手写依赖注入代码",而是"编译期验证依赖图完整性"——漏掉的依赖在 wire_gen.go 生成时就会报错。
  • 生命周期管理:cleanup 函数按 etcd → MySQL → Redis → Kafka 顺序释放,反向关闭可能引发数据不一致(如 Kafka producer 已关但还在写 etcd)。
  • 双协议设计.proto + google.api.http 注解让一份定义同时服务 HTTP 和 gRPC,减少维护成本;但对特殊路由(流式、探活)仍需手动注册。
  • 下一篇讲数据库建模——DDD 的 data 层最终要落地到数据库,理解建模原则才能让分层真正发挥作用。
面试官
为什么 biz 层不能直接 import data 层?违反依赖倒置原则会有什么后果?
候选人
这是 DDD 最核心的规则之一:

违反后果
① biz 依赖了 Redis/MySQL/Kafka 等具体实现,无法单独测试(需要启动外部服务)
② 换数据库要改 biz 代码,违背开闭原则
③ 并行开发受阻:data 团队没完成前,biz 团队无法测试

正确做法
① biz 定义接口(如 UserRepo interface { FindByID(ctx, id) (*User, error) }
② data 实现接口(type userRepo struct { db *gorm.DB }
③ Wire 在入口组装依赖关系

面试加分点:提到"接口在 biz 里定义,实现在 data 里"这个反转是依赖倒置(DIP)的实践核心——高层模块(biz)不依赖低层模块(data),两者都依赖抽象(interface)。

下一章我们可以深入 biz 层的用例编排与事务边界——比如"上传文件"这种跨 repo、跨缓存、跨消息队列的操作,biz 怎么保证一致性,以及 Kafka 事件为什么需要幂等消费。那是另一块面试高频区。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!