summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorMatthew Sackman <matthew@lshift.net>2009-11-10 14:07:28 +0000
committerMatthew Sackman <matthew@lshift.net>2009-11-10 14:07:28 +0000
commit2e5301bb55a8cb7556ec519abc043f868b0b847d (patch)
tree3631c6b3abf399853e008bfdd8aaa9d3949a73b0
parenta6a6707c97e74e586bb6fb5ee1e0172350c64a38 (diff)
downloadrabbitmq-server-git-2e5301bb55a8cb7556ec519abc043f868b0b847d.tar.gz
Fixed. There was a choice of having the qi do its own seqids, which would have been fine, but for the fact that at any point, there is the possibility of the vq deciding to flush everything out to disk. Thus to avoid rewriting in order to preserve order, instead we cope with the fact that there may be partial segments. The solution is to have a dict of segments which received fewer than the expected number of publishes before a different segment was written to. Then, when adding acks, we check to see if the number of acks is now equal to the partial number of segments, if it is, delete. Also, when moving to a different segment, we potentially delete the seg file if the number of acks is equal to the number of publishes. One key aspect of this is that the current segment to which we are publishing never appears in this dict, so it is not possible to delete the current publish segment by having the same number of pubs and acks, without it being really really full. Added and modified tests accordingly.
-rw-r--r--src/rabbit_queue_index.erl133
-rw-r--r--src/rabbit_tests.erl96
2 files changed, 160 insertions, 69 deletions
diff --git a/src/rabbit_queue_index.erl b/src/rabbit_queue_index.erl
index f21f9e17ee..a198ba51ee 100644
--- a/src/rabbit_queue_index.erl
+++ b/src/rabbit_queue_index.erl
@@ -117,7 +117,8 @@
journal_ack_dict,
journal_del_dict,
seg_ack_counts,
- publish_handle
+ publish_handle,
+ partial_segments
}).
-include("rabbit.hrl").
@@ -137,7 +138,8 @@
journal_ack_dict :: dict(),
journal_del_dict :: dict(),
seg_ack_counts :: dict(),
- publish_handle :: hdl_and_count()
+ publish_handle :: hdl_and_count(),
+ partial_segments :: dict()
}).
-spec(init/1 :: (queue_name()) -> {non_neg_integer(), qistate()}).
@@ -401,18 +403,40 @@ get_pub_handle(SegNum, State = #qistate { publish_handle = PubHandle }) ->
{Hdl, State1 #qistate { publish_handle = PubHandle1 }}.
get_counted_handle(SegNum, State, undefined) ->
+ get_counted_handle(SegNum, State, {SegNum, undefined, 0});
+get_counted_handle(SegNum, State = #qistate { partial_segments = Partials },
+ {SegNum, undefined, Count}) ->
{Hdl, State1} = get_seg_handle(SegNum, State),
- {State1, {SegNum, Hdl, 1}};
-get_counted_handle(SegNum, State, {SegNum, undefined, Count}) ->
- {Hdl, State1} = get_seg_handle(SegNum, State),
- {State1, {SegNum, Hdl, Count + 1}};
+ {CountExtra, Partials1} =
+ case dict:find(SegNum, Partials) of
+ {ok, CountExtra1} -> {CountExtra1, dict:erase(SegNum, Partials)};
+ error -> {0, Partials}
+ end,
+ Count1 = Count + 1 + CountExtra,
+ {State1 #qistate { partial_segments = Partials1 }, {SegNum, Hdl, Count1}};
get_counted_handle(SegNum, State, {SegNum, Hdl, Count})
when Count < ?SEGMENT_ENTRIES_COUNT ->
{State, {SegNum, Hdl, Count + 1}};
get_counted_handle(SegNumA, State, {SegNumB, Hdl, ?SEGMENT_ENTRIES_COUNT})
when SegNumA == SegNumB + 1 ->
ok = file_handle_cache:append_write_buffer(Hdl),
- get_counted_handle(SegNumA, State, undefined).
+ get_counted_handle(SegNumA, State, undefined);
+get_counted_handle(SegNumA, State = #qistate { partial_segments = Partials,
+ seg_ack_counts = AckCounts,
+ dir = Dir },
+ {SegNumB, Hdl, Count}) ->
+ %% don't flush here because it's possible SegNumB has been deleted
+ State1 =
+ case dict:find(SegNumB, AckCounts) of
+ {ok, Count} ->
+ %% #acks == #pubs, and we're moving to different
+ %% segment, so delete.
+ delete_segment(SegNumB, State);
+ _ ->
+ State #qistate {
+ partial_segments = dict:store(SegNumB, Count, Partials) }
+ end,
+ get_counted_handle(SegNumA, State1, undefined).
get_seg_handle(SegNum, State = #qistate { dir = Dir, seg_num_handles = SegHdls }) ->
case dict:find(SegNum, SegHdls) of
@@ -425,6 +449,17 @@ get_seg_handle(SegNum, State = #qistate { dir = Dir, seg_num_handles = SegHdls }
State)
end.
+delete_segment(SegNum, State = #qistate { dir = Dir,
+ seg_ack_counts = AckCounts,
+ partial_segments = Partials }) ->
+ State1 = close_handle(SegNum, State),
+ ok = case file:delete(seg_num_to_path(Dir, SegNum)) of
+ ok -> ok;
+ {error, enoent} -> ok
+ end,
+ State1 #qistate {seg_ack_counts = dict:erase(SegNum, AckCounts),
+ partial_segments = dict:erase(SegNum, Partials) }.
+
new_handle(Key, Path, Mode, State = #qistate { seg_num_handles = SegHdls }) ->
{ok, Hdl} = file_handle_cache:open(Path, Mode, [{write_buffer, infinity}]),
{Hdl, State #qistate { seg_num_handles = dict:store(Key, Hdl, SegHdls) }}.
@@ -486,7 +521,8 @@ blank_state(QueueName) ->
journal_ack_dict = dict:new(),
journal_del_dict = dict:new(),
seg_ack_counts = dict:new(),
- publish_handle = undefined
+ publish_handle = undefined,
+ partial_segments = dict:new()
}.
detect_clean_shutdown(Dir) ->
@@ -551,7 +587,8 @@ read_and_prune_segments(State = #qistate { dir = Dir }) ->
{TotalMsgCount, State1} =
lists:foldl(
fun (SegNum, {TotalMsgCount1, StateN =
- #qistate { publish_handle = PublishHandle }}) ->
+ #qistate { publish_handle = PublishHandle,
+ partial_segments = Partials }}) ->
{SDict, PubCount, AckCount, _HighRelSeq, StateM} =
load_segment(SegNum, StateN),
StateL = #qistate { seg_ack_counts = AckCounts } =
@@ -565,19 +602,30 @@ read_and_prune_segments(State = #qistate { dir = Dir }) ->
0 -> AckCounts;
N -> dict:store(SegNum, N, AckCounts)
end,
- %% In the following, there should only be max one
- %% segment that matches the 3rd case. All other
- %% segments should either be full or empty. There
- %% could be no partial segments.
- PublishHandle1 = case PubCount of
- ?SEGMENT_ENTRIES_COUNT -> PublishHandle;
- 0 -> PublishHandle;
- _ when PublishHandle == undefined ->
- {SegNum, undefined, PubCount}
- end,
+ %% In the following, whilst there may be several
+ %% partial segments, we only remember the last
+ %% one. All other partial segments get added into
+ %% the partial_segments dict
+ {PublishHandle1, Partials1} =
+ case PubCount of
+ ?SEGMENT_ENTRIES_COUNT ->
+ {PublishHandle, Partials};
+ 0 ->
+ {PublishHandle, Partials};
+ _ ->
+ {{SegNum, undefined, PubCount},
+ case PublishHandle of
+ undefined ->
+ Partials;
+ {SegNumOld, undefined, PubCountOld} ->
+ dict:store(SegNumOld, PubCountOld,
+ Partials)
+ end}
+ end,
{TotalMsgCount2,
StateL #qistate { seg_ack_counts = AckCounts1,
- publish_handle = PublishHandle1 }}
+ publish_handle = PublishHandle1,
+ partial_segments = Partials1 }}
end, {0, State}, SegNums),
{TotalMsgCount, State1}.
@@ -767,39 +815,40 @@ deliver_or_ack_msg(SDict, AckCount, RelSeq) ->
%%----------------------------------------------------------------------------
append_acks_to_segment(SegNum, Acks,
- State = #qistate { seg_ack_counts = AckCounts }) ->
+ State = #qistate { seg_ack_counts = AckCounts,
+ partial_segments = Partials }) ->
AckCount = case dict:find(SegNum, AckCounts) of
{ok, AckCount1} -> AckCount1;
error -> 0
end,
+ AckTarget = case dict:find(SegNum, Partials) of
+ {ok, PubCount} -> PubCount;
+ error -> ?SEGMENT_ENTRIES_COUNT
+ end,
AckCount2 = AckCount + length(Acks),
- AckCounts1 = case AckCount2 of
- 0 -> AckCounts;
- ?SEGMENT_ENTRIES_COUNT -> dict:erase(SegNum, AckCounts);
- _ -> dict:store(SegNum, AckCount2, AckCounts)
- end,
- append_acks_to_segment(SegNum, AckCount2, Acks,
- State #qistate { seg_ack_counts = AckCounts1 }).
-
-append_acks_to_segment(SegNum, AckCount, _Acks,
- State = #qistate { dir = Dir, publish_handle = PubHdl })
- when AckCount == ?SEGMENT_ENTRIES_COUNT ->
+ append_acks_to_segment(SegNum, AckCount2, Acks, AckTarget, State).
+
+append_acks_to_segment(SegNum, AckCount, _Acks, AckCount, State =
+ #qistate { publish_handle = PubHdl }) ->
PubHdl1 = case PubHdl of
- {SegNum, Hdl, ?SEGMENT_ENTRIES_COUNT} when Hdl /= undefined ->
+ %% If we're adjusting the pubhdl here then there
+ %% will be no entry in partials, thus the target ack
+ %% count must be SEGMENT_ENTRIES_COUNT
+ {SegNum, Hdl, AckCount = ?SEGMENT_ENTRIES_COUNT}
+ when Hdl /= undefined ->
{SegNum + 1, undefined, 0};
_ -> PubHdl
end,
- State1 = close_handle(SegNum, State #qistate { publish_handle = PubHdl1 }),
- ok = case file:delete(seg_num_to_path(Dir, SegNum)) of
- ok -> ok;
- {error, enoent} -> ok
- end,
- State1;
-append_acks_to_segment(SegNum, AckCount, Acks, State)
- when AckCount < ?SEGMENT_ENTRIES_COUNT ->
+ delete_segment(SegNum, State #qistate { publish_handle = PubHdl1 });
+append_acks_to_segment(_SegNum, _AckCount, [], _AckTarget, State) ->
+ State;
+append_acks_to_segment(SegNum, AckCount, Acks, AckTarget, State =
+ #qistate { seg_ack_counts = AckCounts })
+ when AckCount < AckTarget ->
{Hdl, State1} = append_to_segment(SegNum, Acks, State),
ok = file_handle_cache:sync(Hdl),
- State1.
+ State1 #qistate { seg_ack_counts =
+ dict:store(SegNum, AckCount, AckCounts) }.
append_dels_to_segment(SegNum, Dels, State) ->
{_Hdl, State1} = append_to_segment(SegNum, Dels, State),
diff --git a/src/rabbit_tests.erl b/src/rabbit_tests.erl
index d1131ed0af..d74f998e33 100644
--- a/src/rabbit_tests.erl
+++ b/src/rabbit_tests.erl
@@ -1025,6 +1025,23 @@ queue_index_publish(SeqIds, Persistent, Qi) ->
{QiM, [{SeqId, MsgId} | SeqIdsMsgIdsAcc]}
end, {Qi, []}, SeqIds).
+queue_index_deliver(SeqIds, Qi) ->
+ lists:foldl(
+ fun (SeqId, QiN) ->
+ rabbit_queue_index:write_delivered(SeqId, QiN)
+ end, Qi, SeqIds).
+
+queue_index_flush_journal(Qi) ->
+ {_Oks, {false, Qi1}} =
+ rabbit_misc:unfold(
+ fun ({true, QiN}) ->
+ QiM = rabbit_queue_index:flush_journal(QiN),
+ {true, ok, {rabbit_queue_index:can_flush_journal(QiM), QiM}};
+ ({false, _QiN}) ->
+ false
+ end, {true, Qi}),
+ Qi1.
+
verify_read_with_published(_Delivered, _Persistent, [], _) ->
ok;
verify_read_with_published(Delivered, Persistent,
@@ -1071,10 +1088,7 @@ test_queue_index() ->
{LenB, Qi12} = rabbit_queue_index:init(test_queue()),
{0, 20000, Qi13} =
rabbit_queue_index:find_lowest_seq_id_seg_and_next_seq_id(Qi12),
- Qi14 = lists:foldl(
- fun (SeqId, QiN) ->
- rabbit_queue_index:write_delivered(SeqId, QiN)
- end, Qi13, SeqIdsB),
+ Qi14 = queue_index_deliver(SeqIdsB, Qi13),
{ReadC, Qi15} = rabbit_queue_index:read_segment_entries(0, Qi14),
ok = verify_read_with_published(true, true, ReadC,
lists:reverse(SeqIdsMsgIdsB)),
@@ -1094,24 +1108,40 @@ test_queue_index() ->
_Qi21 = rabbit_queue_index:terminate_and_erase(Qi20),
ok = stop_msg_store(),
ok = empty_test_queue(),
- %% this next bit is just to hit the auto deletion of segment files
- SeqIdsC = lists:seq(0,65535),
+
+ %% These next bits are just to hit the auto deletion of segment files.
+ %% First, partials:
+ %% a) partial pub+del+ack, then move to new segment
+ SeqIdsC = lists:seq(0,trunc(SegmentSize/2)),
{0, Qi22} = rabbit_queue_index:init(test_queue()),
{Qi23, _SeqIdsMsgIdsC} = queue_index_publish(SeqIdsC, false, Qi22),
- Qi24 = lists:foldl(
- fun (SeqId, QiN) ->
- rabbit_queue_index:write_delivered(SeqId, QiN)
- end, Qi23, SeqIdsC),
+ Qi24 = queue_index_deliver(SeqIdsC, Qi23),
Qi25 = rabbit_queue_index:write_acks(SeqIdsC, Qi24),
- {_Oks, {false, Qi26}} =
- rabbit_misc:unfold(
- fun ({true, QiN}) ->
- QiM = rabbit_queue_index:flush_journal(QiN),
- {true, ok, {rabbit_queue_index:can_flush_journal(QiM), QiM}};
- ({false, _QiN}) ->
- false
- end, {true, Qi25}),
- _Qi27 = rabbit_queue_index:terminate_and_erase(Qi26),
+ Qi26 = queue_index_flush_journal(Qi25),
+ {Qi27, _SeqIdsMsgIdsC1} = queue_index_publish([SegmentSize], false, Qi26),
+ _Qi28 = rabbit_queue_index:terminate_and_erase(Qi27),
+ ok = stop_msg_store(),
+ ok = empty_test_queue(),
+
+ %% b) partial pub+del, then move to new segment, then ack all in old segment
+ {0, Qi29} = rabbit_queue_index:init(test_queue()),
+ {Qi30, _SeqIdsMsgIdsC2} = queue_index_publish(SeqIdsC, false, Qi29),
+ Qi31 = queue_index_deliver(SeqIdsC, Qi30),
+ {Qi32, _SeqIdsMsgIdsC3} = queue_index_publish([SegmentSize], false, Qi31),
+ Qi33 = rabbit_queue_index:write_acks(SeqIdsC, Qi32),
+ Qi34 = queue_index_flush_journal(Qi33),
+ _Qi35 = rabbit_queue_index:terminate_and_erase(Qi34),
+ ok = stop_msg_store(),
+ ok = empty_test_queue(),
+
+ %% c) just fill up several segments of all pubs, then +dels, then +acks
+ SeqIdsD = lists:seq(0,SegmentSize*4),
+ {0, Qi36} = rabbit_queue_index:init(test_queue()),
+ {Qi37, _SeqIdsMsgIdsD} = queue_index_publish(SeqIdsD, false, Qi36),
+ Qi38 = queue_index_deliver(SeqIdsD, Qi37),
+ Qi39 = rabbit_queue_index:write_acks(SeqIdsD, Qi38),
+ Qi40 = queue_index_flush_journal(Qi39),
+ _Qi41 = rabbit_queue_index:terminate_and_erase(Qi40),
ok = stop_msg_store(),
ok = rabbit_queue_index:start_msg_store([]),
ok = stop_msg_store(),
@@ -1167,14 +1197,26 @@ test_variable_queue_dynamic_duration_change() ->
VQ0 = fresh_variable_queue(),
%% start by sending in a couple of segments worth
Len1 = 2*SegmentSize,
- {_SeqIds, VQ1} = variable_queue_publish(true, Len1, VQ0),
+ {_SeqIds, VQ1} = variable_queue_publish(false, Len1, VQ0),
VQ2 = rabbit_variable_queue:remeasure_egress_rate(VQ1),
- {ok, _TRef} = timer:send_after(1000, {duration, 30, fun erlang:'-'/2}),
+ {ok, _TRef} = timer:send_after(1000, {duration, 60,
+ fun (V) -> (V*0.75)-1 end}),
VQ3 = test_variable_queue_dynamic_duration_change_f(Len1, VQ2),
{VQ4, AckTags} = variable_queue_fetch(Len1, false, false, Len1, VQ3),
VQ5 = rabbit_variable_queue:ack(AckTags, VQ4),
{empty, VQ6} = rabbit_variable_queue:fetch(VQ5),
- rabbit_variable_queue:terminate(VQ6),
+
+ %% just publish and fetch some persistent msgs, this hits the the
+ %% partial segment path in queue_index due to the period when
+ %% duration was 0 and the entire queue was gamma.
+ {_SeqIds1, VQ7} = variable_queue_publish(true, 20, VQ6),
+ {VQ8, AckTags1} = variable_queue_fetch(20, true, false, 20, VQ7),
+ VQ9 = rabbit_variable_queue:ack(AckTags1, VQ8),
+ VQ10 = rabbit_variable_queue:flush_journal(VQ9),
+ VQ11 = rabbit_variable_queue:flush_journal(VQ10),
+ {empty, VQ12} = rabbit_variable_queue:fetch(VQ11),
+
+ rabbit_variable_queue:terminate(VQ12),
passed.
@@ -1183,14 +1225,14 @@ test_variable_queue_dynamic_duration_change_f(Len, VQ0) ->
{{_Msg, false, AckTag, Len}, VQ2} = rabbit_variable_queue:fetch(VQ1),
VQ3 = rabbit_variable_queue:ack([AckTag], VQ2),
receive
- {duration, 30, stop} ->
+ {duration, _, stop} ->
VQ3;
{duration, N, Fun} ->
- N1 = Fun(N, 1),
+ N1 = lists:max([Fun(N), 0]),
Fun1 = case N1 of
- 0 -> fun erlang:'+'/2;
- 30 -> stop;
- _ -> Fun
+ 0 -> fun (V) -> (V+1)/0.75 end;
+ _ when N1 > 400 -> stop;
+ _ -> Fun
end,
{ok, _TRef} = timer:send_after(1000, {duration, N1, Fun1}),
VQ4 = rabbit_variable_queue:remeasure_egress_rate(VQ3),