jimczi commented on code in PR #16571:
URL: https://github.com/apache/lucene/pull/16571#discussion_r3926562637
##########
lucene/core/src/java/org/apache/lucene/index/IndexWriter.java:
##########
@@ -5214,11 +5254,463 @@ private boolean assertSoftDeletesCount(CodecReader
reader, int expectedCount) th
return true;
}
+ /**
+ * Turns one finished merge output into a segment on disk: packs its files
into a compound file if
+ * the codec asks for one, writes its segment info, and warms it.
+ *
+ * <p>Shared by the single-output and the partitioned merge paths, which
differ in how many
+ * outputs they produce and not in what packaging one of them means. The
delete-on-failure and
+ * abort handling here is the delicate part, and is the reason this is one
method rather than two
+ * similar ones.
+ *
+ * @return false if the merge was aborted while this ran, in which case the
caller must abandon it
+ */
+ private boolean packageMergedSegment(
+ MergePolicy.OneMerge merge,
+ MergePolicy mergePolicy,
+ SegmentCommitInfo info,
+ IOContext context)
+ throws IOException {
+ // Very important to do this before opening the reader
+ // because codec must know if prox was written for
+ // this segment:
+ boolean useCompoundFile;
+ synchronized (this) { // Guard segmentInfos
+ useCompoundFile =
+ info.info
+ .getCodec()
+ .compoundFormat()
+ .useCompoundFile(mergePolicy.size(info, this), mergePolicy);
+ }
+
+ if (useCompoundFile) {
+ Collection<String> filesToRemove = info.files();
+ // NOTE: Creation of the CFS file must be performed with the original
+ // directory rather than with the merging directory, so that it is not
+ // subject to merge throttling.
+ TrackingDirectoryWrapper trackingCFSDir = new
TrackingDirectoryWrapper(directory);
+ try {
+ createCompoundFile(infoStream, trackingCFSDir, info.info, context,
this::deleteNewFiles);
+ } catch (Throwable t) {
+ try {
+ synchronized (this) {
+ if (merge.isAborted()) {
+ // This can happen if rollback is called while we were building
+ // our CFS -- fall through to logic below to remove the non-CFS
+ // merged files:
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message(
+ "IW", "hit merge abort exception creating compound file
during merge: " + t);
+ }
+ return false;
+ } else {
+ handleMergeException(t, merge);
+ }
+ }
+ } finally {
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message("IW", "hit exception creating compound file
during merge: " + t);
+ }
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ }
+ }
+
+ synchronized (this) {
+
+ // delete new non cfs files directly: they were never
+ // registered with IFD
+ deleteNewFiles(filesToRemove);
+
+ if (merge.isAborted()) {
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message("IW", "abort merge after building CFS");
+ }
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ return false;
+ }
+ }
+
+ info.info.setUseCompoundFile(true);
+ }
+
+ // Have codec write SegmentInfo. Must do this after
+ // creating CFS so that 1) .si isn't slurped into CFS,
+ // and 2) .si reflects useCompoundFile=true change
+ // above:
+ try {
+ config.getCodec().segmentInfoFormat().write(directory, info.info,
context);
+ } catch (Throwable t) {
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ throw t;
+ }
+
+ // TODO: ideally we would freeze info here!!
+ // because any changes after writing the .si will be
+ // lost...
+
+ final IndexReaderWarmer mergedSegmentWarmer =
config.getMergedSegmentWarmer();
+ if (readerPool.isReaderPoolingEnabled() && mergedSegmentWarmer != null) {
+ final ReadersAndUpdates rld = getPooledInstance(info, true);
+ final SegmentReader sr = rld.getReader(IOContext.DEFAULT);
+ try {
+ mergedSegmentWarmer.warm(sr);
+ } finally {
+ synchronized (this) {
+ rld.release(sr);
+ release(rld);
+ }
+ }
+ }
+ return true;
+ }
+
+ /** Opens a pooled reader for each of a merge's input segments, and holds
their files open. */
+ private void initMergeReaders(MergePolicy.OneMerge merge, IOContext context)
throws IOException {
+ merge.initMergeReaders(
+ sci -> {
+ final ReadersAndUpdates rld = getPooledInstance(sci, true);
+ rld.setIsMerging();
+ synchronized (this) {
+ return rld.getReaderForMerge(
+ context, mr ->
deleter.incRef(mr.reader.getSegmentInfo().files()));
+ }
+ });
+ }
+
+ /** Fresh {@link SegmentCommitInfo} for one output of a merge. */
+ private SegmentCommitInfo newMergeSegmentInfo(MergePolicy.OneMerge merge) {
+ boolean hasBlocks = false;
+ for (SegmentCommitInfo info : merge.segments) {
+ if (info.info.getHasBlocks()) {
+ hasBlocks = true;
+ break;
+ }
+ }
+ SegmentInfo si =
+ new SegmentInfo(
+ directoryOrig,
+ Version.LATEST,
+ null,
+ newSegmentName(),
+ -1,
+ false,
+ hasBlocks,
+ config.getCodec(),
+ Collections.emptyMap(),
+ StringHelper.randomId(),
+ Collections.emptyMap(),
+ config.getIndexSort());
+ Map<String, String> details = new HashMap<>();
+ details.put("mergeMaxNumSegments", "" + merge.maxNumSegments);
+ details.put("mergeFactor", Integer.toString(merge.segments.size()));
+ details.put("mergeOutputs", Integer.toString(merge.getOutputCount()));
+ setDiagnostics(si, SOURCE_MERGE, details);
+ return new SegmentCommitInfo(si, 0, 0, -1L, -1L, -1L,
StringHelper.randomId());
+ }
+
+ /**
+ * Validate a doc-range partition spec. Only the SHAPE is checked -- one
boundary array per input,
+ * all the same length, non-decreasing, starting at 0 and ending at that
input's maxDoc. Given
+ * that shape, the two properties the merge depends on, disjointness and
full coverage, hold by
+ * construction rather than by trusting the caller.
+ *
+ * <p>Plus the one thing that is not about the spec: an index sort.
Correctness rests on the
+ * outputs being contiguous AND in order in the merged document space, which
holds under a sort,
+ * where the merged order is key order and the outputs are key ranges.
Without one, {@link
+ * DocIDMerger} concatenates input by input, so an output's documents land
in one block per input
+ * rather than in a single run, and splitting a term's postings by document
id would quietly hand
+ * documents to the wrong output.
+ */
+ private void validateDocRangePartitions(
+ MergePolicy.OneMerge merge, int[][] partitions, List<CodecReader>
readers)
+ throws IOException {
+ if (partitions.length != merge.segments.size()) {
+ throw new IllegalArgumentException(
+ "docRangePartitions has "
+ + partitions.length
+ + " entries but the merge has "
+ + merge.segments.size()
+ + " input segments");
+ }
+ final int outputs = merge.getOutputCount();
+ if (outputs < 1) {
+ throw new IllegalArgumentException("a partitioned merge must have at
least one output");
+ }
+ for (int i = 0; i < partitions.length; i++) {
+ final int[] b = partitions[i];
+ final int maxDoc = merge.segments.get(i).info.maxDoc();
+ if (b.length != outputs + 1) {
+ throw new IllegalArgumentException(
+ "docRangePartitions[" + i + "] has length " + b.length + ",
expected " + (outputs + 1));
+ }
+ if (b[0] != 0 || b[outputs] != maxDoc) {
+ throw new IllegalArgumentException(
+ "docRangePartitions["
+ + i
+ + "] must span [0, "
+ + maxDoc
+ + "], got ["
+ + b[0]
+ + ", "
+ + b[outputs]
+ + "]");
+ }
+ for (int o = 1; o <= outputs; o++) {
+ if (b[o] < b[o - 1]) {
+ throw new IllegalArgumentException(
+ "docRangePartitions[" + i + "] is not non-decreasing at " + o);
+ }
+ }
+ checkBoundariesRespectBlocks(readers.get(i), b, i);
+ }
+ }
+
+ /**
+ * Refuses a boundary that falls inside a document block.
+ *
+ * <p>A block is a run of documents ending at its parent, and queries over
it find the children by
+ * counting back from the parent. An index sort keeps a block contiguous,
but says nothing about
+ * where a partition may cut, so a boundary landing between a child and its
parent would put them
+ * in different segments and leave both halves quietly wrong.
+ *
+ * <p>Carrying the same partitioning value on the children is not sufficient
on its own: the
+ * boundary is a document offset, and whether the offset a caller derives
from that value lands
+ * before the children or between them and their parent depends on whether
the children carry the
+ * value at all. So the invariant is enforced here, on the offsets
themselves. A policy that wants
+ * to partition a block index has the readers when it chooses its
boundaries, and can align them
+ * by reading the same parent field this does.
+ */
+ private static void checkBoundariesRespectBlocks(CodecReader reader, int[]
b, int input)
+ throws IOException {
+ final String parentField = reader.getFieldInfos().getParentField();
+ if (parentField == null) {
+ return;
+ }
+ final NumericDocValues parents = reader.getNumericDocValues(parentField);
+ if (parents == null) {
+ return;
+ }
+ final int maxDoc = reader.maxDoc();
+ int previous = -1;
+ for (int o = 1; o < b.length - 1; o++) {
+ final int boundary = b[o];
+ // The ends of the space are block boundaries by construction, and a
repeated boundary is an
+ // empty output, which cannot split anything the previous one did not.
+ if (boundary == 0 || boundary == maxDoc || boundary == previous) {
+ continue;
+ }
+ previous = boundary;
+ // A boundary is legal exactly when the document before it ends a block.
+ if (parents.advanceExact(boundary - 1) == false) {
+ throw new IllegalArgumentException(
+ "docRangePartitions["
+ + input
+ + "] cuts inside a document block at "
+ + boundary
+ + ": document "
+ + (boundary - 1)
+ + " is not the last of its block, and a partitioned merge must
not separate a "
+ + "block's documents from their parent");
+ }
+ }
+ }
+
+ /**
+ * Merge producing several output segments, each holding a contiguous doc
range of every input.
+ * Kept separate from {@link #mergeMiddle} so the single-output path stays
untouched.
+ */
+ private int multiOutputMergeMiddle(MergePolicy.OneMerge merge, MergePolicy
mergePolicy)
+ throws IOException {
+ testPoint("mergeMiddleStart");
+ merge.checkAborted();
+
+ final Directory mergeDirectory = mergeScheduler.wrapForMerge(merge,
directory);
+ final IOContext context = IOContext.merge(merge.getStoreMergeInfo());
+
+ boolean success = false;
+ int totalDocs = 0;
+ // Outlive the loop that creates them: the postings of every output are
written between the
+ // phases either side of them, so each merger stays open across all three
phases.
+ final List<SegmentMerger> mergers = new ArrayList<>();
+ // The outputs share one set of input readers, so a reader that hands out
itself as its own
+ // merge instance must still be finished exactly once, whichever path gets
there first.
+ final Set<KnnVectorsReader> finishedVectorReaders =
+ Collections.newSetFromMap(new IdentityHashMap<>());
+ try {
+ initMergeReaders(merge, context);
+
+ // Resolved only now, so the policy can place boundaries on real key
values
+ // from the actual readers rather than guessing from doc counts. They are
+ // passed in rather than read back off the merge, so that a wrapping
merge
+ // can delegate this method.
+ final List<CodecReader> rawReaders = new ArrayList<>();
+ for (MergePolicy.MergeReader mr : merge.getMergeReader()) {
+ rawReaders.add(mr.reader);
+ }
+ final int[][] partitions = merge.getDocRangePartitions(rawReaders);
+ if (partitions == null) {
+ throw new IllegalStateException(
+ "OneMerge.isPartitioned() returned true but
getDocRangePartitions() returned null");
+ }
+ // Before indexing into it: a policy that returns the wrong shape should
be told so, rather
+ // than reaching this with an array index out of bounds.
+ if (partitions.length != merge.segments.size()) {
+ throw new IllegalArgumentException(
+ "docRangePartitions has "
+ + partitions.length
+ + " entries but the merge has "
+ + merge.segments.size()
+ + " input segments");
+ }
+ merge.outputCount = partitions[0].length - 1;
+ validateDocRangePartitions(merge, partitions, rawReaders);
+ final int outputCount = merge.getOutputCount();
+
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message(
+ "IW", "merging " + segString(merge.segments) + " into " +
outputCount + " outputs");
+ }
+
+ // Wrapped once and shared by every output. Wrapping is where a caller
filters a reader, and
+ // every output sees the same inputs, so doing it per output would
repeat that work k times.
+ // This is also where the inputs are verified: each format checksums the
files it is about to
+ // read when its merge begins, which costs a full read of them, and a
partitioned merge runs
+ // those merges once per output. Verifying here instead leaves one check
per input per merge,
+ // and does it before any output has written anything.
+ final List<CodecReader> wrappedReaders = new
ArrayList<>(rawReaders.size());
+ for (MergePolicy.MergeReader mergeReader : merge.getMergeReader()) {
+ merge.checkAborted();
+ final CodecReader wrapped = merge.wrapForMerge(mergeReader.reader);
+ validateMergeReader(wrapped);
+ wrapped.checkIntegrity(merge);
+ wrappedReaders.add(wrapped);
+ }
+ // Only now, so that the checks just above are the ones that read the
files.
+ merge.markInputsVerified();
+
+ final Executor intraMergeExecutor =
mergeScheduler.getIntraMergeExecutor(merge);
+ final List<MergeState.DocMap[]> docMapsPerOutput = new
ArrayList<>(outputCount);
+ merge.mergeStartNS = System.nanoTime();
+
+ final List<SegmentCommitInfo> outInfos = new ArrayList<>(outputCount);
+ final List<TrackingDirectoryWrapper> dirWrappers = new
ArrayList<>(outputCount);
+ // Phase A: merge each output.
+ for (int output = 0; output < outputCount; output++) {
+ merge.checkAborted();
+ final TrackingDirectoryWrapper dirWrapper = new
TrackingDirectoryWrapper(mergeDirectory);
+ final SegmentCommitInfo outInfo = newMergeSegmentInfo(merge);
+ merge.setMergeInfo(output, outInfo);
+
+ final List<CodecReader> mergeReaders = new ArrayList<>();
+ final Counter softDeleteCount = Counter.newCounter(false);
+ int i = 0;
+ for (MergePolicy.MergeReader mergeReader : merge.getMergeReader()) {
+ // Everything outside this output's range looks deleted, so it maps
to -1 in the
+ // resulting DocMap -- which is what routes concurrent deletes to
the right output.
+ CodecReader wrapped =
+ new DocRangeCodecReader(
+ wrappedReaders.get(i), partitions[i][output],
partitions[i][output + 1]);
+ if (softDeletesEnabled) {
+ // Count soft deletes that fall INSIDE this output's range. The
+ // single-output shortcut (softDelCount - numDeletedDocs) cannot be
+ // used here: numDeletedDocs on a range-restricted reader also
counts
+ // every document belonging to the other outputs. hardLiveDocs may
be
+ // null and countSoftDeletes handles that.
+ Counter hardDeleteCounter = Counter.newCounter(false);
Review Comment:
You are right, and it is a correctness bug rather than just a dead counter.
Fixed by excluding those documents once per input instead of per output, since
which ones are hard-deleted does not depend on the range, and the per-output
count is now an assertion that it sees none. The test took two goes: the
document has to be soft-deleted and hard-deleted, otherwise the count is always
zero and nothing shows. With overlapping deletes and the fix removed it fails
on softDeleteCount 8 != 16, so those documents were being counted as soft
deletes and would have been written out as live.
##########
lucene/misc/src/java/org/apache/lucene/misc/index/BalancedSegmentsMergePolicy.java:
##########
@@ -0,0 +1,231 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.lucene.misc.index;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import org.apache.lucene.index.CodecReader;
+import org.apache.lucene.index.FilterMergePolicy;
+import org.apache.lucene.index.MergePolicy;
+import org.apache.lucene.index.NumericDocValues;
+import org.apache.lucene.index.SegmentCommitInfo;
+import org.apache.lucene.index.SegmentInfos;
+import org.apache.lucene.index.TieredMergePolicy;
+import org.apache.lucene.util.Bits;
+
+/**
+ * Makes {@link org.apache.lucene.index.IndexWriter#forceMerge(int)} leave
segments holding about
+ * the same number of documents.
+ *
+ * <p>{@code forceMerge(n)} says how many segments to leave but nothing about
how the documents are
+ * shared between them, and in practice the share is uneven: merging a hundred
equal segments of a
+ * thousand documents into four leaves them holding 97,000, 1,000, 1,000 and
1,000. Since a search
+ * is parallelised by segment, a query then waits on a segment holding 97% of
the index. This policy
+ * leaves four segments of 25,000.
+ *
+ * <pre class="prettyprint">
+ * iwc.setMergePolicy(new BalancedSegmentsMergePolicy(new
TieredMergePolicy()));
+ * ...
+ * writer.forceMerge(numberOfSearchThreads);
+ * </pre>
+ *
+ * <p>The whole forced merge is planned at once: the segments are packed into
groups and each group
+ * is given the number of outputs its size deserves, so the work is several
independent merges the
+ * scheduler can run concurrently rather than one merge of everything. A group
that is already a
+ * single segment of about the right size is left alone, so a nearly balanced
index is balanced
+ * without being rewritten. The number of outputs is per merge rather than
global, so a merge may
+ * also have a single input: {@code forceMerge(16)} on a one-segment index
splits it into sixteen.
+ *
+ * <p>Only forced merges are affected; ordinary merging is left to the wrapped
policy. No index sort
+ * is needed, since documents are shared out by position rather than by any
key. A policy that
+ * placed its boundaries on a key instead would give each output a range of
that key, which this one
+ * does not attempt.
+ *
+ * @lucene.experimental
+ */
+public class BalancedSegmentsMergePolicy extends FilterMergePolicy {
+
+ /** Wraps {@link TieredMergePolicy}. */
+ public BalancedSegmentsMergePolicy() {
+ this(new TieredMergePolicy());
+ }
+
+ /**
+ * @param in the policy that decides ordinary merges; only forced merges are
changed
+ */
+ public BalancedSegmentsMergePolicy(MergePolicy in) {
+ super(in);
+ }
+
+ @Override
+ public MergeSpecification findForcedMerges(
+ SegmentInfos infos,
+ int maxSegmentCount,
+ Map<SegmentCommitInfo, Boolean> segmentsToMerge,
+ MergeContext context)
+ throws IOException {
+ if (maxSegmentCount == Integer.MAX_VALUE || maxSegmentCount < 1) {
+ // Not a request for a particular number of segments, so there is
nothing to balance.
+ return super.findForcedMerges(infos, maxSegmentCount, segmentsToMerge,
context);
+ }
+ final List<SegmentCommitInfo> eligible = new ArrayList<>();
+ for (SegmentCommitInfo info : infos) {
+ if (segmentsToMerge.containsKey(info) == false) {
+ continue;
+ }
+ if (context.getMergingSegments().contains(info)) {
+ // A plan is already running. Since a plan is made for the whole index
at once, planning
+ // again over what is left of it would give the outputs still being
written a second share
+ // of the segment count, so wait for it to finish instead.
+ return null;
+ }
+ eligible.add(info);
+ }
+ if (eligible.isEmpty()) {
+ return null;
+ }
+ long total = 0;
+ for (SegmentCommitInfo info : eligible) {
+ total += liveDocs(info);
+ }
+ if (total == 0) {
+ return super.findForcedMerges(infos, maxSegmentCount, segmentsToMerge,
context);
+ }
+ // Largest first, so a segment bigger than a share of its own takes as
many outputs as it needs
+ // and the small ones pack into what is left.
+
eligible.sort(Comparator.comparingLong(BalancedSegmentsMergePolicy::liveDocs).reversed());
+
+ final double share = (double) total / maxSegmentCount;
+ final MergeSpecification spec = new MergeSpecification();
+ final List<SegmentCommitInfo> group = new ArrayList<>();
+ long groupDocs = 0;
+ int assigned = 0;
+ for (int i = 0; i < eligible.size(); i++) {
+ group.add(eligible.get(i));
+ groupDocs += liveDocs(eligible.get(i));
+ final int left = eligible.size() - i - 1;
+ if (left > 0) {
+ final double withNext = groupDocs + liveDocs(eligible.get(i + 1));
+ if (Math.abs(withNext - share) < Math.abs(groupDocs - share)) {
+ continue; // taking the next segment too gets this group closer to a
share
+ }
+ }
+ // What this group has earned, less one output for every group still to
be formed. Each of
+ // those needs at least one segment, so at most `left` of them remain.
+ final int outputs =
+ left == 0
+ ? maxSegmentCount - assigned
+ : Math.min((int) Math.round(groupDocs / share), maxSegmentCount
- assigned - 1);
+ if (outputs < 1) {
+ continue; // no output left to give it; it joins the next group
+ }
+ if (group.size() > 1 || outputs > 1) {
+ // A lone segment that is already about the right size needs no merge
at all.
Review Comment:
Right, forceMerge(1) on a segment with deletes did nothing. Took your
suggestion as written. The policy is moving to a follow-up PR so this one stays
on the mechanism, the fix goes with it.
##########
lucene/misc/src/java/org/apache/lucene/misc/index/BalancedSegmentsMergePolicy.java:
##########
@@ -0,0 +1,231 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.lucene.misc.index;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import org.apache.lucene.index.CodecReader;
+import org.apache.lucene.index.FilterMergePolicy;
+import org.apache.lucene.index.MergePolicy;
+import org.apache.lucene.index.NumericDocValues;
+import org.apache.lucene.index.SegmentCommitInfo;
+import org.apache.lucene.index.SegmentInfos;
+import org.apache.lucene.index.TieredMergePolicy;
+import org.apache.lucene.util.Bits;
+
+/**
+ * Makes {@link org.apache.lucene.index.IndexWriter#forceMerge(int)} leave
segments holding about
+ * the same number of documents.
+ *
+ * <p>{@code forceMerge(n)} says how many segments to leave but nothing about
how the documents are
+ * shared between them, and in practice the share is uneven: merging a hundred
equal segments of a
+ * thousand documents into four leaves them holding 97,000, 1,000, 1,000 and
1,000. Since a search
+ * is parallelised by segment, a query then waits on a segment holding 97% of
the index. This policy
+ * leaves four segments of 25,000.
+ *
+ * <pre class="prettyprint">
+ * iwc.setMergePolicy(new BalancedSegmentsMergePolicy(new
TieredMergePolicy()));
+ * ...
+ * writer.forceMerge(numberOfSearchThreads);
+ * </pre>
+ *
+ * <p>The whole forced merge is planned at once: the segments are packed into
groups and each group
+ * is given the number of outputs its size deserves, so the work is several
independent merges the
+ * scheduler can run concurrently rather than one merge of everything. A group
that is already a
+ * single segment of about the right size is left alone, so a nearly balanced
index is balanced
+ * without being rewritten. The number of outputs is per merge rather than
global, so a merge may
+ * also have a single input: {@code forceMerge(16)} on a one-segment index
splits it into sixteen.
+ *
+ * <p>Only forced merges are affected; ordinary merging is left to the wrapped
policy. No index sort
+ * is needed, since documents are shared out by position rather than by any
key. A policy that
+ * placed its boundaries on a key instead would give each output a range of
that key, which this one
+ * does not attempt.
+ *
+ * @lucene.experimental
+ */
+public class BalancedSegmentsMergePolicy extends FilterMergePolicy {
+
+ /** Wraps {@link TieredMergePolicy}. */
+ public BalancedSegmentsMergePolicy() {
+ this(new TieredMergePolicy());
+ }
+
+ /**
+ * @param in the policy that decides ordinary merges; only forced merges are
changed
+ */
+ public BalancedSegmentsMergePolicy(MergePolicy in) {
+ super(in);
+ }
+
+ @Override
+ public MergeSpecification findForcedMerges(
+ SegmentInfos infos,
+ int maxSegmentCount,
+ Map<SegmentCommitInfo, Boolean> segmentsToMerge,
+ MergeContext context)
+ throws IOException {
+ if (maxSegmentCount == Integer.MAX_VALUE || maxSegmentCount < 1) {
+ // Not a request for a particular number of segments, so there is
nothing to balance.
+ return super.findForcedMerges(infos, maxSegmentCount, segmentsToMerge,
context);
+ }
+ final List<SegmentCommitInfo> eligible = new ArrayList<>();
+ for (SegmentCommitInfo info : infos) {
+ if (segmentsToMerge.containsKey(info) == false) {
+ continue;
+ }
+ if (context.getMergingSegments().contains(info)) {
+ // A plan is already running. Since a plan is made for the whole index
at once, planning
+ // again over what is left of it would give the outputs still being
written a second share
+ // of the segment count, so wait for it to finish instead.
+ return null;
+ }
+ eligible.add(info);
+ }
+ if (eligible.isEmpty()) {
+ return null;
+ }
+ long total = 0;
+ for (SegmentCommitInfo info : eligible) {
+ total += liveDocs(info);
+ }
+ if (total == 0) {
+ return super.findForcedMerges(infos, maxSegmentCount, segmentsToMerge,
context);
+ }
+ // Largest first, so a segment bigger than a share of its own takes as
many outputs as it needs
+ // and the small ones pack into what is left.
+
eligible.sort(Comparator.comparingLong(BalancedSegmentsMergePolicy::liveDocs).reversed());
+
+ final double share = (double) total / maxSegmentCount;
+ final MergeSpecification spec = new MergeSpecification();
+ final List<SegmentCommitInfo> group = new ArrayList<>();
+ long groupDocs = 0;
+ int assigned = 0;
+ for (int i = 0; i < eligible.size(); i++) {
+ group.add(eligible.get(i));
+ groupDocs += liveDocs(eligible.get(i));
+ final int left = eligible.size() - i - 1;
+ if (left > 0) {
+ final double withNext = groupDocs + liveDocs(eligible.get(i + 1));
+ if (Math.abs(withNext - share) < Math.abs(groupDocs - share)) {
+ continue; // taking the next segment too gets this group closer to a
share
+ }
+ }
+ // What this group has earned, less one output for every group still to
be formed. Each of
+ // those needs at least one segment, so at most `left` of them remain.
+ final int outputs =
+ left == 0
+ ? maxSegmentCount - assigned
+ : Math.min((int) Math.round(groupDocs / share), maxSegmentCount
- assigned - 1);
+ if (outputs < 1) {
+ continue; // no output left to give it; it joins the next group
+ }
+ if (group.size() > 1 || outputs > 1) {
+ // A lone segment that is already about the right size needs no merge
at all.
+ spec.add(new Split(new ArrayList<>(group), outputs));
+ }
+ assigned += outputs;
+ group.clear();
+ groupDocs = 0;
+ }
+ return spec.merges.isEmpty() ? null : spec;
+ }
+
+ private static long liveDocs(SegmentCommitInfo info) {
+ return info.info.maxDoc() - info.getDelCount();
Review Comment:
Agreed, it should balance on what the merge will actually write. Changed,
and moving to the follow-up with the rest of the policy.
##########
lucene/core/src/java/org/apache/lucene/index/IndexWriter.java:
##########
@@ -5214,11 +5254,463 @@ private boolean assertSoftDeletesCount(CodecReader
reader, int expectedCount) th
return true;
}
+ /**
+ * Turns one finished merge output into a segment on disk: packs its files
into a compound file if
+ * the codec asks for one, writes its segment info, and warms it.
+ *
+ * <p>Shared by the single-output and the partitioned merge paths, which
differ in how many
+ * outputs they produce and not in what packaging one of them means. The
delete-on-failure and
+ * abort handling here is the delicate part, and is the reason this is one
method rather than two
+ * similar ones.
+ *
+ * @return false if the merge was aborted while this ran, in which case the
caller must abandon it
+ */
+ private boolean packageMergedSegment(
+ MergePolicy.OneMerge merge,
+ MergePolicy mergePolicy,
+ SegmentCommitInfo info,
+ IOContext context)
+ throws IOException {
+ // Very important to do this before opening the reader
+ // because codec must know if prox was written for
+ // this segment:
+ boolean useCompoundFile;
+ synchronized (this) { // Guard segmentInfos
+ useCompoundFile =
+ info.info
+ .getCodec()
+ .compoundFormat()
+ .useCompoundFile(mergePolicy.size(info, this), mergePolicy);
+ }
+
+ if (useCompoundFile) {
+ Collection<String> filesToRemove = info.files();
+ // NOTE: Creation of the CFS file must be performed with the original
+ // directory rather than with the merging directory, so that it is not
+ // subject to merge throttling.
+ TrackingDirectoryWrapper trackingCFSDir = new
TrackingDirectoryWrapper(directory);
+ try {
+ createCompoundFile(infoStream, trackingCFSDir, info.info, context,
this::deleteNewFiles);
+ } catch (Throwable t) {
+ try {
+ synchronized (this) {
+ if (merge.isAborted()) {
+ // This can happen if rollback is called while we were building
+ // our CFS -- fall through to logic below to remove the non-CFS
+ // merged files:
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message(
+ "IW", "hit merge abort exception creating compound file
during merge: " + t);
+ }
+ return false;
+ } else {
+ handleMergeException(t, merge);
+ }
+ }
+ } finally {
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message("IW", "hit exception creating compound file
during merge: " + t);
+ }
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ }
+ }
+
+ synchronized (this) {
+
+ // delete new non cfs files directly: they were never
+ // registered with IFD
+ deleteNewFiles(filesToRemove);
+
+ if (merge.isAborted()) {
+ if (infoStream.isEnabled("IW")) {
+ infoStream.message("IW", "abort merge after building CFS");
+ }
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ return false;
+ }
+ }
+
+ info.info.setUseCompoundFile(true);
+ }
+
+ // Have codec write SegmentInfo. Must do this after
+ // creating CFS so that 1) .si isn't slurped into CFS,
+ // and 2) .si reflects useCompoundFile=true change
+ // above:
+ try {
+ config.getCodec().segmentInfoFormat().write(directory, info.info,
context);
+ } catch (Throwable t) {
+ // Safe: these files must exist
+ deleteNewFiles(info.files());
+ throw t;
+ }
+
+ // TODO: ideally we would freeze info here!!
+ // because any changes after writing the .si will be
+ // lost...
+
+ final IndexReaderWarmer mergedSegmentWarmer =
config.getMergedSegmentWarmer();
+ if (readerPool.isReaderPoolingEnabled() && mergedSegmentWarmer != null) {
+ final ReadersAndUpdates rld = getPooledInstance(info, true);
+ final SegmentReader sr = rld.getReader(IOContext.DEFAULT);
+ try {
+ mergedSegmentWarmer.warm(sr);
+ } finally {
+ synchronized (this) {
+ rld.release(sr);
+ release(rld);
+ }
+ }
+ }
+ return true;
+ }
+
+ /** Opens a pooled reader for each of a merge's input segments, and holds
their files open. */
+ private void initMergeReaders(MergePolicy.OneMerge merge, IOContext context)
throws IOException {
+ merge.initMergeReaders(
+ sci -> {
+ final ReadersAndUpdates rld = getPooledInstance(sci, true);
+ rld.setIsMerging();
+ synchronized (this) {
+ return rld.getReaderForMerge(
+ context, mr ->
deleter.incRef(mr.reader.getSegmentInfo().files()));
+ }
+ });
+ }
+
+ /** Fresh {@link SegmentCommitInfo} for one output of a merge. */
+ private SegmentCommitInfo newMergeSegmentInfo(MergePolicy.OneMerge merge) {
+ boolean hasBlocks = false;
+ for (SegmentCommitInfo info : merge.segments) {
+ if (info.info.getHasBlocks()) {
+ hasBlocks = true;
+ break;
+ }
+ }
+ SegmentInfo si =
+ new SegmentInfo(
+ directoryOrig,
+ Version.LATEST,
+ null,
+ newSegmentName(),
+ -1,
+ false,
+ hasBlocks,
+ config.getCodec(),
+ Collections.emptyMap(),
+ StringHelper.randomId(),
+ Collections.emptyMap(),
+ config.getIndexSort());
+ Map<String, String> details = new HashMap<>();
+ details.put("mergeMaxNumSegments", "" + merge.maxNumSegments);
+ details.put("mergeFactor", Integer.toString(merge.segments.size()));
+ details.put("mergeOutputs", Integer.toString(merge.getOutputCount()));
+ setDiagnostics(si, SOURCE_MERGE, details);
+ return new SegmentCommitInfo(si, 0, 0, -1L, -1L, -1L,
StringHelper.randomId());
+ }
+
+ /**
+ * Validate a doc-range partition spec. Only the SHAPE is checked -- one
boundary array per input,
+ * all the same length, non-decreasing, starting at 0 and ending at that
input's maxDoc. Given
+ * that shape, the two properties the merge depends on, disjointness and
full coverage, hold by
+ * construction rather than by trusting the caller.
+ *
+ * <p>Plus the one thing that is not about the spec: an index sort.
Correctness rests on the
+ * outputs being contiguous AND in order in the merged document space, which
holds under a sort,
+ * where the merged order is key order and the outputs are key ranges.
Without one, {@link
+ * DocIDMerger} concatenates input by input, so an output's documents land
in one block per input
+ * rather than in a single run, and splitting a term's postings by document
id would quietly hand
+ * documents to the wrong output.
+ */
Review Comment:
It is not required, that paragraph is mine and it outlived the check I
removed. Deleted. testPartitionsAnUnsortedIndex covers it.
##########
lucene/core/src/java/org/apache/lucene/index/MergeState.java:
##########
@@ -287,6 +287,36 @@ static PackedLongValues removeDeletes(final int maxDoc,
final Bits liveDocs) {
return docMapBuilder.build();
}
+ /**
+ * A copy of {@code other} that maps documents differently: same inputs,
same fields, same
+ * everything a format reads, but a different view of where each document
ends up.
+ *
+ * <p>This exists for a merge writing several outputs from one pass, which
needs the same readers
+ * addressed in the whole merged document space rather than in one output's.
Copying keeps that
+ * caller out of the constructor below, whose argument list is long enough
that a caller passing
+ * the fields through by hand is one reordering away from a silent bug.
+ */
+ MergeState(MergeState other, DocMap[] docMaps, boolean needsIndexSort) {
Review Comment:
Left over from the postings push path that came out of this PR. Removed.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]