【问题标题】:Understanding why process stops receiving messages了解进程停止接收消息的原因
【发布时间】: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
  • 我现在明白你的意思了,现在它击中了我。我根据我的工作人员放置了出队逻辑,但工作人员可以是空闲的,并且队列可以继续收集元素而不会触发任何事情。跨度>
  • 我现在明白你的意思了,现在它击中了我。我根据我的工作人员放置了出队逻辑,但工作人员可能是空闲的,队列可以继续收集元素而不会触发任何事情。我猜我需要的是一个阻塞队列,它阻塞进程直到它可以出队。这个进程能够独立运行。

标签: erlang timeout trace


【解决方案1】:

monitor 行为需要依赖于frun。它只需要取决于worker 是否为free。我已更新 monitor 函数以在以下代码中反映这一点。

-module(mq).
-compile(export_all).

-record(monstate,{
    queue,
    qc,
    wpid,
    free=true,
    wref,
    init=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()});

monitor(MState=#monstate{wpid=W,free=F,wref=Ref,queue=Q,qc=C})->
  receive
    {request,{From ,Message}} ->
      %% check whether worker is free or not
      case F of
        true ->
          W ! {From,Message},
          monitor(MState#monstate{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;

    {'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.

用法

Eshell V10.5  (abort with ^G)
1> c(mq).
mq.erl:2: Warning: export_all flag enabled - all functions will be exported
{ok,mq}
2> A=mq:start().
<0.92.0>
3> A ! {self(),aa}.
{<0.85.0>,aa}
4> flush().
Shell got {processed,aa}
ok
5> A ! {self(),aa}.
{<0.85.0>,aa}
6> flush().
Shell got {processed,aa}
ok
7> A ! {self(), aa}, A ! {self(), bb}.
{<0.85.0>,bb}
8> flush().                           
Shell got {processed,aa}
Shell got {processed,bb}
ok
9> 

【讨论】:

  • 非常感谢您的回复!我真的不知道如何在队列上创建一个循环以继续出队。您的响应是正确的。如果队列中有您在{worker,{finished,_}} 模式上循环的项目。
【解决方案2】:

在您的代码中,似乎frun 在第一次运行后始终为假:

            case R of
                true -> self() ! {worker,{finished,R}},
                        monitor(NewState#monstate{frun=false});
                false -> monitor(NewState#monstate{frun=false})
            end;

一旦到达false,将不会传递{worker, {finished, R}} 消息,因此不会从队列中提取任何元素。

更新:死锁序列:

  1. Monitor 接收第一份工作并将其转发给工作人员
  2. Monitor 的frun 现在是假的
  3. 工人执行工作
  4. 工作人员通知监视器作业已完成
  5. 由于监视器的队列是空的,所以没有任何反应。
  6. Monitor 收到第二份作业,因为 frun 为假,作业不会转发给工作人员

【讨论】:

  • 这就是我们想要的行为。第一次运行时,我希望我的monitor 向自己发出信号,表明队列中有一个项目,以便将其发送到worker。从那时起,worker 将充当向@ 发出信号的那个987654331@.frun(第一次运行)用作启动处理循环的标志。
  • worker 接收第二个作业的唯一方法是在监视器中接收worker,{finished,_}} 消息,这仅在监视器具有frun = true 或worker 完成并且有排队的东西时才会发生。如果工作人员在将某些内容添加到队列之前完成,它就会死锁。 (可以A ! {self(), aa}, A ! {self(), bb}.查看)
  • 这在 Erlang 中有些常见,如果不是练习,我建议检查一下是否有任何可用的库适合您的项目。
  • 主管已经是 gen_servers,请小心。过去,我将队列中的第一个位置视为“进行中”位置。因此,当队列从 0 个元素变为 1 个元素时,或者当工作人员完成一项工作并且有排队的工作时,监视器会向工作人员发送一个工作。监视器在完成之前不会从队列中删除第一个作业
  • 虽然它比你想要做的更复杂,但你可以查看this 的想法
猜你喜欢
  • 1970-01-01
  • 2016-05-24
  • 1970-01-01
  • 2015-03-10
  • 1970-01-01
  • 2018-07-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多