惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

P
Palo Alto Networks Blog
Recent Commits to openclaw:main
Recent Commits to openclaw:main
C
CERT Recently Published Vulnerability Notes
C
Cybersecurity and Infrastructure Security Agency CISA
S
Schneier on Security
S
Securelist
酷 壳 – CoolShell
酷 壳 – CoolShell
C
CXSECURITY Database RSS Feed - CXSecurity.com
Cyberwarzone
Cyberwarzone
Apple Machine Learning Research
Apple Machine Learning Research
S
SegmentFault 最新的问题
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
GbyAI
GbyAI
Security Latest
Security Latest
Last Week in AI
Last Week in AI
Microsoft Security Blog
Microsoft Security Blog
云风的 BLOG
云风的 BLOG
Recorded Future
Recorded Future
Webroot Blog
Webroot Blog
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
TaoSecurity Blog
TaoSecurity Blog
C
Cisco Blogs
博客园 - 【当耐特】
Blog — PlanetScale
Blog — PlanetScale
Hugging Face - Blog
Hugging Face - Blog
B
Blog
Hacker News - Newest:
Hacker News - Newest: "LLM"
cs.CV updates on arXiv.org
cs.CV updates on arXiv.org
Attack and Defense Labs
Attack and Defense Labs
The Last Watchdog
The Last Watchdog
U
Unit 42
阮一峰的网络日志
阮一峰的网络日志
Project Zero
Project Zero
WordPress大学
WordPress大学
L
LINUX DO - 最新话题
F
Fortinet All Blogs
L
LINUX DO - 热门话题
PCI Perspectives
PCI Perspectives
Simon Willison's Weblog
Simon Willison's Weblog
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
MongoDB | Blog
MongoDB | Blog
Latest news
Latest news
P
Proofpoint News Feed
T
Threat Research - Cisco Blogs
The Hacker News
The Hacker News
爱范儿
爱范儿
O
OpenAI News
J
Java Code Geeks
T
The Exploit Database - CXSecurity.com
H
Hackread – Cybersecurity News, Data Breaches, AI and More

I'm OWenT

国产大模型(GLM 5.1、Kimi K2.6)真实场景效果和 Coding Plan 额度测试 新版本libatapp的连接管理——从etcd服务发现到拓扑驱动的自动重连 新版本libatbus的设计变更——从树形路由到拓扑驱动 Protobuf又一坑 - C++标准和ABI兼容性 AI真好用-给Blog主题统一加mermaid,chart.js,excalidraw,draw.io的多种引入方式支持 给内网部署Squid-通用HTTP下载缓存 UE使用CodeChecker和clang-tidy生成静态分析报告 找出UE的循环依赖 C++小协程栈和临时变量及作用域的栈溢出问题分析 游戏服务的可观测性能力建设(C++生态) 指标上报的多线程优化和多拉取源点优化 协程(libcopp)的Channel功能和CPU命中率优化 通用RPC代码生成器 实现strong_rc_ptr(比shared_ptr更快的引用计数智能指针) 手夯一个STL allocator和对象内存分析组件 踩坑一处(GCC)STL std::async 实现BUG导致的crash问题 GCC 14的一个warning to error BUG 给xresloader(Excel导表工具)增强UE读表支持(包含蓝图,Blueprint) Opentelemetry社区在gRPC的几个链接问题(静态库和动态库混用,musl工具链,符号裁剪) Excel转表工具(xresloader)的新验证器(验证外部Excel和文本数据,唯一性和自定义规则) protobuf v22和gRPC v1.55版本升级的依赖变化和upb适配 关于protobuf近期版本(v20/v3.20+)和 gRPC v1.54版本在某些编译环境下的一些链接和编译问题 xresloader-Excel导表工具链的近期变更汇总 打通游戏服务端框架的C++20协程改造的最后一环 Opentelemetry-cpp的Logs模块标准更新(涉及近期版本:1.8-1.9的BREAK CHANGES) 给cmake-toolset和工具链(curl等)加HTTP/2和HTTP/3支持 又开新坑之 coredns 插件: nftables和filter 关于opentelemetry-cpp社区对于C++ Head Only组件单例和符号可见性的讨论小记 填个转表工具 xresloader 去年的坑(数组尾部裁剪) 集成 upb 和 lua binding 的踩坑小记 libcopp对C++20协程的接入和接口设计 再度优化GCC、LLVM、Clang、libc++、libc++abi等套件的构建脚本 游戏服务的分布式事务优化(二)- 事务管理 游戏服务的分布式事务优化(一)- Write Ahead Log(WAL) 模块 记录一些bazel适配用编译选项 测试现代化硬件C++浮点数性能和一致性 适配Boringssl和OpenSSL 3.0 近期cmake-toolset的一些适配问题 C++20 Text Formatting/fmtlib 适配问题小记 再次重构LLVM+Clang+libcxx+libc++abi+其他相关工具的构建流程 重构基于CMake的构建工具链 新版GCC和LLVM+Clang终于Release啦 折腾一下nftables下的双拨 [C++20] Module partitions和符号交叉引用(声明和实现分离) [Rust] 实现一个线程安全且迭代器可以保存的链表 基于protobuf的代码生成 几个使用protobuf中C++接口的Arena的坑 Amazon Aurora DB存储引擎论文阅读小记 近期对libatapp的一些优化调整(增加服务发现和连接管理,支持yaml等) xresloader转表工具链增加了一些新功能(map,oneof支持,输出矩阵,基于模板引擎的加载代码生成等) 在游戏服务器中使用分布式事务 libcopp接入C++20 Coroutine和一些过渡期的设计 libatbus 的大幅优化 nftables初体验 容器配置开发环境小计 PALM Tree - 适合多核并发架构的B+树 - 论文阅读小记 跨平台协程库 - libcopp 简介 C++20 Coroutine 性能测试 (附带和libcopp/libco/libgo/goroutine/linux ucontext对比) 尝鲜Github Action 一些xresloader(转表工具)的改进 protobuf、flatbuffer、msgpack 针对小数据包的简单对比 协程框架(libcopp) 小幅优化 Excel转表工具(xresloader) 增加protobuf插件功能和集成 UnrealEngine 支持 Anna(支持任意扩展和超高性能的KV数据库系统)阅读笔记 C++20 Coroutine libcopp merge boost.context 1.69.0 Google去中心化分布式系统论文三件套(Percolator、Spanner、F1)读后感 Rust玩具-企业微信机器人通用服务 使用ELK辅助监控开发测试环境服务质量和问题定位 2018年的新通用伪随机数算法(xoshiro / xoroshiro)的C++(head only)实现 Webpack+vue+boostrap+ejs构建Web版GM工具 Rust的第二次接触-写个小服务器程序 理解和适配AEAD加密套件 atsf4g-co的进化:协程框架v2、对象路由系统和一些其他细节优化 协程框架(libcopp)v2优化、自适应栈池和同类库的Benchmark对比 可执行文件压缩 初识Rust 使用restructedtext编写xresloader文档 atframework的etcd模块化重构 C++的backtrace ECDH椭圆双曲线(比DH快10倍的密钥交换)算法简介和封装 protobuf-net的动态Message实现 pbc的proto3接入 atgateway内置协议流程优化-加密、算法协商和ECDH 整理一波软件源镜像同步工具+DevOps工具 Blog切换到Hugo libcopp v2的第一波优化完成 libcopp(v2) vs goroutine性能测试 libcopp的线程安全、栈池和merge boost.context 1.64.0 GCC 7和LLVM+Clang+libc++abi 4.0的构建脚本 libatbus的几个藏得很深的bug 用cmake交叉编译到iOS和Android 开源项目得一些小维护 atapp的c binding和c#适配 对象路由系统设计 2016年总结 近期的一个协程流程BUG 重写了llvm+clang+libc++和libc++abi的构建脚本 atsf4g完整游戏工程示例|I'm OWenT atframework基本框架已经完成|I'm OWenT
std::condition_variable 的信号丢失问题
owent · 2024-08-02 · via I'm OWenT

blog-website

背景

这篇分享拖更了好久了。问题起源于去年我们项目组接入 opentelemetry-cpp 的时候,在进程优雅退出的时候偶现超时,虽然可以直接kill进程没啥影响但是退出不“优雅”的话总归会破坏发布流程,增加人工介入的成本。这里记录一下问题可能其他的组件有类似的用法也会有相似的问题。

相关代码和使用场景

首先介绍下代码逻辑的背景,在代码结构层面。首先是有一个后台线程处理Batch Process和导出任务。然后caller线程发起 ForceFlush 或者 Shutdown 的时候要等待已经发出的事件执行完毕。 他们之间通过 std::condition_variable 来通知事件完成。

大致代码如下,首先是 bool BatchSpanProcessor::ForceFlush(std::chrono::microseconds timeout) noexcept 接口:

bool BatchSpanProcessor::ForceFlush(std::chrono::microseconds timeout) noexcept
{
  // ...

  // Now wait for the worker thread to signal back from the Export method
  std::unique_lock<std::mutex> lk_cv(synchronization_data_->force_flush_cv_m);

  synchronization_data_->is_force_flush_pending.store(true, std::memory_order_release);
  auto break_condition = [this]() {
    if (synchronization_data_->is_shutdown.load() == true)
    {
      return true;
    }

    // Wake up the worker thread once.
    if (synchronization_data_->is_force_flush_pending.load(std::memory_order_acquire))
    {
      synchronization_data_->is_force_wakeup_background_worker.store(true,
                                                                     std::memory_order_release);
      synchronization_data_->cv.notify_one();
    }

    return synchronization_data_->is_force_flush_notified.load(std::memory_order_acquire);
  };

  // Fix timeout to meet requirement of wait_for
  timeout = opentelemetry::common::DurationUtil::AdjustWaitForTimeout(
      timeout, std::chrono::microseconds::zero());
  bool result;
  if (timeout <= std::chrono::microseconds::zero())
  {
    bool wait_result = false;
    while (!wait_result)
    {
      // When is_force_flush_notified.store(true) and force_flush_cv.notify_all() is called
      // between is_force_flush_pending.load() and force_flush_cv.wait(). We must not wait
      // for ever
      wait_result = synchronization_data_->force_flush_cv.wait_for(lk_cv, schedule_delay_millis_,
                                                                   break_condition);
    }
    result = true;
  }
  else
  {
    result = synchronization_data_->force_flush_cv.wait_for(lk_cv, timeout, break_condition);
  }

  // ...
  return result;
}

然后 void BatchSpanProcessor::DoBackgroundWork() 和相关代码:

void BatchSpanProcessor::DoBackgroundWork()
{
  auto timeout = schedule_delay_millis_;

  while (true)
  {
    // Wait for `timeout` milliseconds
    std::unique_lock<std::mutex> lk(synchronization_data_->cv_m);
    synchronization_data_->cv.wait_for(lk, timeout, [this] {
      if (synchronization_data_->is_force_wakeup_background_worker.load(std::memory_order_acquire))
      {
        return true;
      }

      return !buffer_.empty();
    });
    synchronization_data_->is_force_wakeup_background_worker.store(false,
                                                                   std::memory_order_release);

    if (synchronization_data_->is_shutdown.load() == true)
    {
      DrainQueue();
      return;
    }

    auto start = std::chrono::steady_clock::now();
    Export();
    auto end      = std::chrono::steady_clock::now();
    auto duration = std::chrono::duration_cast<std::chrono::milliseconds>(end - start);

    // Subtract the duration of this export call from the next `timeout`.
    timeout = schedule_delay_millis_ - duration;
  }
}

void BatchSpanProcessor::Export()
{
  do
  {
    // 省略无关代码 ...
    exporter_->Export(nostd::span<std::unique_ptr<Recordable>>(spans_arr.data(), spans_arr.size()));
    NotifyCompletion(notify_force_flush, synchronization_data_);
  } while (true);
}

void BatchSpanProcessor::NotifyCompletion(
    bool notify_force_flush,
    const std::shared_ptr<SynchronizationData> &synchronization_data)
{
  if (!synchronization_data)
  {
    return;
  }

  if (notify_force_flush)
  {
    synchronization_data->is_force_flush_notified.store(true, std::memory_order_release);
    synchronization_data->force_flush_cv.notify_one();
  }
}

简单描述一下就是在退出阶段,首先要设置shutdown flag,然后唤醒后台线程导出数据。之后用一个 std::condition_variable 来等待后台线程通知导出完成。然后自己再退出,时序大概如下:

sequenceDiagram
    Caller->>Caller: 采集当前待上报边界
    Background->>Background: 启动
    Caller->>Background: 通知唤醒后台线程(condition_variable-1)
    Caller->>Caller: 进入等待上报完成通知(condition_variable-2)
    loop 循环执行到退出
    Background->>Background: 导出一批数据
    Background->>Caller: 通知已经导出的数据边界(condition_variable-2)
    Background->>Background: 没有数据则进入等待(condition_variable-1)
    end
    loop 直到当前上报边界已经全部上报完成或超时
      Caller->>Caller: 等待上报完成唤醒(condition_variable-2)
      Caller->>Caller: 检查上报边界或超时
    end
    Caller->>Caller: 退出
    Background->>Background: 后台线程退出

问题分析

上面的流程和代码乍一看似乎没有什么问题,但是在实际项目中偶尔会出现在 BatchSpanProcessor::ForceFlush 里的 synchronization_data_->force_flush_cv.wait_for 陷入了长耗时的等待。有的版本甚至会陷入无尽得等待。

因为在退出流程里,上层要保证最后的数据执行过刷出。并且后台线程不再引用资源之后才能销毁内存对象,所以会ForceFlush 一次,且 schedule_delay_millis_ 一般情况下考虑网络抖动会设置得比较长。某些流程里也有可能会在reload的时候切换到另一个新线程上调用 ForceFlush(max()) 的,这时候如果发生退出也会发生某些Provider一直退不掉的问题。

在后台任务刷出数据完成后,设置完 is_force_flush_notified 并且 force_flush_cv.notify_one() 后下一此循环会发现 is_shutdown 已经是 true 了,就会退出后台线程,之后就不再有 force_flush_cv.notify_one() 了。

那么在调用 ForceFlush() 的线程里 force_flush_cv.wait_for(lk_cv, schedule_delay_millis_, break_condition);break_condition 包含 is_force_flush_notified.load(std::memory_order_acquire); 为什么返回 false 且后面还没收到 force_flush_cv.notify_one() 的通知呢?

我们先来看一看 std::condition_variable::wait_for 的实现。(几个主流STL库实现大同小异,这里贴Linux下GCC的实现作为案例)

template<typename _Rep, typename _Period, typename _Predicate>
  bool
  wait_for(unique_lock<mutex>& __lock,
      const chrono::duration<_Rep, _Period>& __rtime,
      _Predicate __p)
  {
using __dur = typename steady_clock::duration;
return wait_until(__lock,
    steady_clock::now() +
    chrono::__detail::ceil<__dur>(__rtime),
    std::move(__p));
  }

template<typename _Clock, typename _Duration, typename _Predicate>
  bool
  wait_until(unique_lock<mutex>& __lock,
  const chrono::time_point<_Clock, _Duration>& __atime,
  _Predicate __p)
  {
while (!__p())
if (wait_until(__lock, __atime) == cv_status::timeout)
  return __p();
return true;
  }

  // 其他的 wait_until 嵌套都是设计模式相关,不在展示,最后是调用下面的代码,pthread的接口
  void
  wait_until(mutex& __m, clockid_t __clock, timespec& __abs_time)
  {
    pthread_cond_clockwait(&_M_cond, __m.native_handle(), __clock,
          &__abs_time);
  }

可以看到,这里关键点在这个代码:

while (!__p())
if (wait_until(__lock, __atime) == cv_status::timeout)
  return __p();
return true;

可能会存在一个极小的临界区,在调用 __p() 的时候 is_force_flush_notified.store(true, std::memory_order_release)force_flush_cv.notify_one() 还没执行到,然后准备进入 wait_until(__lock, __atime), 但是在最终调用 pthread_cond_clockwaitis_force_flush_notified.store(true, std::memory_order_release)force_flush_cv.notify_one() 执行掉了。然后再进入wait就没收到 notify ,并且由于退出阶段,后台线程已经退出不会再收到新的notify,所以必须等待到超时。如果这时候超时时间很长就会陷入一个非常长时间的等待。

解决方案

找到问题之后解决方法就比较简单了。主要有两种,一种是再加一个 mutex 锁,把 is_force_flush_notifiedforce_flush_cvstd::condition_variable::wait_for 也放进这个锁的临界区里保护,这样会多一个锁,多一些额外开销。另一种方案则是缩短等待时间后重试,缺点是发生这个情况的时候可能要多等个这个等待时间后才能完成。

因为这个毕竟是在退出阶段才出现且很小概率出现,所以我选了后一种不带来额外开销的方案。PR去年已经合入,小伙伴门可以放心使用。

后续问题

后面其实还有线程安全的问题也和这个相关,虽然我们项目里没用到,但是社区收到了issue。因为上面的 is_force_flush_notified 是个atomic的bool值,在多线程调用 ForceFlush() 的时候有概率某个线程的 is_force_flush_notified 会被其他线程吞掉。相关详情见: [SDK] BatchSpanProcessor::ForceFlush appears to be hanging forever[SDK] BatchLogRecordProcessor::ForceFlush hangs for 10 seconds 。 现在的解决法方法是不再使用bool值判定多次调用的先后关系,而是使用exported序号来进行。PR见: [SDK] Fix forceflush may wait for ever

欢迎有兴趣的小伙伴互相交流研究。