diff options
| -rw-r--r-- | src/rabbit_disk_queue.erl | 47 | ||||
| -rw-r--r-- | src/rabbit_mixed_queue.erl | 37 | ||||
| -rw-r--r-- | src/rabbit_queue_mode_manager.erl | 6 |
3 files changed, 62 insertions, 28 deletions
diff --git a/src/rabbit_disk_queue.erl b/src/rabbit_disk_queue.erl index d8c6580f01..9c7d35eb4e 100644 --- a/src/rabbit_disk_queue.erl +++ b/src/rabbit_disk_queue.erl @@ -41,7 +41,8 @@ -export([publish/3, deliver/1, phantom_deliver/1, ack/2, tx_publish/1, tx_commit/3, tx_cancel/1, requeue/2, requeue_with_seqs/2, purge/1, delete_queue/1, - delete_non_durable_queues/1, auto_ack_next_message/1 + delete_non_durable_queues/1, auto_ack_next_message/1, + requeue_next_n/2 ]). -export([length/1, filesync/0, cache_info/0]). @@ -264,6 +265,7 @@ -spec(requeue_with_seqs/2 :: (queue_name(), [{{msg_id(), seq_id()}, {seq_id_or_next(), bool()}}]) -> 'ok'). +-spec(requeue_next_n/2 :: (queue_name(), non_neg_integer()) -> 'ok'). -spec(purge/1 :: (queue_name()) -> non_neg_integer()). -spec(delete_non_durable_queues/1 :: (set()) -> 'ok'). -spec(stop/0 :: () -> 'ok'). @@ -315,6 +317,9 @@ requeue(Q, MsgSeqIds) when is_list(MsgSeqIds) -> requeue_with_seqs(Q, MsgSeqSeqIds) when is_list(MsgSeqSeqIds) -> gen_server2:cast(?SERVER, {requeue_with_seqs, Q, MsgSeqSeqIds}). +requeue_next_n(Q, N) when is_integer(N) -> + gen_server2:cast(?SERVER, {requeue_next_n, Q, N}). + purge(Q) -> gen_server2:call(?SERVER, {purge, Q}, infinity). @@ -509,6 +514,9 @@ handle_cast({requeue, Q, MsgSeqIds}, State) -> handle_cast({requeue_with_seqs, Q, MsgSeqSeqIds}, State) -> {ok, State1} = internal_requeue(Q, MsgSeqSeqIds, State), noreply(State1); +handle_cast({requeue_next_n, Q, N}, State) -> + {ok, State1} = internal_requeue_next_n(Q, N, State), + noreply(State1); handle_cast({delete_queue, Q}, State) -> {ok, State1} = internal_delete_queue(Q, State), noreply(State1); @@ -747,13 +755,19 @@ find_next_seq_id(CurrentSeq, NextSeqId) when NextSeqId > CurrentSeq -> NextSeqId. +%% the queue is empty, and we've just written exactly where we +%% expected, so read it back determine_next_read_id(CurrentReadWrite, CurrentReadWrite, CurrentReadWrite) -> CurrentReadWrite; +%% we've just written in the next slot, so the next read pos is unaltered determine_next_read_id(CurrentRead, _CurrentWrite, next) -> CurrentRead; +%% queue is empty, but we've written somewhere else - a gap has formed +%% - so read back from where we wrote, after the gap determine_next_read_id(CurrentReadWrite, CurrentReadWrite, NextWrite) when NextWrite > CurrentReadWrite -> NextWrite; +%% queue is not empty, and we've created a gap, so the read pos is unaltered determine_next_read_id(CurrentRead, CurrentWrite, NextWrite) when NextWrite >= CurrentWrite -> CurrentRead. @@ -1200,6 +1214,37 @@ requeue_message({{{MsgId, SeqIdOrig}, {SeqIdTo, NewIsDelivered}}, decrement_cache(MsgId, State), {NextSeqIdTo1, Q, State}. +%% move the next N messages from the front of the queue to the back. +internal_requeue_next_n(Q, N, State = #dqstate { sequences = Sequences }) -> + {ReadSeqId, WriteSeqId, Length} = sequence_lookup(Sequences, Q), + ReadSeqId1 = determine_next_read_id(ReadSeqId, WriteSeqId, next), + if N >= Length -> {ok, State}; + true -> + {atomic, {ReadSeqIdN, WriteSeqIdN}} = + mnesia:transaction( + fun() -> + ok = mnesia:write_lock_table(rabbit_disk_queue), + requeue_next_messages(Q, N, ReadSeqId1, WriteSeqId) + end + ), + true = ets:insert(Sequences, {Q, ReadSeqIdN, WriteSeqIdN, Length}), + {ok, State} + end. + +requeue_next_messages(_Q, 0, ReadSeq, WriteSeq) -> + {ReadSeq, WriteSeq}; +requeue_next_messages(Q, N, ReadSeq, WriteSeq) -> + WriteSeq1 = adjust_last_msg_seq_id(Q, WriteSeq, next, write), + NextWriteSeq = find_next_seq_id(WriteSeq1, next), + [Obj = #dq_msg_loc { next_seq_id = NextSeqIdOrig }] = + mnesia:read(rabbit_disk_queue, {Q, ReadSeq}, write), + ok = mnesia:write(rabbit_disk_queue, + Obj #dq_msg_loc {queue_and_seq_id = {Q, WriteSeq1}, + next_seq_id = NextWriteSeq + }, write), + ok = mnesia:delete(rabbit_disk_queue, {Q, ReadSeq}, write), + requeue_next_messages(Q, N - 1, NextSeqIdOrig, NextWriteSeq). + internal_purge(Q, State = #dqstate { sequences = Sequences }) -> case ets:lookup(Sequences, Q) of [] -> {ok, 0, State}; diff --git a/src/rabbit_mixed_queue.erl b/src/rabbit_mixed_queue.erl index f415472719..4a2803a4ef 100644 --- a/src/rabbit_mixed_queue.erl +++ b/src/rabbit_mixed_queue.erl @@ -126,7 +126,7 @@ to_disk_only_mode(TxnMessages, State = %% message on disk. %% Note we also batch together messages on disk so that we minimise %% the calls to requeue. - ok = send_messages_to_disk(Q, MsgBuf, [], 0, []), + ok = send_messages_to_disk(Q, MsgBuf, 0, 0, []), %% tx_publish txn messages. Some of these will have been already %% published if they really are durable and persistent which is %% why we can't just use our own tx_publish/2 function (would end @@ -141,47 +141,36 @@ to_disk_only_mode(TxnMessages, State = garbage_collect(), {ok, State #mqstate { mode = disk, msg_buf = queue:new() }}. -send_messages_to_disk(Q, Queue, Requeue, PublishCount, Commit) -> +send_messages_to_disk(Q, Queue, RequeueCount, PublishCount, Commit) -> case queue:out(Queue) of {empty, Queue} -> ok = flush_messages_to_disk_queue(Q, Commit), - [] = flush_requeue_to_disk_queue(Q, Requeue, []), + [] = flush_requeue_to_disk_queue(Q, RequeueCount, []), ok; - {{value, {Msg = #basic_message { guid = MsgId }, IsDelivered, OnDisk}}, + {{value, {Msg = #basic_message { guid = MsgId }, _IsDelivered, OnDisk}}, Queue1} -> case OnDisk of true -> - ok = flush_messages_to_disk_queue (Q, Commit), - {MsgId, IsDelivered, AckTag, _PersistRemaining} = - rabbit_disk_queue:phantom_deliver(Q), + ok = flush_messages_to_disk_queue(Q, Commit), send_messages_to_disk( - Q, Queue1, [{AckTag, {next, IsDelivered}} | Requeue], - 0, []); + Q, Queue1, 1 + RequeueCount, 0, []); false -> Commit1 = - flush_requeue_to_disk_queue(Q, Requeue, Commit), + 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(Q, Queue1, [], 1, [MsgId]); + send_messages_to_disk(Q, Queue1, 0, 1, [MsgId]); false -> send_messages_to_disk - (Q, Queue1, [], PublishCount + 1, + (Q, Queue1, 0, PublishCount + 1, [MsgId | Commit1]) end end; {{value, {disk, Count}}, Queue2} -> ok = flush_messages_to_disk_queue(Q, Commit), - {Requeue1, 0} = - rabbit_misc:unfold( - fun (0) -> false; - (N) -> - {_MsgId, IsDelivered, AckTag, _PersistRemaining} - = rabbit_disk_queue:phantom_deliver(Q), - {true, {AckTag, {next, IsDelivered}}, N - 1} - end, Count), - send_messages_to_disk(Q, Queue2, Requeue1 ++ Requeue, 0, []) + send_messages_to_disk(Q, Queue2, RequeueCount + Count, 0, []) end. flush_messages_to_disk_queue(Q, Commit) -> @@ -189,10 +178,10 @@ flush_messages_to_disk_queue(Q, Commit) -> true -> rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), []) end. -flush_requeue_to_disk_queue(Q, Requeue, Commit) -> - if [] == Requeue -> Commit; +flush_requeue_to_disk_queue(Q, RequeueCount, Commit) -> + if 0 == RequeueCount -> Commit; true -> ok = rabbit_disk_queue:tx_commit(Q, lists:reverse(Commit), []), - rabbit_disk_queue:requeue_with_seqs(Q, lists:reverse(Requeue)), + rabbit_disk_queue:requeue_next_n(Q, RequeueCount), [] end. diff --git a/src/rabbit_queue_mode_manager.erl b/src/rabbit_queue_mode_manager.erl index 99f6e408d9..f5cc32b445 100644 --- a/src/rabbit_queue_mode_manager.erl +++ b/src/rabbit_queue_mode_manager.erl @@ -201,20 +201,20 @@ handle_call({pin_to_disk, Pid}, _From, disk_mode_pins = Pins }) -> {Res, State1} = case sets:is_element(Pid, Pins) of - true -> {already_pinned, State}; + true -> {ok, State}; false -> case find_queue(Pid, Mixed) of {mixed, {OAlloc, _OActivity}} -> {Module, Function, Args} = dict:fetch(Pid, Callbacks), ok = erlang:apply(Module, Function, Args ++ [disk]), - {convert_to_disk_mode, + {ok, State #state { mixed_queues = dict:erase(Pid, Mixed), available_tokens = Avail + OAlloc, disk_mode_pins = sets:add_element(Pid, Pins) }}; disk -> - {already_disk, + {ok, State #state { disk_mode_pins = sets:add_element(Pid, Pins) }} end |
