爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      influxdb集群写数据writeToShard解析

      首页 知识中心 大数据 文章详情页

      influxdb集群写数据writeToShard解析

      2023-08-07 07:03:07 阅读次数:433

      influxdb,集群

      摘要:

      解析influxdb集群写数据writeToShard, 分析出influxdb集群如何处理数据.

      源码:

      顶层writeToShard

      // writeToShards writes points to a shard and ensures a write consistency level has been met.  If the write
      // partially succeeds, ErrPartialWrite is returned.
      func (w *PointsWriter) writeToShard(shard *meta.ShardInfo, database, retentionPolicy string,
      consistency ConsistencyLevel, points []models.Point) error {
      atomic.AddInt64(&w.stats.PointWriteReqLocal, int64(len(points)))

      // The required number of writes to achieve the requested consistency level
      required := len(shard.Owners)
      switch consistency {
      case ConsistencyLevelAny, ConsistencyLevelOne:
      required = 1
      case ConsistencyLevelQuorum:
      required = required/2 + 1
      }

      // response channel for each shard writer go routine
      type AsyncWriteResult struct {
      Owner meta.ShardOwner
      Err error
      }
      ch := make(chan *AsyncWriteResult, len(shard.Owners))

      for _, owner := range shard.Owners {
      go func(shardID uint64, owner meta.ShardOwner, points []models.Point) {
      if w.Node.ID == owner.NodeID {
      atomic.AddInt64(&w.stats.PointWriteReqLocal, int64(len(points)))
      err := w.TSDBStore.WriteToShard(shardID, points)
      // If we've written to shard that should exist on the current node, but the store has
      // not actually created this shard, tell it to create it and retry the write
      if err == tsdb.ErrShardNotFound {
      err = w.TSDBStore.CreateShard(database, retentionPolicy, shardID, true)
      if err != nil {
      ch <- &AsyncWriteResult{owner, err}
      return
      }
      err = w.TSDBStore.WriteToShard(shardID, points)
      }
      ch <- &AsyncWriteResult{owner, err}
      return
      }

      atomic.AddInt64(&w.stats.PointWriteReqRemote, int64(len(points)))
      err := w.ShardWriter.WriteShard(shardID, owner.NodeID, points)
      if err != nil && tsdb.IsRetryable(err) {
      // The remote write failed so queue it via hinted handoff
      atomic.AddInt64(&w.stats.WritePointReqHH, int64(len(points)))
      hherr := w.HintedHandoff.WriteShard(shardID, owner.NodeID, points)
      if hherr != nil {
      ch <- &AsyncWriteResult{owner, hherr}
      return
      }

      // If the write consistency level is ANY, then a successful hinted handoff can
      // be considered a successful write so send nil to the response channel
      // otherwise, let the original error propagate to the response channel
      if hherr == nil && consistency == ConsistencyLevelAny {
      ch <- &AsyncWriteResult{owner, nil}
      return
      }
      }
      ch <- &AsyncWriteResult{owner, err}

      }(shard.ID, owner, points)
      }

      var wrote int
      timeout := time.After(w.WriteTimeout)
      var writeError error
      for range shard.Owners {
      select {
      case <-w.closing:
      return ErrWriteFailed
      case <-timeout:
      atomic.AddInt64(&w.stats.WriteTimeout, 1)
      // return timeout error to caller
      return ErrTimeout
      case result := <-ch:
      // If the write returned an error, continue to the next response
      if result.Err != nil {
      atomic.AddInt64(&w.stats.WriteErr, 1)
      w.Logger.Info("write failed", zap.Uint64("shard", shard.ID), zap.Uint64("node", result.Owner.NodeID), zap.Error(result.Err))

      // Keep track of the first error we see to return back to the client
      if writeError == nil {
      writeError = result.Err
      }
      continue
      }

      wrote++

      // We wrote the required consistency level
      if wrote >= required {
      atomic.AddInt64(&w.stats.WriteOK, 1)
      return nil
      }
      }
      }

      if wrote > 0 {
      atomic.AddInt64(&w.stats.WritePartial, 1)
      return ErrPartialWrite
      }

      if writeError != nil {
      return fmt.Errorf("write failed: %v", writeError)
      }

      return ErrWriteFailed

      }

      向远端shard写入:

      // WriteShard writes time series points to a shard
      func (w *ShardWriter) WriteShard(shardID, ownerID uint64, points []models.Point) error {
      c, err := w.dial(ownerID)
      if err != nil {
      return err
      }

      conn, ok := c.(*pooledConn)
      if !ok {
      panic("wrong connection type")
      }
      defer func(conn net.Conn) {
      conn.Close() // return to pool
      }(conn)

      // Determine the location of this shard and whether it still exists
      db, rp, sgi := w.MetaClient.ShardOwner(shardID)
      if sgi == nil {
      // If we can't get the shard group for this shard, then we need to drop this request
      // as it is no longer valid. This could happen if writes were queued via
      // hinted handoff and we're processing the queue after a shard group was deleted.
      return nil
      }

      // Build write request.
      var request WriteShardRequest
      request.SetShardID(shardID)
      request.SetDatabase(db)
      request.SetRetentionPolicy(rp)
      request.AddPoints(points)

      // Marshal into protocol buffers.
      buf, err := request.MarshalBinary()
      if err != nil {
      return err
      }

      // Write request.
      conn.SetWriteDeadline(time.Now().Add(w.timeout))
      if err := WriteTLV(conn, writeShardRequestMessage, buf); err != nil {
      conn.MarkUnusable()
      return err
      }

      // Read the response.
      conn.SetReadDeadline(time.Now().Add(w.timeout))
      _, buf, err = ReadTLV(conn)
      if err != nil {
      conn.MarkUnusable()
      return err
      }

      // Unmarshal response.
      var response WriteShardResponse
      if err := response.UnmarshalBinary(buf); err != nil {
      return err
      }

      if response.Code() != 0 {
      return fmt.Errorf("error code %d: %s", response.Code(), response.Message())
      }

      return nil
      }

      寻找到shard的映射:

      // MapShards maps the points contained in wp to a ShardMapping.  If a point
      // maps to a shard group or shard that does not currently exist, it will be
      // created before returning the mapping.
      func (w *PointsWriter) MapShards(wp *WritePointsRequest) (*ShardMapping, error) {
      rp, err := w.MetaClient.RetentionPolicy(wp.Database, wp.RetentionPolicy)
      if err != nil {
      return nil, err
      } else if rp == nil {
      return nil, freetsdb.ErrRetentionPolicyNotFound(wp.RetentionPolicy)
      }

      // Holds all the shard groups and shards that are required for writes.
      list := make(sgList, 0, 8)
      min := time.Unix(0, models.MinNanoTime)
      if rp.Duration > 0 {
      min = time.Now().Add(-rp.Duration)
      }

      for _, p := range wp.Points {
      // Either the point is outside the scope of the RP, or we already have
      // a suitable shard group for the point.
      if p.Time().Before(min) || list.Covers(p.Time()) {
      continue
      }

      // No shard groups overlap with the point's time, so we will create
      // a new shard group for this point.
      sg, err := w.MetaClient.CreateShardGroup(wp.Database, wp.RetentionPolicy, p.Time())
      if err != nil {
      return nil, err
      }

      if sg == nil {
      return nil, errors.New("nil shard group")
      }
      list = list.Append(*sg)
      }

      mapping := NewShardMapping(len(wp.Points))
      for _, p := range wp.Points {
      sg := list.ShardGroupAt(p.Time())
      if sg == nil {
      // We didn't create a shard group because the point was outside the
      // scope of the RP.
      mapping.Dropped = append(mapping.Dropped, p)
      atomic.AddInt64(&w.stats.WriteDropped, 1)
      continue
      }

      sh := sg.ShardFor(p.HashID())
      mapping.MapPoint(&sh, p)
      }
      return mapping, nil
      }

      时序图:

      2022-02-14 influxdb集群写数据writeToShard解析

      分析:

      从逻辑流程上:

      1. 到该函数时, 已经完成了key -> shard之间映射关系的查找
      1. 需要注意shard中的每个shard都开创一个协程去处理
      1. 如果shard是自己, 调用tsdb写
      1. 向远端写, 如果写失败
      2. 写入内存中的hinted handoff
      1. 收集处理结果

      从架构上:

      1. 可以看出将key换算成shardID, 分布到整个集群中存储
      2. 一旦远端shard写失败, 使用HH保存在本shard
      3. 需要注意该函数对协程的使用技巧
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.51cto.com/adofsauron/5644255,作者:帝尊悟世,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:raid5数据恢复成功案例

      下一篇:openssl关于RSA算法

      相关文章

      2025-05-19 09:05:01

      【Linux】HDP集群日志配置和日志删除脚本

      HDP 集群 默认安装的,日志放在数据盘,但是 namenode和snamenode的数据盘本身不大只有 500G,在不经意间 数据盘被日志装满,首先从集群配置着手。

      2025-05-19 09:05:01
      log4j , 日志 , 集群
      2025-05-19 09:04:44

      spark控制台没显示其他机器

      spark控制台没显示其他机器

      2025-05-19 09:04:44
      Spark , 节点 , 集群
      2025-05-14 10:02:48

      YARN与HBase任务

      YARN与HBase任务

      2025-05-14 10:02:48
      HBase , 任务 , 应用程序 , 资源 , 集群
      2025-05-06 09:19:21

      【Linux 从基础到进阶】Kubernetes 集群搭建与管理

      Kubernetes(简称 K8s)是一个用于自动化部署、扩展和管理容器化应用程序的开源平台。它提供了容器编排功能,能够管理大量的容器实例,并支持应用的自动扩展、高可用性和自愈能力。

      2025-05-06 09:19:21
      Kubernetes , Pod , 容器 , 节点 , 集群
      2025-05-06 09:19:12

      redis高可用集群搭建

      redis高可用集群搭建

      2025-05-06 09:19:12
      master , redis , 服务器 , 节点 , 集群
      2025-03-28 06:55:13

      配置集群免密登录

      配置集群免密登录

      2025-03-28 06:55:13
      登录 , 节点 , 集群
      2025-03-17 07:49:59

      redis-cluster分布式集群安装部署

      redis-cluster分布式集群安装部署

      2025-03-17 07:49:59
      redis , Ruby , 安装 , 实例 , 集群
      2025-03-12 09:32:22

      大数据平台的运维与管理技巧

      大数据平台的运维与管理技巧

      2025-03-12 09:32:22
      平台 , 故障 , 数据 , 运维 , 集群
      2025-02-18 07:29:23

      【图论】【 割边】【C++算法】1192. 查找集群内的关键连接

      【图论】【 割边】【C++算法】1192. 查找集群内的关键连接

      2025-02-18 07:29:23
      https , 服务器 , 算法 , 连接 , 集群
      2024-12-27 08:00:32

      oracle 11.2.0.4 asm单实例不随系统启动而自动开启

      oracle 11.2.0.4 asm单实例不随系统启动而自动开启

      2024-12-27 08:00:32
      Oracle , 开启 , 自动 , 集群
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5239678

      查看更多

      最新文章

      使用outflux 导入influxdb 的数据到timescaledb

      2024-05-14 08:53:33

      griddb 集群大小评估算法

      2023-04-07 06:47:51

      查看更多

      热门文章

      griddb 集群大小评估算法

      2023-04-07 06:47:51

      使用outflux 导入influxdb 的数据到timescaledb

      2024-05-14 08:53:33

      查看更多

      热门标签

      算法 leetcode python 数据 java 数组 节点 大数据 i++ 链表 golang c++ 排序 django 数据类型
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      使用outflux 导入influxdb 的数据到timescaledb

      griddb 集群大小评估算法

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号