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]