【发布时间】:2020-07-03 20:43:50
【问题描述】:
我有一个情况,我有 3 个进程:
- 一个进程充当消息的调度程序:
server - 一个进程作为主管(对于工人):
monitor - 一个进程作为工作人员,在完成时通知主管:
worker
server 发送一个监控请求,monitor 首先检查worker 是否忙。如果忙,monitor 将消息排入队列(如果容量未达到),否则将其转发给worker .
当工作人员完成处理时,它会通知客户端和monitor
我的问题是我的工作进程在处理完第一条消息后停止响应
-module(mq).
-compile(export_all).
-record(monstate,{
queue,
qc,
wpid,
free=true,
wref,
init=false,
frun=false
}).
-record(sstate,{
init=false,
mpid=null,
mref=null
}).
-define(QUEUE_SIZE,5).
-define(PROC_SLEEP,2000).
createProcess({M,F,A})->
Pid=spawn(M,F,[A]),
Ref=erlang:monitor(process,Pid),
{Pid,Ref}.
start()->
spawn(?MODULE,server,[#sstate{init=false}]).
server(State=#sstate{init=I})when I=:=false ->
{MPid,MRef}=createProcess({?MODULE,monitor,#monstate{init=false}}),
server(State#sstate{init=true,mpid=MPid,mref=MRef});
server(State=#sstate{mpid=MPid,mref=MRef})->
receive
{From,state}->From ! State,
server(State);
{From,Message}-> MPid ! {request,{From,Message}},
server(State);
{'DOWN',MRef,process,MPid,_}-> {NewMPid,NewMRef}=createProcess({?MODULE,monitor,#monstate{init=false}}),
server(State#sstate{mpid=NewMPid,mref=NewMRef});
_ ->exit(invalid_message)
end.
tryEnqueue(Message,MState=#monstate{queue=Q,qc=C}) when C<?QUEUE_SIZE->
NewQueue=queue:in(Message,Q),
{queued,MState#monstate{qc=C+1,queue=NewQueue}};
tryEnqueue(_,MState)->{queue_full,MState}.
monitor(MState=#monstate{wpid=_,wref=_,init=I}) when I=:= false ->
{WorkerPid,WorkerRef}=createProcess({?MODULE,worker,self()}),
monitor(MState#monstate{wpid=WorkerPid,wref=WorkerRef,init=true,qc=0,queue=queue:new(),frun=true});
monitor(MState=#monstate{wpid=W,free=F,wref=Ref,queue=Q,qc=C,frun=R})->
receive
{request,{From ,Message}} ->
{Result,NewState}=tryEnqueue({From,Message},MState),
case Result of
queue_full -> From ! {queue_full,Message};
_ -> ok
end,
case R of
true -> self() ! {worker,{finished,R}},
monitor(NewState#monstate{frun=false});
false -> monitor(NewState#monstate{frun=false})
end;
{worker,{finished,_}}-> case queue:out(Q) of
{{_,Element},Rest} -> W ! Element,
monitor(MState#monstate{free=false,queue=Rest,qc=C-1});
{empty,Rest} -> monitor(MState#monstate{free=true,queue=Rest})
end;
{'DOWN',Ref,process,_,_}->
{NewWorkerPid,NewWorkerRef}=createProcess({?MODULE,worker,self()}),
monitor(MState#monstate{wpid=NewWorkerPid,wref=NewWorkerRef,free=true});
_->exit(invalid_message)
end.
worker(MPid)->
receive
{From,MSG} ->
timer:sleep(?PROC_SLEEP),
From ! {processed,MSG},
MPid ! {worker,{finished,MSG}},
worker(MPid);
_ ->exit(bad_msg)
end.
用法
2> A=mq:start().
<0.83.0>
3> A ! {self(),aa}.
{<0.76.0>,aa}
4> flush().
Shell got {processed,aa}
ok
5> A ! {self(),aa}.
{<0.76.0>,aa}
6> flush().
ok
我添加了一个跟踪器来查看发生了什么:
10> dbg:tracer().
{ok,<0.96.0>}
11> dbg:p(new,[sos,m]).
{ok,[{matched,nonode@nohost,0}]}
首次运行:
14> A ! {self(),aa}.
(<0.100.0>) << {<0.76.0>,aa} // message received my server
(<0.100.0>) <0.101.0> ! {request,{<0.76.0>,aa}} //message forwarded by server to monitor
{<0.76.0>,aa}
(<0.101.0>) << {request,{<0.76.0>,aa}}
15> (<0.101.0>) <0.101.0> ! {worker,{finished,true}} //monitor starting the cycle
15> (<0.101.0>) << {worker,{finished,true}}
15> (<0.101.0>) <0.102.0> ! {<0.76.0>,aa} // monitor sending message to worker
15> (<0.102.0>) << {<0.76.0>,aa}
15> (<0.105.0>) <0.62.0> ! {io_request,<0.105.0>,
#Ref<0.3226054513.2760638467.167990>,
{get_until,unicode,
["15",62,32],
erl_scan,tokens,
[1,[text]]}}
15> (<0.102.0>) << timeout //worker getting timeout ??
15> (<0.102.0>) <0.76.0> ! {processed,aa} //worker sends to self() thje message
15> (<0.102.0>) <0.101.0> ! {worker,{finished,aa}} //worker notifies monitor to update state
15> (<0.101.0>) << {worker,{finished,aa}}
第二次运行:
15> A ! {self(),aa}.
(<0.100.0>) << {<0.76.0>,aa}
(<0.100.0>) <0.101.0> ! {request,{<0.76.0>,aa}} //monitor receiveing message
{<0.76.0>,aa}
(<0.101.0>) << {request,{<0.76.0>,aa}}
16> (<0.106.0>) <0.62.0> ! {io_request,<0.106.0>,
#Ref<0.3226054513.2760638467.168007>,
{get_until,unicode,
["16",62,32],
erl_scan,tokens,
[1,[text]]}}
从我的跟踪中可以看出,在第一次通话中我不明白会发生什么。我的worker 是否超时?如果是,为什么?
P.S frun 变量用作标志,仅在第一次 monitor 迭代时为真,这样当第一个项目到达时,进程将调用自己来处理它(发送它给工人),因为工人是免税的。
在第一次运行后,monitor 将在worker 发出信号他空闲时从队列中取出项目。
更新
所以在有用的 cmets 之后,我在 monitor 中稍微改变了我的逻辑,以便 worker 在第一次运行时收到消息,或者在他完成并通知 monitor 之后,仍然有monitor 的队列中的项目。
我还是不行。死锁在哪里?
monitor(MState=#monstate{wpid=W,free=F,wref=Ref,queue=Q,qc=C,frun=FirstRun})->
receive
{request,{From ,Message}} -> case FirstRun of
true -> W ! {From,Message},
monitor(MState#monstate{frun=false,free=false});
false ->
St=case tryEnqueue({From,Message},MState) of
{queue_full,S} -> From ! {queue_full,Message},
S;
{queued,S} -> S
end,
monitor(St)
end;
{worker,{finished,_}}-> case queue:out(Q) of
{{_,Element},Rest} -> W ! Element,
monitor(MState#monstate{free=false,queue=Rest,qc=C-1});
{empty,Rest} -> monitor(MState#monstate{free=true,queue=Rest})
end;
end.
【问题讨论】:
-
关于更新,死锁也是一样的。我已经用死锁序列更新了响应。此外,
W!Message缺少From -
我现在明白你的意思了,现在它击中了我。我根据我的工作人员放置了出队逻辑,但工作人员可以是空闲的,并且队列可以继续收集元素而不会触发任何事情。跨度>
-
我现在明白你的意思了,现在它击中了我。我根据我的工作人员放置了出队逻辑,但工作人员可能是空闲的,队列可以继续收集元素而不会触发任何事情。我猜我需要的是一个阻塞队列,它阻塞进程直到它可以出队。这个进程能够独立运行。