summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorMatthew Sackman <matthew@lshift.net>2009-06-10 21:02:54 +0100
committerMatthew Sackman <matthew@lshift.net>2009-06-10 21:02:54 +0100
commita99b541be0332defa32c8ce7548acd9830c56892 (patch)
treec2b2601b141f3296ec8465d0a278132e2c84e015
parentcf5e1ae896b8aceca66c776b3b98806eb3bfefcf (diff)
downloadrabbitmq-server-git-a99b541be0332defa32c8ce7548acd9830c56892.tar.gz
well, I've made the acking for messages which are on disk but are not persistent/durable async, and it has improved some issues. But, if you switch to disk only mode, then allow, say 10k messages to build up (use MulticastMain) then switch back to ram mode, then it won't recover - the receive rate will stay very low, and rabbitmqctl list_queues will continue to grow insanely. This is very very odd, because querying the disk_queue directly for the queue length shows it drops to 0, but at least one CPU is maxed out at 100% use, messages continue to arrive, but the delivery rate never goes back up. Mysterious.
-rw-r--r--src/rabbit_disk_queue.erl25
-rw-r--r--src/rabbit_mixed_queue.erl21
2 files changed, 32 insertions, 14 deletions
diff --git a/src/rabbit_disk_queue.erl b/src/rabbit_disk_queue.erl
index 3fc208df7e..2b6f7b003b 100644
--- a/src/rabbit_disk_queue.erl
+++ b/src/rabbit_disk_queue.erl
@@ -41,7 +41,7 @@
-export([publish/4, publish_with_seq/5, deliver/1, phantom_deliver/1, ack/2,
tx_publish/2, tx_commit/3, tx_commit_with_seqs/3, tx_cancel/1,
requeue/2, requeue_with_seqs/2, purge/1, delete_queue/1,
- dump_queue/1, delete_non_durable_queues/1
+ dump_queue/1, delete_non_durable_queues/1, auto_ack_next_message/1
]).
-export([length/1, is_empty/1, next_write_seq/1]).
@@ -287,6 +287,9 @@ phantom_deliver(Q) ->
ack(Q, MsgSeqIds) when is_list(MsgSeqIds) ->
gen_server2:cast(?SERVER, {ack, Q, MsgSeqIds}).
+auto_ack_next_message(Q) ->
+ gen_server2:cast(?SERVER, {auto_ack_next_message, Q}).
+
tx_publish(MsgId, Msg) when is_binary(Msg) ->
gen_server2:cast(?SERVER, {tx_publish, MsgId, Msg}).
@@ -315,7 +318,7 @@ delete_queue(Q) ->
gen_server2:cast(?SERVER, {delete_queue, Q}).
dump_queue(Q) ->
- gen_server2:pcall(?SERVER, {dump_queue, Q}, infinity).
+ gen_server2:call(?SERVER, {dump_queue, Q}, infinity).
delete_non_durable_queues(DurableQueues) ->
gen_server2:call(?SERVER, {delete_non_durable_queues, DurableQueues}, infinity).
@@ -422,10 +425,10 @@ handle_call({publish_with_seq, Q, MsgId, SeqId, MsgBody}, _From, State) ->
internal_publish(Q, MsgId, SeqId, MsgBody, true, State),
{reply, MsgSeqId, State1};
handle_call({deliver, Q}, _From, State) ->
- {ok, Result, State1} = internal_deliver(Q, true, State),
+ {ok, Result, State1} = internal_deliver(Q, true, false, State),
{reply, Result, State1};
handle_call({phantom_deliver, Q}, _From, State) ->
- {ok, Result, State1} = internal_deliver(Q, false, State),
+ {ok, Result, State1} = internal_deliver(Q, false, false, State),
{reply, Result, State1};
handle_call({tx_commit, Q, PubMsgIds, AckSeqIds}, _From, State) ->
PubMsgSeqIds = zip_with_tail(PubMsgIds, {duplicate, next}),
@@ -499,6 +502,9 @@ handle_cast({publish_with_seq, Q, MsgId, SeqId, MsgBody}, State) ->
handle_cast({ack, Q, MsgSeqIds}, State) ->
{ok, State1} = internal_ack(Q, MsgSeqIds, State),
{noreply, State1};
+handle_cast({auto_ack_next_message, Q}, State) ->
+ {ok, State1} = internal_auto_ack(Q, State),
+ {noreply, State1};
handle_cast({tx_publish, MsgId, MsgBody}, State) ->
{ok, State1} = internal_tx_publish(MsgId, MsgBody, State),
{noreply, State1};
@@ -696,14 +702,14 @@ sequence_lookup(Sequences, Q) ->
%% ---- INTERNAL RAW FUNCTIONS ----
-internal_deliver(Q, ReadMsg, State = #dqstate { sequences = Sequences }) ->
+internal_deliver(Q, ReadMsg, FakeDeliver, State = #dqstate { sequences = Sequences }) ->
case ets:lookup(Sequences, Q) of
[] -> {ok, empty, State};
[{Q, SeqId, SeqId, 0}] -> {ok, empty, State};
[{Q, ReadSeqId, WriteSeqId, Length}] when Length > 0 ->
Remaining = Length - 1,
{ok, Result, NextReadSeqId, State1} =
- internal_read_message(Q, ReadSeqId, false, ReadMsg, State),
+ internal_read_message(Q, ReadSeqId, FakeDeliver, ReadMsg, State),
true = ets:insert(Sequences,
{Q, NextReadSeqId, WriteSeqId, Remaining}),
{ok,
@@ -739,6 +745,13 @@ internal_read_message(Q, ReadSeqId, FakeDeliver, ReadMsg, State) ->
{ok, {MsgId, Delivered, {MsgId, ReadSeqId}}, NextReadSeqId, State}
end.
+internal_auto_ack(Q, State) ->
+ case internal_deliver(Q, false, true, State) of
+ {ok, empty, State1} -> {ok, State1};
+ {ok, {_MsgId, _Delivered, MsgSeqId, _Remaining}, State1} ->
+ remove_messages(Q, [MsgSeqId], true, State1)
+ end.
+
internal_ack(Q, MsgSeqIds, State) ->
remove_messages(Q, MsgSeqIds, true, State).
diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl
index 8aedc9eb9f..4dce52e70f 100644
--- a/src/rabbit_mixed_queue.erl
+++ b/src/rabbit_mixed_queue.erl
@@ -103,7 +103,6 @@ to_mixed_mode(State = #mqstate { mode = disk, queue = 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),
- rabbit_log:info("Queue length: ~p ~w~n", [Q, erlang:length(QList)]),
{MsgBuf1, NextSeq1} =
lists:foldl(
fun ({MsgId, MsgBin, _Size, IsDelivered, _AckTag, SeqId}, {Buf, NSeq})
@@ -111,8 +110,12 @@ to_mixed_mode(State = #mqstate { mode = disk, queue = Q }) ->
Msg = #basic_message { guid = MsgId } = bin_to_msg(MsgBin),
{queue:in({SeqId, Msg, IsDelivered, true}, Buf), SeqId + 1}
end, {queue:new(), 0}, QList),
- {ok, State #mqstate { mode = mixed, msg_buf = MsgBuf1,
- next_write_seq = NextSeq1 }}.
+ 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)]),
+ {ok, State1}.
purge_non_persistent_messages(State = #mqstate { mode = disk, queue = Q,
is_durable = IsDurable }) ->
@@ -222,11 +225,13 @@ deliver(State = #mqstate { mode = mixed, queue = Q, is_durable = IsDurable,
IsDelivered, OnDisk}} ->
AckTag =
if OnDisk ->
- {MsgId, IsDelivered, AckTag2, _PersistRemaining} =
- rabbit_disk_queue:phantom_deliver(Q),
- if IsPersistent andalso IsDurable -> AckTag2;
- true -> ok = rabbit_disk_queue:ack(Q, [AckTag2]),
- noack
+ 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,