diff options
| author | Matthew Sackman <matthew@lshift.net> | 2009-06-10 21:02:54 +0100 |
|---|---|---|
| committer | Matthew Sackman <matthew@lshift.net> | 2009-06-10 21:02:54 +0100 |
| commit | a99b541be0332defa32c8ce7548acd9830c56892 (patch) | |
| tree | c2b2601b141f3296ec8465d0a278132e2c84e015 | |
| parent | cf5e1ae896b8aceca66c776b3b98806eb3bfefcf (diff) | |
| download | rabbitmq-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.erl | 25 | ||||
| -rw-r--r-- | src/rabbit_mixed_queue.erl | 21 |
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, |
