爆款云主机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云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      2022-03-16 k8s的operator接收数据到数据队列的过程

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

      2022-03-16 k8s的operator接收数据到数据队列的过程

      2023-02-23 09:20:15 阅读次数:498

      epoll,queue,operator,k8s,sed

      摘要:

      operator使用多路复用从socket接收数据后写入队列, 由独立的处理协程读出数据处理.

      本文记录operator读取数据到数据队列的过程

      调用堆栈:

      (gdb) bt
      #0  k8s.io/client-go/util/workqueue.(*Type).Add (q=0xc00007e420, item=...) at /root/work/hello/vendor/k8s.io/client-go/util/workqueue/queue.go:121
      #1  0x0000000000b9139d in k8s.io/client-go/util/workqueue.(*delayingType).Add (.this=0xc00007e6c0, item=...) at <autogenerated>:1
      #2  0x0000000000b919dd in k8s.io/client-go/util/workqueue.(*rateLimitingType).Add (.this=0xc00013c240, item=...) at <autogenerated>:1
      #3  0x00000000019ce55a in sigs.k8s.io/controller-runtime/pkg/handler.(*EnqueueRequestForObject).Create (e=0x2b9eb70 <runtime.zerobase>, evt=..., q=...)
          at /root/work/hello/vendor/sigs.k8s.io/controller-runtime/pkg/handler/enqueue.go:44
      #4  0x0000000001a18e59 in sigs.k8s.io/controller-runtime/pkg/source/internal.EventHandler.OnAdd (e=..., obj=...) at /root/work/hello/vendor/sigs.k8s.io/controller-runtime/pkg/source/internal/eventsource.go:63
      #5  0x0000000001a1a08f in sigs.k8s.io/controller-runtime/pkg/source/internal.(*EventHandler).OnAdd (.this=0xc0001de7c0, obj=...) at <autogenerated>:1
      #6  0x00000000019b65b8 in k8s.io/client-go/tools/cache.(*processorListener).run.func1 () at /root/work/hello/vendor/k8s.io/client-go/tools/cache/shared_informer.go:787
      #7  0x0000000000b8631c in k8s.io/apimachinery/pkg/util/wait.BackoffUntil.func1 (f={void (void)} 0xc00006dd88) at /root/work/hello/vendor/k8s.io/apimachinery/pkg/util/wait/wait.go:155
      #8  0x0000000000b860c9 in k8s.io/apimachinery/pkg/util/wait.BackoffUntil (f={void (void)} 0xc00006de48, backoff=..., sliding=true, stopCh=0xc0000222a0)
          at /root/work/hello/vendor/k8s.io/apimachinery/pkg/util/wait/wait.go:156
      #9  0x0000000000b85fa5 in k8s.io/apimachinery/pkg/util/wait.JitterUntil (f={void (void)} 0xc00006dea0, period=1000000000, jitterFactor=0, sliding=true, stopCh=0xc0000222a0)
          at /root/work/hello/vendor/k8s.io/apimachinery/pkg/util/wait/wait.go:133
      #10 0x0000000000b85ed3 in k8s.io/apimachinery/pkg/util/wait.Until (f={void (void)} 0xc00006ded8, period=1000000000, stopCh=0xc0000222a0) at /root/work/hello/vendor/k8s.io/apimachinery/pkg/util/wait/wait.go:90
      #11 0x00000000019b645a in k8s.io/client-go/tools/cache.(*processorListener).run (p=0xc00021a280) at /root/work/hello/vendor/k8s.io/client-go/tools/cache/shared_informer.go:781
      #12 0x00000000019bcf4b in k8s.io/client-go/tools/cache.(*processorListener).run-fm () at /root/work/hello/vendor/k8s.io/client-go/tools/cache/shared_informer.go:775
      #13 0x0000000000b85df8 in k8s.io/apimachinery/pkg/util/wait.(*Group).Start.func1 () at /root/work/hello/vendor/k8s.io/apimachinery/pkg/util/wait/wait.go:73
      #14 0x000000000046afe1 in runtime.goexit () at /usr/local/go/src/runtime/asm_amd64.s:1581
      #15 0x0000000000000000 in ?? ()

      核心函数:

      wait.Until

      // Until loops until stop channel is closed, running f every period.
      //
      // Until is syntactic sugar on top of JitterUntil with zero jitter factor and
      // with sliding = true (which means the timer for period starts after the f
      // completes).
      func Until(f func(), period time.Duration, stopCh <-chan struct{}) {
        JitterUntil(f, period, 0.0, true, stopCh)
      }
      // JitterUntil loops until stop channel is closed, running f every period.
      //
      // If jitterFactor is positive, the period is jittered before every run of f.
      // If jitterFactor is not positive, the period is unchanged and not jittered.
      //
      // If sliding is true, the period is computed after f runs. If it is false then
      // period includes the runtime for f.
      //
      // Close stopCh to stop. f may not be invoked if stop channel is already
      // closed. Pass NeverStop to if you don't want it stop.
      func JitterUntil(f func(), period time.Duration, jitterFactor float64, sliding bool, stopCh <-chan struct{}) {
        BackoffUntil(f, NewJitteredBackoffManager(period, jitterFactor, &clock.RealClock{}), sliding, stopCh)
      }
      // BackoffUntil loops until stop channel is closed, run f every duration given by BackoffManager.
      //
      // If sliding is true, the period is computed after f runs. If it is false then
      // period includes the runtime for f.
      func BackoffUntil(f func(), backoff BackoffManager, sliding bool, stopCh <-chan struct{}) {
        var t clock.Timer
        for {
          select {
          case <-stopCh:
            return
          default:
          }
      
          if !sliding {
            t = backoff.Backoff()
          }
      
          func() {
            defer runtime.HandleCrash()
            f()
          }()
      
          if sliding {
            t = backoff.Backoff()
          }
      
          // NOTE: b/c there is no priority selection in golang
          // it is possible for this to race, meaning we could
          // trigger t.C and stopCh, and t.C select falls through.
          // In order to mitigate we re-check stopCh at the beginning
          // of every loop to prevent extra executions of f().
          select {
          case <-stopCh:
            if !t.Stop() {
              <-t.C()
            }
            return
          case <-t.C():
          }
        }
      }

      processorListener:run

      func (p *processorListener) run() {
        // this call blocks until the channel is closed.  When a panic happens during the notification
        // we will catch it, **the offending item will be skipped!**, and after a short delay (one second)
        // the next notification will be attempted.  This is usually better than the alternative of never
        // delivering again.
        stopCh := make(chan struct{})
        wait.Until(func() {
          for next := range p.nextCh {
            switch notification := next.(type) {
            case updateNotification:
              p.handler.OnUpdate(notification.oldObj, notification.newObj)
            case addNotification:
              p.handler.OnAdd(notification.newObj)
            case deleteNotification:
              p.handler.OnDelete(notification.oldObj)
            default:
              utilruntime.HandleError(fmt.Errorf("unrecognized notification: %T", next))
            }
          }
          // the only way to get here is if the p.nextCh is empty and closed
          close(stopCh)
        }, 1*time.Second, stopCh)
      }

      EventHandler:OnAdd

      // OnAdd creates CreateEvent and calls Create on EventHandler.
      func (e EventHandler) OnAdd(obj interface{}) {
        c := event.CreateEvent{}
      
        // Pull Object out of the object
        if o, ok := obj.(client.Object); ok {
          c.Object = o
        } else {
          log.Error(nil, "OnAdd missing Object",
            "object", obj, "type", fmt.Sprintf("%T", obj))
          return
        }
      
        for _, p := range e.Predicates {
          if !p.Create(c) {
            return
          }
        }
      
        // Invoke create handler
        e.EventHandler.Create(c, e.Queue)
      }

      EnqueueRequestForObject:Create

      // Create implements EventHandler.
      func (e *EnqueueRequestForObject) Create(evt event.CreateEvent, q workqueue.RateLimitingInterface) {
        if evt.Object == nil {
          enqueueLog.Error(nil, "CreateEvent received with no metadata", "event", evt)
          return
        }
        q.Add(reconcile.Request{NamespacedName: types.NamespacedName{
          Name:      evt.Object.GetName(),
          Namespace: evt.Object.GetNamespace(),
        }})
      }

      sharedIndexInformer:AddEventHandlerWithResyncPeriod

      func (s *sharedIndexInformer) AddEventHandlerWithResyncPeriod(handler ResourceEventHandler, resyncPeriod time.Duration) {
        s.startedLock.Lock()
        defer s.startedLock.Unlock()
      
        if s.stopped {
          klog.V(2).Infof("Handler %v was not added to shared informer because it has stopped already", handler)
          return
        }
      
        if resyncPeriod > 0 {
          if resyncPeriod < minimumResyncPeriod {
            klog.Warningf("resyncPeriod %v is too small. Changing it to the minimum allowed value of %v", resyncPeriod, minimumResyncPeriod)
            resyncPeriod = minimumResyncPeriod
          }
      
          if resyncPeriod < s.resyncCheckPeriod {
            if s.started {
              klog.Warningf("resyncPeriod %v is smaller than resyncCheckPeriod %v and the informer has already started. Changing it to %v", resyncPeriod, s.resyncCheckPeriod, s.resyncCheckPeriod)
              resyncPeriod = s.resyncCheckPeriod
            } else {
              // if the event handler's resyncPeriod is smaller than the current resyncCheckPeriod, update
              // resyncCheckPeriod to match resyncPeriod and adjust the resync periods of all the listeners
              // accordingly
              s.resyncCheckPeriod = resyncPeriod
              s.processor.resyncCheckPeriodChanged(resyncPeriod)
            }
          }
        }
      
        listener := newProcessListener(handler, resyncPeriod, determineResyncPeriod(resyncPeriod, s.resyncCheckPeriod), s.clock.Now(), initialBufferSize)
      
        if !s.started {
          s.processor.addListener(listener)
          return
        }
      
        // in order to safely join, we have to
        // 1. stop sending add/update/delete notifications
        // 2. do a list against the store
        // 3. send synthetic "Add" events to the new handler
        // 4. unblock
        s.blockDeltas.Lock()
        defer s.blockDeltas.Unlock()
      
        s.processor.addListener(listener)
        for _, item := range s.indexer.List() {
          listener.add(addNotification{newObj: item})
        }
      }
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.51cto.com/adofsauron/5644372,作者:帝尊悟世,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:2022-04-01 访问k8s内的etcd的数据

      下一篇:数据加密

      相关文章

      2025-05-09 08:50:35

      STL:Stack和Queue的模拟实现

      适配器是一种设计模式(设计模式是一套被反复使用的、多数人知晓的、经过分类编目的、代码设计经验的总结),该种模式是将一个类的接口转换成客户希望的另外一个接口。

      2025-05-09 08:50:35
      deque , queue , stack , 元素 , 容器 , 底层 , 适配器
      2025-05-07 09:09:52

      【C++/STL】stack/queue的使用及底层剖析&&双端队列&&容器适配器

      适配器是一种设计模式(设计模式是一套被反复使用的、多数人知晓的、经过分类编目的、代码设计经验的总结),该种模式是将一个类的接口转换成客户希望的另外一个接口。

      2025-05-07 09:09:52
      deque , queue , stack , STL , 容器
      2025-04-22 09:28:31

      Tcp的三次握手及netty和实际开发如何设置全连接队列参数

      Tcp的三次握手及netty和实际开发如何设置全连接队列参数

      2025-04-22 09:28:31
      queue , server , 队列
      2025-04-22 09:28:19

      【C++】stack、queue的模拟实现与deque的介绍

      【C++】stack、queue的模拟实现与deque的介绍

      2025-04-22 09:28:19
      deque , queue , stack , 容器 , 适配器
      2025-04-18 07:11:19

      k8s不同role级别的服务发现

      k8s不同role级别的服务发现

      2025-04-18 07:11:19
      k8s , pod , prometheus , role
      2025-04-18 07:11:11

      使用k8s的sdk编写一个项目获取pod和node信息

      使用k8s的sdk编写一个项目获取pod和node信息

      2025-04-18 07:11:11
      client , k8s , node , pod
      2025-04-18 07:11:11

      k8s-apiserver监控源码解读

      本节重点介绍 :k8s代码库和模块地址下载 apiserver源码apiserver中监控源码阅读k8s源码地址分布k8s代码库访问github上k8s仓库,readme中给出了k8s 模块的代码地址举例图片组件仓库列表 地址Reposit

      2025-04-18 07:11:11
      io , k8s , 源码
      2025-04-18 07:10:38

      shell编程-sed命令详解(超详细)

      在Shell编程中,对文本进行处理和转换是一项常见的任务。sed命令作为一种流式文本编辑器,提供了强大的文本处理能力,可以通过简单的命令实现复杂的文本操作。掌握sed命令的基本用法,可以极大地提高文本处理的效率和灵活性。

      2025-04-18 07:10:38
      sed , 命令 , 插入 , 文件 , 文本 , 替换
      2025-04-16 09:26:39

      为什么说k8s中监控更复杂了

      为什么说k8s中监控更复杂了

      2025-04-16 09:26:39
      k8s , 对象 , 监控
      2025-04-16 09:26:39

      prometheus为k8s做的4大适配工作

      prometheus为k8s做的4大适配工作

      2025-04-16 09:26:39
      k8s , pod , 指标 , 组件 , 适配
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5245858

      查看更多

      最新文章

      【epoll】epoll的水平触发和边沿触发,及为什么边沿触发必须使用非阻塞?

      2025-02-26 07:20:49

      k8s 数据卷需要很长时间才能挂载成功

      2024-05-16 09:52:01

      给定一个数组arr,代表每个人的能力值。再给定一个非负数k,如果两个人能力差值正好为k,那么可以凑在一起比赛。一局比赛只有两个人,返回最多可以同时有多少场比赛。

      2024-05-08 07:02:21

      给定一个正数数组arr,代表若干人的体重。再给定一个正数limit,表示所有船共同拥有的载重量。每艘船最多坐两人,且不能超过载重,想让所有的人同时过河,并且用最好的分配方法让船尽

      2024-04-18 09:15:34

      sed 命令详解(增删该查)

      2024-03-28 08:17:27

      2022-04-01 访问k8s内的etcd的数据

      2023-02-23 07:38:36

      查看更多

      热门文章

      2022-04-01 访问k8s内的etcd的数据

      2023-02-23 07:38:36

      sed 命令详解(增删该查)

      2024-03-28 08:17:27

      给定一个正数数组arr,代表若干人的体重。再给定一个正数limit,表示所有船共同拥有的载重量。每艘船最多坐两人,且不能超过载重,想让所有的人同时过河,并且用最好的分配方法让船尽

      2024-04-18 09:15:34

      k8s 数据卷需要很长时间才能挂载成功

      2024-05-16 09:52:01

      给定一个数组arr,代表每个人的能力值。再给定一个非负数k,如果两个人能力差值正好为k,那么可以凑在一起比赛。一局比赛只有两个人,返回最多可以同时有多少场比赛。

      2024-05-08 07:02:21

      【epoll】epoll的水平触发和边沿触发,及为什么边沿触发必须使用非阻塞?

      2025-02-26 07:20:49

      查看更多

      热门标签

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

      相关产品

      弹性云主机

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

      天翼云电脑(公众版)

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

      对象存储

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

      云硬盘

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

      查看更多

      随机文章

      k8s 数据卷需要很长时间才能挂载成功

      给定一个正数数组arr,代表若干人的体重。再给定一个正数limit,表示所有船共同拥有的载重量。每艘船最多坐两人,且不能超过载重,想让所有的人同时过河,并且用最好的分配方法让船尽

      给定一个数组arr,代表每个人的能力值。再给定一个非负数k,如果两个人能力差值正好为k,那么可以凑在一起比赛。一局比赛只有两个人,返回最多可以同时有多少场比赛。

      【epoll】epoll的水平触发和边沿触发,及为什么边沿触发必须使用非阻塞?

      sed 命令详解(增删该查)

      2022-04-01 访问k8s内的etcd的数据

      • 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号