Erlang OTP中GenServer handle_call内调用函数触发超时错误
Erlang OTP聊天服务同步消息超时问题解决
问题描述
基于Erlang OTP实现聊天架构:chat_server作为中心服务运行在Alice节点,Bob、Clark等客户端节点通过chat_client:start_link/0启动。chat_client提供send_message/3和receive_message/3两个GenServer同步调用接口,分别对应chat_server的send_message_server/3和receive_message_server/3。单独调用两个接口均正常,但在chat_client的send_message的handle_call回调中调用receive_message/3时,会出现调用超时错误,接收节点(Clark)还会出现进程冻结现象。需要实现发送消息时接收节点同步接收消息的功能。
代码片段
chat_server.erl
-module(chat_server). -behaviour(gen_server). -export([start_link/0, send_message_server/3, receive_message_server/3, stop/0]). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -define(SERVER, ?MODULE). -record(chat_server_state, {messages,receivers,senders,sent,received}). start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). send_message_server(From,To,Msg)-> io:format("message sent!~n"), gen_server:call({?MODULE,'alice@DESKTOP-RD414DV'}, {send_message_server,From,To,Msg}). receive_message_server(From,To,Msg)-> io:format("message received!~n"), gen_server:call({?MODULE,'alice@DESKTOP-RD414DV'}, {receive_message_server,From,To,Msg}). stop() -> gen_server:stop(?MODULE). init([]) -> io:format("~p starting...",[?MODULE]), {ok, #chat_server_state{ messages = [], sent=[], received = [], receivers = [], senders=[] }}. handle_call({send_message_server,From,To,Msg}, _From, State = #chat_server_state{receivers = Receivers,messages = Messages, sent =Sent, senders = Senders}) -> io:format("message came to server and then conveyed by the server~n"), {reply, ok, State#chat_server_state{receivers = [To|Receivers], messages = [Msg|Messages], sent=[Msg|Sent], senders = [From|Senders]}}; handle_call({receive_message_server,From,To,Msg}, _From, State = #chat_server_state{received = Received }) -> io:format("A message sent by ~p, it received to ~p~n",[From,To]), {reply, ok, State#chat_server_state{received = [Msg|Received]}}. handle_cast(_Request, State = #chat_server_state{}) -> {noreply, State}. handle_info(_Info, State = #chat_server_state{}) -> {noreply, State}. terminate(_Reason, _State = #chat_server_state{}) -> ok. code_change(_OldVsn, State = #chat_server_state{}, _Extra) -> {ok, State}.
chat_client.erl
-module(chat_client). -behaviour(gen_server). -export([start_link/0, stop/0, send_message/3, receive_message/3]). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -define(SERVER, ?MODULE). -record(chat_server_state, {messages,receivers,senders,sent,received}). start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). send_message(From,To, Msg)-> gen_server:call({?MODULE, node()},{send_message,From,To,Msg}). receive_message(From, To, Msg)-> gen_server:call({?MODULE, list_to_atom(To)},{receive_message,From,To,Msg}). stop()-> gen_server:stop(?MODULE). init([]) -> io:format("~p connected...",[node()]), {ok, #chat_server_state{ messages = [], sent=[], received = [], receivers = [], senders = [] }}. handle_call({send_message,From,To,Msg}, _From, State = #chat_server_state{receivers = Receivers,messages = Messages, sent =Sent, senders = Senders}) -> chat_server:send_message_server(From,To,Msg), %% receive_message(From,To,Msg), //添加此行后出现超时 {reply, ok, State#chat_server_state{receivers = [To|Receivers], messages = [Msg|Messages], sent=[Msg|Sent], senders = [From|Senders]}}; handle_call({receive_message,From,To,Msg}, _From, State = #chat_server_state{messages = Messages, received = Received}) -> chat_server:receive_message_server(From,To,Msg), io:format("i am ~p~n and sent this by ~p",[To, From]), {reply, ok, State#chat_server_state{ messages = [Msg|Messages], sent=[Msg|Received] }}. handle_cast(_Request, State = #chat_server_state{}) -> {noreply, State}. handle_info(_Info, State = #chat_server_state{}) -> {noreply, State}. terminate(_Reason, _State = #chat_server_state{}) -> ok. code_change(_OldVsn, State = #chat_server_state{}, _Extra) -> {ok, State}.
现象复现
单独调用正常
- 发送消息(Bob节点):
(bob@DESKTOP-RD414DV)44> chat_client:send_message("bob@DESKTOP-RD414DV","clark@DESKTOP-RD414DV","heee"). message sent! ok
- 接收消息(Clark节点):
(clark@DESKTOP-RD414DV)38> chat_client:receive_message("bob@DESKTOP-RD414DV","clark@DESKTOP-RD414DV","heee"). message received! i am "clark@DESKTOP-RD414DV" and sent this by "bob@DESKTOP-RD414DV" ok
嵌套调用出现问题
在chat_client的send_message回调中添加receive_message调用后,执行发送消息超时:
(bob@DESKTOP-RD414DV)38> chat_client:send_message("bob@DESKTOP-RD414DV","clark@DESKTOP-RD414DV","heee"). message sent! ** exception exit: {timeout, {gen_server,call, [{chat_client,'bob@DESKTOP-RD414DV'}, {send_message,"bob@DESKTOP-RD414DV", "clark@DESKTOP-RD414DV","heee"}]}} in function gen_server:call/2 (gen_server.erl, line 367)
接收节点Clark后续操作冻结,执行代码编译时才输出消息:
(clark@DESKTOP-RD414DV)5> c(chat_client). message received! i am "clark@DESKTOP-RD414DV" and sent this by "bob@DESKTOP-RD414DV" {ok,chat_client}
问题原因
核心问题是GenServer同步调用嵌套导致的阻塞死锁:
- Bob节点的
chat_client进程在处理send_message同步调用时,发起对Clark节点chat_client的同步调用,导致自身进程阻塞,等待Clark的回复。 - Clark节点的
chat_client进程在处理该同步调用时,又发起对Alice节点chat_server的同步调用,若chat_server此时有延迟或Clark进程因其他原因无法及时回复,Bob的进程会因等待超时抛出错误。 - GenServer是单进程模型,阻塞的进程无法处理其他请求,导致接收节点出现冻结现象。
解决方案
将接收消息的同步调用改为异步通知,由中心服务chat_server主动推送消息给接收方,避免发送方进程阻塞等待接收方回复。
修改后的代码
chat_server.erl(新增异步推送逻辑)
-module(chat_server). -behaviour(gen_server). -export([start_link/0, send_message_server/3, receive_message_server/3, stop/0]). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -define(SERVER, ?MODULE). -record(chat_server_state, {messages,receivers,senders,sent,received}). start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). send_message_server(From,To,Msg)-> io:format("message sent!~n"), gen_server:call({?MODULE,'alice@DESKTOP-RD414DV'}, {send_message_server,From,To,Msg}). receive_message_server(From,To,Msg)-> io:format("message received!~n"), gen_server:call({?MODULE,'alice@DESKTOP-RD414DV'}, {receive_message_server,From,To,Msg}). stop() -> gen_server:stop(?MODULE). % 新增:异步通知接收方客户端 send_to_receiver(From, To, Msg) -> gen_server:cast({chat_client, list_to_atom(To)}, {receive_message, From, To, Msg}). init([]) -> io:format("~p starting...",[?MODULE]), {ok, #chat_server_state{ messages = [], sent=[], received = [], receivers = [], senders=[] }}. handle_call({send_message_server,From,To,Msg}, _From, State = #chat_server_state{receivers = Receivers,messages = Messages, sent =Sent, senders = Senders}) -> io:format("message came to server and then conveyed by the server~n"), % 调用异步推送通知接收方 send_to_receiver(From, To, Msg), {reply, ok, State#chat_server_state{receivers = [To|Receivers], messages = [Msg|Messages], sent=[Msg|Sent], senders = [From|Senders]}}; handle_call({receive_message_server,From,To,Msg}, _From, State = #chat_server_state{received = Received }) -> io:format("A message sent by ~p, it received to ~p~n",[From,To]), {reply, ok, State#chat_server_state{received = [Msg|Received]}}. handle_cast(_Request, State = #chat_server_state{}) -> {noreply, State}. handle_info(_Info, State = #chat_server_state{}) -> {noreply, State}. terminate(_Reason, _State = #chat_server_state{}) -> ok. code_change(_OldVsn, State = #chat_server_state{}, _Extra) -> {ok, State}.
chat_client.erl(修改为异步接收消息)
-module(chat_client). -behaviour(gen_server). -export([start_link/0, stop/0, send_message/3, receive_message/3]). -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). -define(SERVER, ?MODULE). % 修正:使用自己的状态记录,避免与chat_server混淆 -record(chat_client_state, {messages,receivers,senders,sent,received}). start_link() -> gen_server:start_link({local, ?MODULE}, ?MODULE, [], []). send_message(From,To, Msg)-> gen_server:call({?MODULE, node()},{send_message,From,To,Msg}). % 修改:将同步调用改为异步cast receive_message(From, To, Msg)-> gen_server:cast({?MODULE, list_to_atom(To)},{receive_message,From,To,Msg}). stop()-> gen_server:stop(?MODULE). init([]) -> io:format("~p connected...",[node()]), {ok, #chat_client_state{ messages = [], sent=[], received = [], receivers = [], senders = [] }}. handle_call({send_message,From,To,Msg}, _From, State = #chat_client_state{receivers = Receivers,messages = Messages, sent =Sent, senders = Senders}) -> chat_server:send_message_server(From,To,Msg), % 移除直接调用receive_message的代码,由chat_server异步推送 {reply, ok, State#chat_client_state{receivers = [To|Receivers], messages = [Msg|Messages], sent=[Msg|Sent], senders = [From|Senders]}}; % 新增:处理异步推送的消息 handle_cast({receive_message,From,To,Msg}, State = #chat_client_state{messages = Messages, received = Received}) -> chat_server:receive_message_server(From,To,Msg), io:format("i am ~p~n and received this from ~p~n",[To, From]), {noreply, State#chat_client_state{ messages = [Msg|Messages], received=[Msg|Received] }}; handle_cast(_Request, State = #chat_client_state{}) -> {noreply, State}. handle_info(_Info, State = #chat_client_state{}) -> {noreply, State}. terminate(_Reason, _State = #chat_client_state{}) -> ok. code_change(_OldVsn, State = #chat_client_state{}, _Extra) -> {ok, State}.
说明
- 中心服务
chat_server在处理发送请求后,通过gen_server:cast异步通知接收方客户端,避免发送方进程阻塞。 - 接收方客户端通过
handle_cast处理异步消息,不会阻塞自身进程,也不会导致发送方超时。 - 修正了
chat_client的状态记录命名,避免与chat_server的状态混淆,提高代码可读性。
内容的提问来源于stack exchange,提问作者Ke_Sandaru
相关产品推荐
相关产品推荐

