Lucas Kot-Zaniewski created SOLR-18507:
------------------------------------------

             Summary: Make Index Fingerprint Faster in Leader Election Path
                 Key: SOLR-18507
                 URL: https://issues.apache.org/jira/browse/SOLR-18507
             Project: Solr
          Issue Type: Wish
            Reporter: Lucas Kot-Zaniewski


Leader election causes the new candidate to run PeerSync against all replicas. 
The reason is it must ensure all replicas under its stewardship are "document 
equivalent" before it can safely pronounce its leadership. Document equivalence 
is ascertained with the help of a checksum called the index fingerprint. The 
fingerprint is computed over all the document versions in the index and is O(# 
of docs in the index) in the worst case. Yes, there is a patchwork of caching 
and parallelization to speed up the average case. However, our alarming informs 
us about the \~worst case often enough for us to consider it a major pain 
point. The issue is especially noticeable when you have a lot of docs relative 
to your allocated compute. For instance, if you have low billions of documents 
on a node with ~16 cores, this computation can block leader elections for tens 
of seconds (at least in our environment). This amounts to a write-outage that 
makes Solr look bad.

The slowness stems ~99% (literally) from decoding Solr's \_version\_ field. 
Reading doc-values at this scale is slow; it requires increment doc Id iterator 
and, more expensively, decompressing the long value from the DV long-point 
format billions of times. Note this is CPU time and not waiting for mmap to 
load into resident memory. We run on clusters with ample RAM but I imagine this 
problem may be even worse when you run more memory constrained processes.

I have thought about this problem for a bit and can volunteer some general 
directions (obviously not exhaustive) where improvements could be made.

1. Keep existing in-memory cache but compute it more eagerly, i.e. on start-up 
and after a full recovery. I find these scenarios cause cache invalidation the 
most often. The downside is you make initial node-start-up a function of index 
size. You have to consider deletions also invalidate the cache so you'd want to 
subtract deletions here in order to avoid total re-computation after 
delete-heavy workloads. Which brings me to my next suggestion.

2. Persisting a gross fingerprint during segment creation and then subtracting 
deletions during fingerprint calculation. Deletions are easily invertible with 
regards to the fingerprint. You can compute a gross fingerprint when you create 
the segment and store it, perhaps along with commit info or in the segment 
attributes themselves. Then when you need the fingerprint you subtract the 
checksum of the deleted docs from this gross checksum. This allows you to scan 
only the deleted docs. I've ran a PoC of this kind of solution and find that it 
results in a ~7X improvement when the deletion ratio is ~10%. The downside is 
your fingerprint computation cost is now a function of your deletion ratio 
which, although much better than a full index scan, is still typically O(# of 
docs in the index).

3. Carrying the idea even further, you could incrementally update the 
fingerprint with each generation of deletes as you index. So you can store not 
only the gross-fingerprint but the net fingerprint as well. This net 
fingerprint could track the generation of your live docs. I have a PoC of this 
kind of solution as well and find it reduces the peer sync cost dramatically 
(~100X). The downsides are you need a custom codec and you pay a bigger tax 
during indexing than 2. My profiling found this tax to still be really 
miniscule (less than ~1% of indexing CPU samples) but it exists nonetheless. 
Also, it adds a bit of complexity, most prominently in the form of a custom 
codec. The PoC also has no optimizations for in-place DV updates which would 
require even _more_ complexity to optimize (did I mention dislike in-place dv 
updates?).

4. Disable fingerprint computation for tlog+pull replicas. Tlog+pull replicas 
don't really benefit from tracking document-equivalence via a checksum over the 
individual document versions. They converge to the same underlying segments so 
I think this is totally wasted work (unless we do something along the lines of 
the custom codec proposed in points 2 or 3. There you could see how a leader 
tlog replica could keep track of fingerprint just in case a benefitting NRT 
replica ever gets created in the future though this seems like a rather small 
corner case).

5. (more of an architectural shift) scrap the fingerprint altogether. The only 
reason it is needed is because Solr's "update identifier" of choice (the 
document \_version\_) is sparse. It's effectively a Lamport timestamp. But if 
it were a dense sequence number instead, you could track the high watermark of 
the cluster via a simple checkpoint. Elastic and open search do this for their 
"document replication" mode. Each replica tracks a checkpoint which is the 
point up to which it knows the full history "densely". The minimum checkpoint 
of all active replicas is a high watermark, below which you don't need to check 
_anything_.

I welcome any other discussions or suggestions you may have regarding this 
ticket!



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to