[
https://issues.apache.org/jira/browse/SOLR-18507?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Lucas Kot-Zaniewski updated SOLR-18507:
---------------------------------------
Description:
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
the version long-values of 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!
was:
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!
> 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
> Priority: Major
>
> 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 the version long-values of 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]