diff options
| author | Matthew Sackman <matthew@lshift.net> | 2009-08-17 12:23:10 +0100 |
|---|---|---|
| committer | Matthew Sackman <matthew@lshift.net> | 2009-08-17 12:23:10 +0100 |
| commit | b99edb893e39ac72d8796611952ec74c36ad7e6b (patch) | |
| tree | 62cf05268b75330e3e4e7a901228fd9b97bcc950 | |
| parent | 376e6350bbbeaa4c891fd10c7d577413cf9d96be (diff) | |
| download | rabbitmq-server-git-b99edb893e39ac72d8796611952ec74c36ad7e6b.tar.gz | |
Reworking of mixed queue.
| -rw-r--r-- | src/rabbit_mixed_queue.erl | 122 |
1 files changed, 48 insertions, 74 deletions
diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl index f798b36951..4b0810a8dd 100644 --- a/src/rabbit_mixed_queue.erl +++ b/src/rabbit_mixed_queue.erl @@ -165,12 +165,11 @@ to_disk_only_mode(TxnMessages, State = send_messages_to_disk(IsDurable, Q, Queue, PublishCount, RequeueCount, Commit, MsgBuf) -> case queue:out(Queue) of - {empty, Queue} -> + {empty, _Queue} -> ok = flush_messages_to_disk_queue(Q, Commit), [] = flush_requeue_to_disk_queue(Q, RequeueCount, []), {ok, MsgBuf}; - {{value, {Msg = #basic_message { guid = MsgId, - is_persistent = IsPersistent }, + {{value, {Msg = #basic_message { is_persistent = IsPersistent }, IsDelivered}}, Queue1} -> case IsDurable andalso IsPersistent of true -> %% it's already in the Q @@ -178,61 +177,47 @@ send_messages_to_disk(IsDurable, Q, Queue, PublishCount, RequeueCount, IsDurable, Q, Queue1, PublishCount, RequeueCount + 1, Commit, inc_queue_length(Q, MsgBuf, 1)); false -> - Commit1 = flush_requeue_to_disk_queue - (Q, RequeueCount, Commit), - ok = rabbit_disk_queue:tx_publish(Msg), - case PublishCount == ?TO_DISK_MAX_FLUSH_SIZE of - true -> - ok = flush_messages_to_disk_queue(Q, Commit1), - send_messages_to_disk( - IsDurable, Q, Queue1, 1, 0, - [{MsgId, IsDelivered}], - inc_queue_length(Q, MsgBuf, 1)); - false -> - send_messages_to_disk( - IsDurable, Q, Queue1, PublishCount + 1, 0, - [{MsgId, IsDelivered} | Commit1], - inc_queue_length(Q, MsgBuf, 1)) - end + republish_message_to_disk_queue( + IsDurable, Q, Queue1, PublishCount, RequeueCount, Commit, + MsgBuf, Msg, IsDelivered) end; - {{value, {Msg = #basic_message { guid = MsgId }, IsDelivered, _AckTag}}, - Queue1} -> + {{value, {Msg, IsDelivered, _AckTag}}, Queue1} -> %% these have come via the prefetcher, so are no longer in %% the disk queue so they need to be republished - Commit1 = flush_requeue_to_disk_queue(Q, RequeueCount, Commit), - ok = rabbit_disk_queue:tx_publish(Msg), - case PublishCount == ?TO_DISK_MAX_FLUSH_SIZE of - true -> - ok = flush_messages_to_disk_queue(Q, Commit1), - send_messages_to_disk(IsDurable, Q, Queue1, 1, 0, - [{MsgId, IsDelivered}], - inc_queue_length(Q, MsgBuf, 1)); - false -> - send_messages_to_disk(IsDurable, Q, Queue1, PublishCount+1, - 0, [{MsgId, IsDelivered} | Commit1], - inc_queue_length(Q, MsgBuf, 1)) - end; + republish_message_to_disk_queue(IsDelivered, Q, Queue1, + PublishCount, RequeueCount, Commit, + MsgBuf, Msg, IsDelivered); {{value, {Q, Count}}, Queue1} -> send_messages_to_disk(IsDurable, Q, Queue1, PublishCount, RequeueCount + Count, Commit, inc_queue_length(Q, MsgBuf, Count)) end. +republish_message_to_disk_queue(IsDurable, Q, Queue, PublishCount, RequeueCount, + Commit, MsgBuf, Msg = + #basic_message { guid = MsgId }, IsDelivered) -> + Commit1 = flush_requeue_to_disk_queue(Q, RequeueCount, Commit), + ok = rabbit_disk_queue:tx_publish(Msg), + {PublishCount1, Commit2} = + case PublishCount == ?TO_DISK_MAX_FLUSH_SIZE of + true -> ok = flush_messages_to_disk_queue(Q, Commit1), + {1, [{MsgId, IsDelivered}]}; + false -> {PublishCount + 1, [{MsgId, IsDelivered} | Commit1]} + end, + send_messages_to_disk(IsDurable, Q, Queue, PublishCount1, 0, + Commit2, inc_queue_length(Q, MsgBuf, 1)). + +flush_messages_to_disk_queue(_Q, []) -> + ok; flush_messages_to_disk_queue(Q, Commit) -> - ok = if [] == Commit -> ok; - true -> rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), []) - end. + rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), []). +flush_requeue_to_disk_queue(_Q, 0, Commit) -> + Commit; flush_requeue_to_disk_queue(Q, RequeueCount, Commit) -> - if 0 == RequeueCount -> Commit; - true -> - ok = if [] == Commit -> ok; - true -> rabbit_disk_queue:tx_commit - (Q, lists:reverse(Commit), []) - end, - rabbit_disk_queue:requeue_next_n(Q, RequeueCount), - [] - end. + ok = flush_messages_to_disk_queue(Q, Commit), + ok = rabbit_disk_queue:requeue_next_n(Q, RequeueCount), + []. to_mixed_mode(_TxnMessages, State = #mqstate { mode = mixed }) -> {ok, State}; @@ -266,14 +251,13 @@ to_mixed_mode(TxnMessages, State = #mqstate { mode = disk, queue = Q, inc_queue_length(_Q, MsgBuf, 0) -> MsgBuf; inc_queue_length(Q, MsgBuf, Count) -> - case queue:out_r(MsgBuf) of - {empty, MsgBuf} -> - queue:in({Q, Count}, MsgBuf); - {{value, {Q, Len}}, MsgBuf1} -> - queue:in({Q, Len + Count}, MsgBuf1); - {{value, _}, _MsgBuf1} -> - queue:in({Q, Count}, MsgBuf) - end. + {NewCount, MsgBufTail} = + case queue:out_r(MsgBuf) of + {empty, MsgBuf1} -> {Count, MsgBuf1}; + {{value, {Q, Len}}, MsgBuf1} -> {Len + Count, MsgBuf1}; + {{value, _}, _MsgBuf1} -> {Count, MsgBuf} + end, + queue:in({Q, NewCount}, MsgBufTail). dec_queue_length(Count, State = #mqstate { queue = Q, msg_buf = MsgBuf }) -> case queue:out(MsgBuf) of @@ -314,8 +298,7 @@ publish(Msg = #basic_message { is_persistent = IsPersistent }, State = #mqstate { queue = Q, mode = mixed, is_durable = IsDurable, msg_buf = MsgBuf, length = Length, memory_size = QSize, memory_gain = Gain }) -> - Persist = IsDurable andalso IsPersistent, - ok = case Persist of + ok = case IsDurable andalso IsPersistent of true -> rabbit_disk_queue:publish(Q, Msg, false); false -> ok end, @@ -333,12 +316,11 @@ publish_delivered(Msg = queue = Q, length = 0, memory_size = QSize, memory_gain = Gain }) when Mode =:= disk orelse (IsDurable andalso IsPersistent) -> - Persist = IsDurable andalso IsPersistent, ok = rabbit_disk_queue:publish(Q, Msg, true), MsgSize = size_of_message(Msg), State1 = State #mqstate { memory_size = QSize + MsgSize, memory_gain = Gain + MsgSize }, - case Persist of + case IsDurable andalso IsPersistent of true -> %% must call phantom_deliver otherwise the msg remains at %% the head of the queue. This is synchronous, but @@ -386,14 +368,7 @@ deliver(State = #mqstate { msg_buf = MsgBuf, queue = Q, %% message has come via the prefetcher, thus it's been %% delivered. If it's not persistent+durable, we should %% ack it now - AckTag1 = - case IsDurable andalso IsPersistent of - true -> - AckTag; - false -> - ok = rabbit_disk_queue:ack(Q, [AckTag]), - noack - end, + AckTag1 = maybe_ack(Q, IsDurable, IsPersistent, AckTag), {{Msg, IsDelivered, AckTag1, Rem}, State1 #mqstate { msg_buf = MsgBuf1 }}; _ when Prefetcher == undefined -> @@ -401,14 +376,7 @@ deliver(State = #mqstate { msg_buf = MsgBuf, queue = Q, {Msg = #basic_message { is_persistent = IsPersistent }, _Size, IsDelivered, AckTag, _PersistRem} = rabbit_disk_queue:deliver(Q), - AckTag1 = - case IsDurable andalso IsPersistent of - true -> - AckTag; - false -> - ok = rabbit_disk_queue:ack(Q, [AckTag]), - noack - end, + AckTag1 = maybe_ack(Q, IsDurable, IsPersistent, AckTag), {{Msg, IsDelivered, AckTag1, Rem}, State2}; _ -> case rabbit_queue_prefetcher:drain(Prefetcher) of @@ -425,6 +393,12 @@ deliver(State = #mqstate { msg_buf = MsgBuf, queue = Q, end end. +maybe_ack(_Q, true, true, AckTag) -> + AckTag; +maybe_ack(Q, _, _, AckTag) -> + ok = rabbit_disk_queue:ack(Q, [AckTag]), + noack. + remove_noacks(MsgsWithAcks) -> lists:foldl( fun ({Msg, noack}, {AccAckTags, AccSize}) -> |
