summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/rabbit_amqqueue_process.erl2
-rw-r--r--src/rabbit_control.erl11
-rw-r--r--src/rabbit_mixed_queue.erl165
-rw-r--r--src/rabbit_queue_mode_manager.erl9
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),