Erlang生产者-消费者程序异常行为原因排查求助
咱们来一步步捋清楚你的程序为啥会出现性能崩盘甚至搞崩系统的情况,核心是两个致命小bug加上不合理的请求发送逻辑,直接把Erlang的进程和消息队列玩坏了:
1. 消费者生成函数的低级错误:错调了生产者生成函数
先看你的spawnConsumers代码:
spawnConsumers(Number, ServerPid) -> case Number of 0 -> io:format("Spawned producers"); % 这里输出都写错了,应该是consumers N -> spawn(zad2,consumer,[ServerPid]), spawnProducers(N - 1,ServerPid) end.
递归的时候你调用的是spawnProducers而不是spawnConsumers!这意味着:
- 当你调用
start(10,10)时,原本要生成10个消费者,结果每生成1个消费者,就会递归生成N-1个生产者 - 最终生产者数量是
10(初始调用的spawnProducers) + 10+9+8+...+1 = 65个,消费者只有10个 - 到
start(100,100)时,生产者数量直接变成100 + 100+99+...+1 = 5150个,加上100个消费者,总共5250个进程!这么多进程疯狂发消息,系统CPU和内存直接被榨干,不重启才怪。
修复方法:把递归里的spawnProducers(N - 1,ServerPid)改成spawnConsumers(N - 1,ServerPid),顺便把输出的Spawned producers改成Spawned consumers,别再误导自己了。
2. 生产者/消费者无节制发送请求,导致服务器消息队列积压
你的producer和consumer是无限循环,发完一个请求立刻递归发下一个,根本不管服务器有没有回复:
producer(ServerPid) -> X = rand:uniform(9), ToProduce = [rand:uniform(500) || _ <- lists:seq(1, X)], ServerPid ! {self(),produce,ToProduce}, producer(ServerPid). % 发完立刻发下一个,完全不等服务器处理
这就导致每个生产者/消费者在极短时间内向服务器发海量消息,服务器的消息队列会迅速膨胀到几万甚至几十万条。Erlang进程的消息队列是存在内存里的,消息越多占内存越大,而且服务器要逐个处理这些积压的消息,CPU全用来干这个了,后续请求自然处理越来越慢,最后直接把内存撑爆。
修复方法:让生产者/消费者等服务器回复后再发下一个请求,比如:
producer(ServerPid) -> X = rand:uniform(9), ToProduce = [rand:uniform(500) || _ <- lists:seq(1, X)], ServerPid ! {self(),produce,ToProduce}, receive % 等服务器回复完再继续 ok -> ok; tryagain -> ok end, producer(ServerPid). consumer(ServerPid) -> X = rand:uniform(9), ServerPid ! {self(),consume,X}, receive {ok, _Data} -> ok; tryagain -> ok end, consumer(ServerPid).
这样每个进程会等上一个请求处理完再发下一个,避免消息队列爆炸。
3. 顺便提个逻辑错误(不影响性能,但影响结果正确性)
服务器处理消费请求时用了lists:subtract(Buffer, Data),这个函数会删掉Buffer里所有和Data重复的元素,而不是删掉前Number个元素。比如Buffer是[1,2,1,3],消费2个元素得到Data是[1,2],lists:subtract会返回[3],但正确的剩余Buffer应该是[1,3]。
修复方法:用lists:nthtail(Number, Buffer)替代lists:subtract(Buffer, Data),这样就能正确保留消费后的剩余元素。
额外的性能优化建议
- 可以给生产者/消费者加个微小延迟,比如
timer:sleep(10),模拟真实场景的生产/消费间隔,避免短时间内请求太密集 - 把缓冲区从列表改成Erlang的
queue模块(队列),因为列表的++操作是O(n)复杂度,队列的入队出队都是O(1),缓冲区大的时候性能提升很明显:
% 初始化服务器时用queue:new()替代[] server(Buffer, Capacity, CountPid) -> receive {Pid, produce, InputList} -> NumberProduce = length(InputList), case queue:len(Buffer) + NumberProduce =< Capacity of true -> % 把InputList的元素逐个入队 NewBuffer = lists:foldl(fun(Item, Q) -> queue:in(Item, Q) end, Buffer, InputList), CountPid ! NumberProduce, Pid ! ok, server(NewBuffer,Capacity, CountPid); false -> Pid ! tryagain, server(Buffer, Capacity, CountPid) end; {Pid, consume, Number} -> case queue:len(Buffer) >= Number of true -> % 取出前Number个元素 {Data, NewBuffer} = lists:foldl(fun(_, {Acc, Q}) -> {[queue:get(Q)|Acc], queue:drop(Q)} end, {[], Buffer}, lists:seq(1, Number)), Pid ! {ok, lists:reverse(Data)}, server(NewBuffer, Capacity,CountPid); false -> Pid ! tryagain, server(Buffer, Capacity, CountPid) end end.
最后总结
修复完消费者生成的bug,再让生产者/消费者等待服务器回复,你的程序就能达到预期的线性增长效果,进程越多生产速度也会越快,再也不会出现系统资源耗尽的情况了。
内容的提问来源于stack exchange,提问作者Dawid

