diff options
| author | Matthew Sackman <matthew@lshift.net> | 2009-06-12 00:54:48 +0100 |
|---|---|---|
| committer | Matthew Sackman <matthew@lshift.net> | 2009-06-12 00:54:48 +0100 |
| commit | be5837815fad0f5073031c4c5c377836267cd702 (patch) | |
| tree | 28593824cf4bc97c65c14f1e97b871d229e1e4ab /src | |
| parent | 45bc68c0f95116be2a54c7cab546e7cee607db0b (diff) | |
| download | rabbitmq-server-git-be5837815fad0f5073031c4c5c377836267cd702.tar.gz | |
Made mixed_queue track its length by itself. This avoids synchronous calls to the disk_queue when operating in disk only mode and seems to have substantially improved performance (in addition to avoiding a sync call, repeated lasting for the length of a queue (erlang stdlib) with a million+ items in it can't have been cheap). It now seems to be very much the case that when coming out of disk only mode, huge back logs are recovered reliably.
Also, added reduce_memory_footprint and increase_memory_footprint to control. Both can be run twice and alter whether the disk_queue changes mode or the individual queues.
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), |
