diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/rabbit_amqqueue_process.erl | 2 | ||||
| -rw-r--r-- | src/rabbit_control.erl | 11 | ||||
| -rw-r--r-- | src/rabbit_mixed_queue.erl | 165 | ||||
| -rw-r--r-- | src/rabbit_queue_mode_manager.erl | 9 |
4 files changed, 106 insertions, 81 deletions
diff --git a/src/rabbit_amqqueue_process.erl b/src/rabbit_amqqueue_process.erl index ff0cc56b22..a701fa4d1d 100644 --- a/src/rabbit_amqqueue_process.erl +++ b/src/rabbit_amqqueue_process.erl @@ -237,7 +237,7 @@ deliver_queue(Fun, FunAcc0, end. deliver_from_queue(is_message_ready, undefined, #q { mixed_state = MS }) -> - 0 /= rabbit_mixed_queue:length(MS); + not rabbit_mixed_queue:is_empty(MS); deliver_from_queue(AckRequired, Acc = undefined, State = #q { mixed_state = MS }) -> {Res, MS2} = rabbit_mixed_queue:deliver(MS), MS3 = case {Res, AckRequired} of diff --git a/src/rabbit_control.erl b/src/rabbit_control.erl index 6649899ade..0ead9533c9 100644 --- a/src/rabbit_control.erl +++ b/src/rabbit_control.erl @@ -137,6 +137,9 @@ Available commands: list_bindings [-p <VHostPath>] list_connections [<ConnectionInfoItem> ...] + reduce_memory_footprint + increase_memory_footprint + Quiet output mode is selected with the \"-q\" flag. Informational messages are suppressed when quiet mode is in effect. @@ -276,6 +279,14 @@ action(list_connections, Node, Args, Inform) -> [ArgAtoms]), ArgAtoms); +action(reduce_memory_footprint, Node, _Args, Inform) -> + Inform("Reducing memory footprint", []), + call(Node, {rabbit_queue_mode_manager, reduce_memory_usage, []}); + +action(increase_memory_footprint, Node, _Args, Inform) -> + Inform("Reducing memory footprint", []), + call(Node, {rabbit_queue_mode_manager, increase_memory_usage, []}); + action(Command, Node, Args, Inform) -> {VHost, RemainingArgs} = parse_vhost_flag(Args), action(Command, Node, VHost, RemainingArgs, Inform). diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl index a950584a10..74e47a00db 100644 --- a/src/rabbit_mixed_queue.erl +++ b/src/rabbit_mixed_queue.erl @@ -45,14 +45,15 @@ msg_buf, next_write_seq, queue, - is_durable + is_durable, + length } ). start_link(Queue, IsDurable, disk) -> purge_non_persistent_messages( #mqstate { mode = disk, msg_buf = queue:new(), queue = Queue, - next_write_seq = 0, is_durable = IsDurable }); + next_write_seq = 0, is_durable = IsDurable, length = 0 }); start_link(Queue, IsDurable, mixed) -> {ok, State} = start_link(Queue, IsDurable, disk), to_mixed_mode(State). @@ -98,23 +99,21 @@ to_disk_only_mode(State = #mqstate { mode = mixed, queue = Q, msg_buf = MsgBuf, to_mixed_mode(State = #mqstate { mode = mixed }) -> {ok, State}; -to_mixed_mode(State = #mqstate { mode = disk, queue = Q }) -> +to_mixed_mode(State = #mqstate { mode = disk, queue = Q, length = Length }) -> rabbit_log:info("Converting queue to mixed mode: ~p~n", [Q]), %% load up a new queue with everything that's on disk. %% don't remove non-persistent messages that happen to be on disk QList = rabbit_disk_queue:dump_queue(Q), - {MsgBuf1, NextSeq1} = + {MsgBuf1, NextSeq1, Length} = lists:foldl( - fun ({MsgId, MsgBin, _Size, IsDelivered, _AckTag, SeqId}, {Buf, NSeq}) + fun ({MsgId, MsgBin, _Size, IsDelivered, _AckTag, SeqId}, + {Buf, NSeq, L}) when SeqId >= NSeq -> Msg = #basic_message { guid = MsgId } = bin_to_msg(MsgBin), - {queue:in({SeqId, Msg, IsDelivered, true}, Buf), SeqId + 1} - end, {queue:new(), 0}, QList), + {queue:in({SeqId, Msg, IsDelivered, true}, Buf), SeqId+1, L+1} + end, {queue:new(), 0, 0}, QList), State1 = State #mqstate { mode = mixed, msg_buf = MsgBuf1, - next_write_seq = NextSeq1 }, - rabbit_log:info("Queue length: ~p ~w ~w~n", - [Q, rabbit_mixed_queue:length(State), - rabbit_mixed_queue:length(State1)]), + next_write_seq = NextSeq1 }, {ok, State1}. purge_non_persistent_messages(State = #mqstate { mode = disk, queue = Q, @@ -131,7 +130,8 @@ purge_non_persistent_messages(State = #mqstate { mode = disk, queue = Q, ok = if Acks == [] -> ok; true -> rabbit_disk_queue:ack(Q, lists:reverse(Acks)) end, - {ok, State #mqstate { next_write_seq = NextSeq2 }}. + Length = NextSeq2 - NextSeq, + {ok, State #mqstate { next_write_seq = NextSeq2, length = Length }}. deliver_all_messages(Q, IsDurable, Acks, Requeue, NextSeq) -> case rabbit_disk_queue:deliver(Q) of @@ -158,12 +158,13 @@ bin_to_msg(MsgBin) -> binary_to_term(MsgBin). publish(Msg = #basic_message { guid = MsgId }, - State = #mqstate { mode = disk, queue = Q }) -> + State = #mqstate { mode = disk, queue = Q, length = Length }) -> ok = rabbit_disk_queue:publish(Q, MsgId, msg_to_bin(Msg), false), - {ok, State}; + {ok, State #mqstate { length = Length + 1 }}; publish(Msg = #basic_message { guid = MsgId, is_persistent = IsPersistent }, State = #mqstate { queue = Q, mode = mixed, is_durable = IsDurable, - next_write_seq = NextSeq, msg_buf = MsgBuf }) -> + next_write_seq = NextSeq, msg_buf = MsgBuf, + length = Length }) -> OnDisk = IsDurable andalso IsPersistent, ok = if OnDisk -> rabbit_disk_queue:publish_with_seq(Q, MsgId, NextSeq, @@ -172,7 +173,8 @@ publish(Msg = #basic_message { guid = MsgId, is_persistent = IsPersistent }, end, {ok, State #mqstate { next_write_seq = NextSeq + 1, msg_buf = queue:in({NextSeq, Msg, false, OnDisk}, - MsgBuf) + MsgBuf), + length = Length + 1 }}. %% Assumption here is that the queue is empty already (only called via @@ -180,67 +182,69 @@ publish(Msg = #basic_message { guid = MsgId, is_persistent = IsPersistent }, %% the disk queue could well not be the same as the NextSeq (true = %% NextSeq >= disk_queue_write_seq_for_queue(Q)) , but this doesn't %% matter because the AckTag will still be correct (AckTags for -%% non-persistent messages don't exist). (next_write_seq is actually -%% only used to calculate how many messages are in the queue). +%% non-persistent messages don't exist). publish_delivered(Msg = #basic_message { guid = MsgId, is_persistent = IsPersistent}, State = #mqstate { mode = Mode, is_durable = IsDurable, - next_write_seq = NextSeq, queue = Q }) + next_write_seq = NextSeq, queue = Q, + length = 0 }) when Mode =:= disk orelse (IsDurable andalso IsPersistent) -> rabbit_disk_queue:publish(Q, MsgId, msg_to_bin(Msg), false), + State1 = if Mode =:= disk -> State; + true -> State #mqstate { next_write_seq = NextSeq + 1 } + end, if IsDurable andalso IsPersistent -> %% must call phantom_deliver otherwise the msg remains at %% the head of the queue. This is synchronous, but %% unavoidable as we need the AckTag {MsgId, false, AckTag, 0} = rabbit_disk_queue:phantom_deliver(Q), - {ok, AckTag, State}; + {ok, AckTag, State1}; true -> %% in this case, we don't actually care about the ack, so %% auto ack it (asynchronously). ok = rabbit_disk_queue:auto_ack_next_message(Q), - {ok, noack, State #mqstate { next_write_seq = NextSeq + 1 }} + {ok, noack, State1} end; -publish_delivered(_Msg, State = #mqstate { mode = mixed, msg_buf = MsgBuf }) -> - true = queue:is_empty(MsgBuf), +publish_delivered(_Msg, State = #mqstate { mode = mixed, length = 0 }) -> {ok, noack, State}. -deliver(State = #mqstate { mode = disk, queue = Q, is_durable = IsDurable }) -> - case rabbit_disk_queue:deliver(Q) of - empty -> {empty, State}; - {MsgId, MsgBin, _Size, IsDelivered, AckTag, Remaining} -> - #basic_message { guid = MsgId, is_persistent = IsPersistent } = - Msg = bin_to_msg(MsgBin), - AckTag2 = if IsPersistent andalso IsDurable -> AckTag; - true -> ok = rabbit_disk_queue:ack(Q, [AckTag]), - noack - end, - {{Msg, IsDelivered, AckTag2, Remaining}, State} - end; +deliver(State = #mqstate { length = 0 }) -> + {empty, State}; +deliver(State = #mqstate { mode = disk, queue = Q, is_durable = IsDurable, + length = Length }) -> + {MsgId, MsgBin, _Size, IsDelivered, AckTag, Remaining} + = rabbit_disk_queue:deliver(Q), + #basic_message { guid = MsgId, is_persistent = IsPersistent } = + Msg = bin_to_msg(MsgBin), + AckTag2 = if IsPersistent andalso IsDurable -> AckTag; + true -> ok = rabbit_disk_queue:ack(Q, [AckTag]), + noack + end, + {{Msg, IsDelivered, AckTag2, Remaining}, + State #mqstate { length = Length - 1}}; deliver(State = #mqstate { mode = mixed, queue = Q, is_durable = IsDurable, - next_write_seq = NextWrite, msg_buf = MsgBuf }) -> - {Result, MsgBuf2} = queue:out(MsgBuf), - case Result of - empty -> - {empty, State}; - {value, {Seq, Msg = #basic_message { guid = MsgId, - is_persistent = IsPersistent }, - IsDelivered, OnDisk}} -> - AckTag = - if OnDisk -> - if IsPersistent andalso IsDurable -> - {MsgId, IsDelivered, AckTag2, _PersistRem} = - rabbit_disk_queue:phantom_deliver(Q), - AckTag2; - true -> - ok = rabbit_disk_queue:auto_ack_next_message(Q), - noack - end; - true -> noack - end, - {{Msg, IsDelivered, AckTag, (NextWrite - 1 - Seq)}, - State #mqstate { msg_buf = MsgBuf2 }} - end. + next_write_seq = NextWrite, msg_buf = MsgBuf, + length = Length }) -> + {{value, {Seq, Msg = #basic_message { guid = MsgId, + is_persistent = IsPersistent }, + IsDelivered, OnDisk}}, MsgBuf2} + = queue:out(MsgBuf), + AckTag = + if OnDisk -> + if IsPersistent andalso IsDurable -> + {MsgId, IsDelivered, AckTag2, _PersistRem} = + rabbit_disk_queue:phantom_deliver(Q), + AckTag2; + true -> + ok = rabbit_disk_queue:auto_ack_next_message(Q), + noack + end; + true -> noack + end, + Rem = Length - 1, + {{Msg, IsDelivered, AckTag, Rem}, + State #mqstate { msg_buf = MsgBuf2, length = Rem }}. remove_noacks(Acks) -> lists:filter(fun (A) -> A /= noack end, Acks). @@ -268,17 +272,19 @@ tx_publish(_Msg, State = #mqstate { mode = mixed }) -> only_msg_ids(Pubs) -> lists:map(fun (Msg) -> Msg #basic_message.guid end, Pubs). -tx_commit(Publishes, Acks, State = #mqstate { mode = disk, queue = Q }) -> +tx_commit(Publishes, Acks, State = #mqstate { mode = disk, queue = Q, + length = Length }) -> RealAcks = remove_noacks(Acks), ok = if ([] == Publishes) andalso ([] == RealAcks) -> ok; true -> rabbit_disk_queue:tx_commit(Q, only_msg_ids(Publishes), RealAcks) end, - {ok, State}; + {ok, State #mqstate { length = Length + erlang:length(Publishes) }}; tx_commit(Publishes, Acks, State = #mqstate { mode = mixed, queue = Q, msg_buf = MsgBuf, next_write_seq = NextSeq, - is_durable = IsDurable + is_durable = IsDurable, + length = Length }) -> {PersistentPubs, MsgBuf2, NextSeq2} = lists:foldl(fun (Msg = #basic_message { is_persistent = IsPersistent }, @@ -302,7 +308,8 @@ tx_commit(Publishes, Acks, State = #mqstate { mode = mixed, queue = Q, rabbit_disk_queue:tx_commit_with_seqs( Q, lists:reverse(PersistentPubs), RealAcks) end, - {ok, State #mqstate { msg_buf = MsgBuf2, next_write_seq = NextSeq2 }}. + {ok, State #mqstate { msg_buf = MsgBuf2, next_write_seq = NextSeq2, + length = Length + erlang:length(Publishes) }}. only_persistent_msg_ids(Pubs) -> lists:reverse( @@ -327,7 +334,8 @@ tx_cancel(Publishes, %% [{Msg, AckTag}] requeue(MessagesWithAckTags, State = #mqstate { mode = disk, queue = Q, - is_durable = IsDurable }) -> + is_durable = IsDurable, + length = Length }) -> %% here, we may have messages with no ack tags, because of the %% fact they are not persistent, but nevertheless we want to %% requeue them. This means publishing them delivered. @@ -346,11 +354,12 @@ requeue(MessagesWithAckTags, State = #mqstate { mode = disk, queue = Q, [] end, [], MessagesWithAckTags), ok = rabbit_disk_queue:requeue(Q, lists:reverse(Requeue)), - {ok, State}; + {ok, State #mqstate {length = Length + erlang:length(MessagesWithAckTags)}}; requeue(MessagesWithAckTags, State = #mqstate { mode = mixed, queue = Q, msg_buf = MsgBuf, next_write_seq = NextSeq, - is_durable = IsDurable + is_durable = IsDurable, + length = Length }) -> {PersistentPubs, MsgBuf2, NextSeq2} = lists:foldl( @@ -368,27 +377,25 @@ requeue(MessagesWithAckTags, State = #mqstate { mode = mixed, queue = Q, true -> rabbit_disk_queue:requeue_with_seqs( Q, lists:reverse(PersistentPubs)) end, - {ok, State #mqstate { msg_buf = MsgBuf2, next_write_seq = NextSeq2 }}. + {ok, State #mqstate {msg_buf = MsgBuf2, next_write_seq = NextSeq2, + length = Length + erlang:length(MessagesWithAckTags)}}. -purge(State = #mqstate { queue = Q, mode = disk }) -> +purge(State = #mqstate { queue = Q, mode = disk, length = Count }) -> Count = rabbit_disk_queue:purge(Q), - {Count, State}; -purge(State = #mqstate { queue = Q, msg_buf = MsgBuf, mode = mixed }) -> + {Count, State #mqstate { length = 0 }}; +purge(State = #mqstate { queue = Q, mode = mixed, length = Length }) -> rabbit_disk_queue:purge(Q), - Count = queue:len(MsgBuf), - {Count, State #mqstate { msg_buf = queue:new() }}. + {Length, State #mqstate { msg_buf = queue:new(), length = 0 }}. delete_queue(State = #mqstate { queue = Q, mode = disk }) -> rabbit_disk_queue:delete_queue(Q), - {ok, State}; + {ok, State #mqstate { length = 0 }}; delete_queue(State = #mqstate { queue = Q, mode = mixed }) -> rabbit_disk_queue:delete_queue(Q), - {ok, State #mqstate { msg_buf = queue:new() }}. + {ok, State #mqstate { msg_buf = queue:new(), length = 0 }}. -length(#mqstate { queue = Q, mode = disk }) -> - rabbit_disk_queue:length(Q); -length(#mqstate { mode = mixed, msg_buf = MsgBuf }) -> - queue:len(MsgBuf). +length(#mqstate { length = Length }) -> + Length. -is_empty(State) -> - 0 == rabbit_mixed_queue:length(State). +is_empty(#mqstate { length = Length }) -> + 0 == Length. diff --git a/src/rabbit_queue_mode_manager.erl b/src/rabbit_queue_mode_manager.erl index 080607bbad..32ad6b4cfe 100644 --- a/src/rabbit_queue_mode_manager.erl +++ b/src/rabbit_queue_mode_manager.erl @@ -38,7 +38,8 @@ -export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2, code_change/3]). --export([register/1, change_memory_usage/2]). +-export([register/1, change_memory_usage/2, + reduce_memory_usage/0, increase_memory_usage/0]). -define(SERVER, ?MODULE). @@ -54,6 +55,12 @@ register(Pid) -> change_memory_usage(_Pid, Conserve) -> gen_server2:cast(?SERVER, {change_memory_usage, Conserve}). + +reduce_memory_usage() -> + gen_server2:cast(?SERVER, {change_memory_usage, true}). + +increase_memory_usage() -> + gen_server2:cast(?SERVER, {change_memory_usage, false}). init([]) -> process_flag(trap_exit, true), |
