← 返回首页
学习思考

MIT6.5840分布式系统 lab2:Key/Value Server

实验要求

在这个实验中,你将为单台机器构建一个键/值服务器,该服务器确保每个操作在网络故障的情况下只执行一次,并且这些操作是线性一致的。后续的实验将会复制这样的服务器以应对服务器崩溃的情况。

客户端可以向键/值服务器发送三种不同的RPC:Put(key, value)Append(key, arg)Get(key)。服务器维护一个内存中的键/值对映射。键和值都是字符串。Put(key, value)在映射中安装或替换特定键的值,Append(key, arg)将arg附加到键的值后并返回旧值,Get(key)获取键的当前值。对于不存在的键,Get应该返回一个空字符串。对于不存在的键,Append应该像现有值是零长度字符串一样操作。每个客户端通过具有Put/Append/Get方法的Clerk与服务器通信。Clerk管理与服务器的RPC交互。

你的服务器必须安排应用程序对ClerkGet/Put/Append方法调用是线性一致的。如果客户端请求不是并发的,每个客户端的Get/Put/Append调用应该观察到之前一系列调用所暗示的状态修改。对于并发调用,返回值和最终状态必须与操作按某种顺序一次执行的结果相同。如果调用在时间上重叠,则这些调用是并发的:例如,如果客户端X调用Clerk.Put(),客户端Y调用Clerk.Append(),然后客户端X的调用返回。一个调用必须观察到在调用开始之前完成的所有调用的效果。

线性一致性对于应用程序是方便的,因为它的行为就像一个单一的服务器一次处理一个请求。例如,如果一个客户端从服务器获得了一个更新请求的成功响应,随后其他客户端发起的读取请求保证能够看到该更新的效果。对于单个服务器,提供线性一致性相对容易。

Pt.1:Key/value server with no network failures

这部分很简单,我们只需要知道这个 GetPutAppend 的逻辑即可

go
func (kv *KVServer) Get(args *GetArgs, reply *GetReply) { // Your code here. kv.mu.Lock() defer kv.mu.Unlock() reply.Value = kv.data[args.Key] } func (kv *KVServer) Put(args *PutAppendArgs, reply *PutAppendReply) { // Your code here. kv.mu.Lock() defer kv.mu.Unlock() kv.data[args.Key] = args.Value } func (kv *KVServer) Append(args *PutAppendArgs, reply *PutAppendReply) { // Your code here. kv.mu.Lock() defer kv.mu.Unlock() oVal := kv.data[args.Key] kv.data[args.Key] += args.Value reply.Value = oVal }
go
func (ck *Clerk) Get(key string) string { // You will have to modify this function. args := GetArgs{ Key: key, Id: ck.id, } reply := GetReply{} ck.server.Call("KVServer.Get", &args, &reply) return reply.Value } func (ck *Clerk) PutAppend(key string, value string, op string) string { // You will have to modify this function. args := PutAppendArgs{ Key: key, Value: value, } reply := PutAppendReply{} ck.server.Call("KVServer."+op, &args, &reply) return reply.Value } func (ck *Clerk) Put(key string, value string) { ck.PutAppend(key, value, "Put") }

Pt.2:Key/value server with dropped messages

任务要求

现在,您应该修改您的解决方案,以便在遇到丢失的消息(例如 RPC 请求和 RPC 回复)时继续工作。如果消息丢失,则客户端的 ck.server.Call() 将返回 false (更准确地说, Call() 等待响应直至超市,如果在此时间内没有响应就返回false)。您将面临的一个问题是 Clerk 可能需要多次发送 RPC,直到成功为止。但是,每次调用 Clerk.Put() 或 Clerk.Append() 应该只会导致一次执行,因此您必须确保重新发送不会导致服务器执行请求两次。

你的任务是在 Clerk 中添加重试逻辑,并且在 server.go 中来过滤重复请求。

Hint: - 您需要唯一地标识client操作,以确保KV Server仅执行每个操作一次。 - 您必须仔细考虑server必须维持什么状态来处理重复的 Get() 、 Put() 和 Append() 请求(如果有话)。 - 您的重复检测方案应该快速释放服务器内存,例如让每个 RPC 暗示client已看到其前一个 RPC 的回复。可以假设client一次只向Clerk发起一次调用。

思考

首先,我们修改每个结构体,

go
package kvsrv // Put or Append type PutAppendArgs struct { Key string Value string // You'll have to add definitions here. // Field names must start with capital letters, // otherwise RPC will break. Id int64 Req int } type PutAppendReply struct { Value string } type GetArgs struct { Key string // You'll have to add definitions here. Id int64 Req int } type GetReply struct { Value string }
go
type Clerk struct { server *labrpc.ClientEnd // You will have to modify this struct. id int64 req int retryTime time.Duration requestTimeout time.Duration mtx sync.Mutex }
go
type Client struct { req int res string } type KVServer struct { mu sync.Mutex // Your definitions here. data map[string]string // kv数据存放处 client map[int64]Client // }

接着,为每个 rpc调用的部分添加超时检查,并补全参数的初始化

go
func (ck *Clerk) Get(key string) string { // You will have to modify this function. args := GetArgs{ Key: key, Id: ck.id, Req: ck.req, } reply := GetReply{} // 超时检查,如果超时进行告知并处理,防止一直占用内存 for { select { case <-time.After(ck.requestTimeout): { log.Fatal("get rpc timeout") return "" } default: ok := ck.server.Call("KVServer.Get", &args, &reply) if ok { // Get操作不需要进行修改,只需要保证Get在上一个PutAppend之后即可 return reply.Value } time.Sleep(ck.retryTime) } } return reply.Value }
go
func (ck *Clerk) PutAppend(key string, value string, op string) string { // You will have to modify this function. args := PutAppendArgs{ Key: key, Value: value, Id: ck.id, Req: ck.req, } reply := PutAppendReply{} for { select { case <-time.After(ck.requestTimeout): { log.Fatal("putAppend rpc timeout") return "" } default: ok := ck.server.Call("KVServer."+op, &args, &reply) if ok { // 记录操作的顺序 ck.req += 1; return reply.Value } time.Sleep(ck.retryTime) } } return reply.Value }

下边在 server 中,需要实现线性一致:

go
func (kv *KVServer) Put(args *PutAppendArgs, reply *PutAppendReply) { // Your code here. kv.mu.Lock() defer kv.mu.Unlock() client, _ := kv.client[args.Id] // 下边的判断条件是为了确保当前执行的操作在server上一个执行的操作之后 if args.Req >= client.req { client.req = args.Req + 1 // 更新序列号,保证后边的操作都在该操作之后 delete(kv.data, args.Key) // 删除原本的kv kv.data[args.Key] = args.Value kv.client[args.Id] = client } } func (kv *KVServer) Append(args *PutAppendArgs, reply *PutAppendReply) { // Your code here. kv.mu.Lock() defer kv.mu.Unlock() client, _ := kv.client[args.Id] if args.Req >= client.req { client.req = args.Req + 1 // 更新 client.res = kv.data[args.Key] kv.data[args.Key] += args.Value kv.client[args.Id] = client } reply.Value = client.res }

总结

本文由 GJJ 创作,内容来源于 Notion 数据库,随时可在 Notion 中编辑更新。 本站由 DeepSeek-v4-flash 辅助构建,项目参考 NotionNext

← 返回首页
61
文章
6
标签
3
分类
962
运行天数