You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Crystal语言是否内置Queue?如何基于Channel实现生产者-消费者模式?

Crystal语言的Queue与生产者-消费者模式实现

一、Crystal有没有内置的Queue?

当然有!Crystal标准库提供了线程安全的Queue(T),设计思路和Ruby的Queue高度相似,完全能满足你想要的线程间安全存取元素的需求。

你可以直接通过Queue(T).new创建指定类型的队列,常用方法和Ruby的Queue对应:

  • #push(item):往队列里加元素(对应Ruby的#enq)
  • #pop:从队列取元素,空队列时会阻塞等待(对应Ruby的#deq)
  • #size:查看当前队列元素数量
  • #empty?:判断队列是否为空

二、用内置Queue实现生产者-消费者模式

用内置Queue来实现你要的模式非常直观,生产者可以持续跑耗时任务,把结果推入队列就行,消费者则从队列里取出来处理:

require "queue"

queue = Queue(Int32).new

# 生产者线程
spawn do
  15.times do |i|
    # 模拟耗时操作
    sleep 0.1
    puts "生产者: 生成并推送 #{i}"
    queue.push(i)
  end
end

# 消费者线程
spawn do
  loop do
    item = queue.pop
    puts "消费者: 接收并处理 #{item}"
    sleep 0.5
  end
end

# 让主线程保持运行,避免程序直接退出
sleep 10

这里生产者的耗时操作不会被阻塞,因为Queue#push是非阻塞的(只要队列没到内存上限),推完元素就可以立刻继续下一次任务,完全符合你的需求。

三、用Channel实现非阻塞的生产者-消费者模式

你之前遇到Channel#send阻塞的问题,是因为默认创建的是无缓冲Channel——这种情况下send必须等消费者receive之后才能继续。要解决这个问题,你可以创建带缓冲的Channel,指定一个缓冲容量,这样当缓冲没满时,send会立刻返回,生产者就能继续执行后续任务:

# 创建一个容量为5的带缓冲Channel,最多存5个未被接收的元素
channel = Channel(Int32).new(5)

spawn do
  15.times do |i|
    # 模拟耗时操作
    sleep 0.1
    puts "生产者: 发送 #{i}"
    channel.send(i) # 缓冲未满时不会阻塞,直接放入缓冲
  end
  channel.close # 生产者完成所有任务后关闭Channel,告诉消费者不用等了
end

spawn do
  loop do
    # 用receive?判断Channel是否已关闭且为空
    if item = channel.receive?
      puts "消费者: 接收 #{item}"
      sleep 0.5
    else
      break # Channel关闭且没元素了,退出循环
    end
  end
end

sleep 10

带缓冲的Channel内部会维护一个队列,缓冲未满时生产者完全不用等。如果后面消费者处理速度跟不上,缓冲被填满的话,send会再次阻塞,这其实是个贴心的流量控制机制,防止生产者无限制生成数据导致内存爆掉。

四、两种方案怎么选?

  • 内置Queue:API和Ruby的Queue几乎一致,上手快,适合习惯Ruby写法的场景,线程安全还不用手动管关闭。
  • Channel:是Crystal并发模型的核心组件,带缓冲的Channel同样能实现需求,而且可以通过close明确终止消费者循环,适合更复杂的并发协作场景。

内容的提问来源于stack exchange,提问作者qpi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 07:23:01