Hi everyone,
I recently added async to my channel, but I’m seeing a ~2.5x performance drop:
Sync: 4 × 10M msg in 172 ms (~231M msg/s)
Async: 4 × 10M msg in 500 ms (~78M msg/s)
I expected async to be faster due to lower context switching overhead. I’ve already tried removing the backoff logic in poll_send (returning Poll::Ready / Poll::Pending directly), but the result remained the same.
What typically causes this kind of overhead in async channels? Any insights or suggestions would be greatly appreciated
Here is the async :
Atomic waker :
use core::task::Waker;
use core::sync::atomic::{AtomicBool, Ordering};
use core::cell::UnsafeCell;
#[cfg_attr(feature = "cache_64", repr(align(64)))]
#[cfg_attr(feature = "cache_128", repr(align(128)))]
pub struct AtomicWaker {
waker: UnsafeCell<Option<Waker>>,
locked: AtomicBool,
has_waker: AtomicBool,
}
impl AtomicWaker {
pub const fn new() -> Self {
Self {
waker: UnsafeCell::new(None),
locked: AtomicBool::new(false),
has_waker: AtomicBool::new(false),
}
}
pub fn register(&self, waker: &Waker) {
while self.locked.compare_exchange(
false, true, Ordering::Acquire, Ordering::Relaxed
).is_err() {
core::hint::spin_loop();
}
unsafe {
let slot = &mut *self.waker.get();
match slot {
Some(existing) if existing.will_wake(waker) => {}
_ => {
*slot = Some(waker.clone());
self.has_waker.store(true, Ordering::Relaxed);
}
}
}
self.locked.store(false, Ordering::Release);
}
pub fn wake(&self) {
if !self.has_waker.load(Ordering::Acquire) {
return;
}
while self.locked.compare_exchange(
false, true, Ordering::Acquire, Ordering::Relaxed
).is_err() {
core::hint::spin_loop();
}
let waker = unsafe { (*self.waker.get()).take() };
self.locked.store(false, Ordering::Release);
if let Some(w) = waker {
w.wake();
}
}
}
unsafe impl Send for AtomicWaker {}
unsafe impl Sync for AtomicWaker {}
The async component :
#[cfg(feature = "async")]
pub struct SendReady<'a, T, const CAPACITY_P2: u8, const BATCH: u8> {
ring: &'a Ring<T, CAPACITY_P2, BATCH>,
}
#[cfg(feature = "async")]
impl<'a, T, const CAPACITY_P2: u8, const BATCH: u8> Future
for SendReady<'a, T, CAPACITY_P2, BATCH>
{
type Output = &'a mut [MaybeUninit<T>];
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
unsafe {
match self.ring.poll_send(cx) {
Poll::Ready(slot) => {
Poll::Ready(&mut *(slot as *mut _))
}
Poll::Pending => Poll::Pending,
}
}
}
}
#[inline(always)]
pub fn get_all<F>(&self, mut handler: F) -> usize
where
F: FnMut(&T),
{
let head = self.consumer.head.load(Ordering::Relaxed) as usize;
let tail = self.producer.tail.load(Ordering::Acquire) as usize;
let ready_to_read = tail.wrapping_sub(head);
if ready_to_read == 0 {
return 0;
}
let buf = unsafe { slice::from_raw_parts(self.buf_ptr.as_ptr() as *const MaybeUninit<T>, Self::CAPACITY) };
let pairs = ready_to_read / 2;
let mut pos = head;
for _ in 0..pairs {
let id0 = pos & Self::MASK;
let id1 = pos.wrapping_add(1) & Self::MASK;
let next_id = pos.wrapping_add(2) & Self::MASK;
unsafe {
crate::prefetch_read(self.buf_ptr.as_ptr().add(next_id) as *const T);
handler(buf.get_unchecked(id0).assume_init_ref());
handler(buf.get_unchecked(id1).assume_init_ref());
}
pos = pos.wrapping_add(2);
}
if ready_to_read & 1 != 0 {
let id = pos & Self::MASK;
unsafe { handler(buf.get_unchecked(id).assume_init_ref()); }
pos = pos.wrapping_add(1);
}
self.consumer.head.store(pos as u64, Ordering::Release);
unsafe { *self.consumer.cached_tail.get() = tail as u64; }
#[cfg(not(feature = "async"))]
if self.producer.futex.swap(FUTEX_OPEN) == FUTEX_WAITING {
self.producer.futex.wake_one();
}
pos.wrapping_sub(head)
}
#[cfg(feature = "async")]
pub fn get_all_async<F>(&self, handler: F) -> usize
where
F: FnMut(&T),
{
let n = self.get_all(handler);
if n > 0 {
self.producer.atomic_waker.wake();
}
n
}
#[cfg(feature = "async")]
pub fn poll_send(
&self,
cx: &mut Context<'_>,
) -> Poll<&mut [MaybeUninit<T>]> {
if let Some(slot) = self.sender() {
return Poll::Ready(slot);
}
loop {
let mut spin = 6;
let outer = 12;
for _ in 0..outer {
for _ in 0..spin {
core::hint::spin_loop();
}
if let Some(b) = self.sender() {
return Poll::Ready(b);
}
spin <<= 1;
/*unsafe {
libc::sched_yield();
}*/
self.producer.atomic_waker.register(cx.waker());
match self.sender() {
Some(slot) => return Poll::Ready(slot),
None => return Poll::Pending,
}
}
}
}
The benchmarker :
const BATCH_SIZE: u8 = 1;
let messages_per_producer: usize = 10_000_000;
let ring_bits = 12u32;
let max_producers = 4usize;
let total_messages_expected = max_producers * messages_per_producer;
// heap async benchmark
{
let channel = std::sync::Arc::new(heap::Channel::<u64, 12, BATCH_SIZE>::new(max_producers));
let source = std::sync::Arc::new(AtomicU64::new(0));
let total = std::sync::Arc::new(AtomicU64::new(0));
let done_producers = std::sync::Arc::new(AtomicUsize::new(0));
let start_time = Instant::now();
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(max_producers + 1)
.build()
.unwrap();
rt.block_on(async {
let mut producer_handles = Vec::new();
for id in 0..max_producers {
let channel = channel.clone();
let source = source.clone();
let done_producers = done_producers.clone();
producer_handles.push(tokio::spawn(async move {
let sender = channel.sender().unwrap();
let mut sent: u64 = 0;
let limit = messages_per_producer as u64;
let mut local_sum: u64 = 0;
while sent < limit {
let res = sender.send_ready().await;
let n = res.len().min((limit - sent) as usize);
for (j, slot) in res[..n].iter_mut().enumerate() {
let val = sent + j as u64;
slot.write(val);
local_sum = local_sum.wrapping_add(val);
}
sender.send_n(n as u64);
sent += n as u64;
}
source.fetch_add(local_sum, Ordering::Relaxed);
done_producers.fetch_add(1, Ordering::Release);
}));
}
let channel_recv = channel.clone();
let done_producers_recv = done_producers.clone();
let total_recv = total.clone();
let consumer_handle = tokio::spawn(async move {
let rings: Vec<_> = (0..max_producers)
.map(|i| channel_recv.receiver(i).unwrap())
.collect();
let mut local_total: u64 = 0;
let mut total_processed = 0;
loop {
let mut got = 0;
for r in &rings {
got += r.get_all_async(|val| {
local_total = local_total.wrapping_add(*val);
});
}
total_processed += got;
if got == 0 {
tokio::task::yield_now().await;
}
if done_producers_recv.load(Ordering::Acquire) == max_producers
&& rings.iter().all(|r| r.is_empty())
{
break;
}
}
total_recv.store(local_total, Ordering::Relaxed);
assert_eq!(total_processed, total_messages_expected, "heap async: wrong message count");
});
for h in producer_handles {
h.await.unwrap();
}
consumer_handle.await.unwrap();
});
let duration = start_time.elapsed();
channel.close();
let s = source.load(Ordering::Relaxed);
let t = total.load(Ordering::Relaxed);
println!("── heap async ──");
println!("source: {}, total: {}, match: {}", s, t, s == t);
println!(" Ring Bits : {}", ring_bits);
println!(" Ring Capacity : {} elemen", 1u32 << ring_bits);
println!(" Jumlah Producer : {}", max_producers);
println!(" Total Messages : {} item", total_messages_expected);
println!(" Waktu Eksekusi : {:?}", duration);
println!(" Throughput : {:.2} msg/sec", total_messages_expected as f64 / duration.as_secs_f64());
println!();
}
For the full code is in here : GitHub - fuji-184/F_Channel · GitHub
What is the causes of the performance drop in async?
I also have tried to remove the backoff logic in poll_send, to return Poll::Ready or Poll::Pending directly, the result is still same

