当前位置: 首页 > news >正文

【ETCD】[源码阅读]深度解析 EtcdServer 的 processInternalRaftRequestOnce 方法

在分布式系统中,etcd 的一致性与高效性得益于其强大的 Raft 协议模块。而 processInternalRaftRequestOnce 是 etcd 服务器处理内部 Raft 请求的核心方法之一。本文将从源码角度解析这个方法的逻辑流程,帮助读者更好地理解 etcd 的内部实现。

方法源码

func (s *EtcdServer) processInternalRaftRequestOnce(ctx context.Context, r pb.InternalRaftRequest) (*applyResult, error) {ai := s.getAppliedIndex()ci := s.getCommittedIndex()if ci > ai+maxGapBetweenApplyAndCommitIndex {return nil, ErrTooManyRequests}r.Header = &pb.RequestHeader{ID: s.reqIDGen.Next(),}// check authinfo if it is not InternalAuthenticateRequestif r.Authenticate == nil {authInfo, err := s.AuthInfoFromCtx(ctx)if err != nil {return nil, err}if authInfo != nil {r.Header.Username = authInfo.Usernamer.Header.AuthRevision = authInfo.Revision}}data, err := r.Marshal()if err != nil {return nil, err}if len(data) > int(s.Cfg.MaxRequestBytes) {return nil, ErrRequestTooLarge}id := r.IDif id == 0 {id = r.Header.ID}ch := s.w.Register(id)cctx, cancel := context.WithTimeout(ctx, s.Cfg.ReqTimeout())defer cancel()start := time.Now()err = s.r.Propose(cctx, data)if err != nil {proposalsFailed.Inc()s.w.Trigger(id, nil) // GC waitreturn nil, err}proposalsPending.Inc()defer proposalsPending.Dec()select {case x := <-ch:return x.(*applyResult), nilcase <-cctx.Done():proposalsFailed.Inc()s.w.Trigger(id, nil) // GC waitreturn nil, s.parseProposeCtxErr(cctx.Err(), start)case <-s.done:return nil, ErrStopped}
}

方法解析

1. 校验状态与索引

ai := s.getAppliedIndex()
ci := s.getCommittedIndex()
if ci > ai+maxGapBetweenApplyAndCommitIndex {return nil, ErrTooManyRequests
}

getAppliedIndexgetCommittedIndex 分别获取当前节点的已应用索引和已提交索引。如果两者的差值过大,说明节点存在过多未应用的日志条目,可能导致性能问题,因此直接返回错误。

  • maxGapBetweenApplyAndCommitIndex:定义了允许的最大索引差距。
  • 防止机制:避免提交速度过快导致内存积压。

2. 生成请求头

r.Header = &pb.RequestHeader{ID: s.reqIDGen.Next(),
}

每个请求分配一个唯一的 ID,以便后续跟踪和处理。

3. 身份验证检查

if r.Authenticate == nil {authInfo, err := s.AuthInfoFromCtx(ctx)if err != nil {return nil, err}if authInfo != nil {r.Header.Username = authInfo.Usernamer.Header.AuthRevision = authInfo.Revision}
}
  • 目的:除认证请求外,其他请求需要验证用户身份。
  • 逻辑
    1. 调用 AuthInfoFromCtx 从上下文中提取用户身份。
    2. 将身份信息写入请求头,供后续处理。

4. 请求大小检查

if len(data) > int(s.Cfg.MaxRequestBytes) {return nil, ErrRequestTooLarge
}
  • 目的:防止超大请求导致内存或网络问题。
  • 机制:检查请求序列化后的大小是否超过配置的最大限制。

5. 注册请求等待通道

id := r.ID
if id == 0 {id = r.Header.ID
}
ch := s.w.Register(id)
  • 注册通道:使用请求 IDs.w(wait 组件)中注册一个等待通道,用于异步获取结果。

6. 发起 Raft 提案

cctx, cancel := context.WithTimeout(ctx, s.Cfg.ReqTimeout())
defer cancel()start := time.Now()
err = s.r.Propose(cctx, data)
  • 发起提案:调用 s.r.Propose 将请求数据交给 Raft 模块进行分布式一致性处理。
  • 超时控制:通过 Context.WithTimeout 设置提案的最大执行时间,避免长期阻塞。
  • 错误处理:如果提案失败,增加失败计数,并触发通道清理。

7. 等待提案结果

select {
case x := <-ch:return x.(*applyResult), nil
case <-cctx.Done():proposalsFailed.Inc()s.w.Trigger(id, nil) // GC waitreturn nil, s.parseProposeCtxErr(cctx.Err(), start)
case <-s.done:return nil, ErrStopped
}
  • 等待逻辑
    1. 通道 ch:正常返回应用结果。
    2. 上下文超时:处理超时错误,并清理等待通道。
    3. 服务关闭:直接返回停止错误。
  • 触发机制:使用 Trigger 清理通道,避免资源泄露。

8. 性能指标统计

  • proposalsPending.Inc():增加当前挂起的提案计数。
  • proposalsFailed.Inc():统计失败提案次数。

关键逻辑总结

processInternalRaftRequestOnce 方法的核心逻辑可分为以下几个阶段:

  1. 预检查:检查索引状态、请求大小和用户认证。
  2. 请求处理:序列化请求并将其提交到 Raft 模块。
  3. 结果等待:通过通道或超时控制获取提案的处理结果。

流程图

超出限制
正常
超出限制
正常
正常
超时
服务关闭
收到内部请求
检查已应用索引与已提交索引
返回 ErrTooManyRequests
生成请求头并检查认证信息
检查请求大小
返回 ErrRequestTooLarge
注册等待通道
调用 Raft 提案
等待结果
返回提案结果
返回超时错误
返回 ErrStopped

zz总结

processInternalRaftRequestOnce 是 etcd 服务端处理内部 Raft 请求的重要方法,它结合了请求校验、身份认证、Raft 提案以及结果返回的完整逻辑链条。理解其实现,可以帮助我们深入掌握 etcd 的核心一致性协议和服务端处理流程。

http://www.lryc.cn/news/502558.html

相关文章:

  • 【RabbitMQ】RabbitMQ中核心概念交换机(Exchange)、队列(Queue)和路由键(Routing Key)等详细介绍
  • 【AI知识】过拟合、欠拟合和正则化
  • 计算机毕设-基于springboot的航空散货调度系统的设计与实现(附源码+lw+ppt+开题报告)
  • 视图、转发与重定向、静态资源处理
  • 优选算法——分治(快排)
  • 【Linux系统】文件系统
  • javaweb的基础
  • 家里养几条金鱼比较好?
  • 写作词汇积累:差池、一体两面、切实可行极简理解
  • 移远EC200A-CN的OPENCPU使用GO开发嵌入式程序TBOX
  • LEED绿色建筑认证最新消息
  • SpringBoot中集成常见邮箱中容易出现的问题
  • webstorm开发uniapp(从安装到项目运行)
  • C# 探险之旅:第七节 - 条件判断(三元判断符):? : 的奇妙冒险
  • FlinkCDC实战:将 MySQL 数据同步至 ES
  • debug小记
  • Qt C++ 显示多级结构体,包括结构体名、变量名和值
  • 【JAVA】旅游行业中大数据的使用
  • 【AI+网络/仿真数据集】1分钟搭建云原生端到端5G网络
  • 微服务-01【续】
  • 测试工程师八股文01|Linux系统操作
  • 【Qt】qt基础
  • UniScene:Video、LiDAR 和Occupancy全面SOTA
  • TensorFlow深度学习实战(1)——神经网络与模型训练过程详解
  • 03篇--二值化与自适应二值化
  • 基于python的一个简单的压力测试(DDoS)脚本
  • 基于 Spring Boot 实现图片的服务器本地存储及前端回显
  • 深入 TCP VJ-Style
  • go高性能单机缓存项目
  • 数据结构绪论