[
https://issues.apache.org/jira/browse/SOLR-18408?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Andrzej Bialecki updated SOLR-18408:
------------------------------------
Description:
{{PartitionManager.getOffsetForPartition}} returns
{{partitionRecords.get(partitionRecords.size() - 1).offset() + 1}}, which is
the last offset of the _entire_ poll batch for that partition. Every
{{WorkUnit}} created while iterating that batch therefore receives the same
{{nextOffset}}. The first WorkUnit to complete commits an offset past all the
later WorkUnits that are still in flight.
Example failure scenario: a poll returns records at offsets 100–199 for one
partition and the collapse logic splits them into three WorkUnits (e.g. because
some records contained only additions and some others contained also deletes,
which enforces a buffer flush + new WorkUnit). All three units now carry
{{{}nextOffset = 200{}}}. The first completes and commits 200. The consumer
crashes while the other two are still executing. On restart the consumer
resumes at 200 and the records in the two unfinished WorkUnits are never
replayed. Data is lost.
Additionally, {{PartitionManager.checkForOffsetUpdates}} synchronizes on the
method argument {{TopicPartition}} which is an "equal-but-distinct" object
returned by Kafka from different threads, so the synchronization doesn't work.
was:
{{PartitionManager.getOffsetForPartition}} returns
{{partitionRecords.get(partitionRecords.size() - 1).offset() + 1 }}which is the
last offset of the _entire_ poll batch for that partition. Every {{WorkUnit}}
created while iterating that batch therefore receives the same
{{{}nextOffset{}}}. The first WorkUnit to complete commits an offset past all
the later WorkUnits that are still in flight.
Example failure scenario: a poll returns records at offsets 100–199 for one
partition and the collapse logic splits them into three WorkUnits (e.g. because
some records contained only additions and some others contained also deletes,
which enforces a buffer flush + new WorkUnit). All three units now carry
{{{}nextOffset = 200{}}}. The first completes and commits 200. The consumer
crashes while the other two are still executing. On restart the consumer
resumes at 200 and the records in the two unfinished WorkUnits are never
replayed. Data is lost.
Additionally, {{PartitionManager.checkForOffsetUpdates}} synchronizes on the
method argument {{TopicPartition}} which is an "equal-but-distinct" object
returned by Kafka from different threads, so the synchronization doesn't work.
> CrossDC Consumer: incorrect nextOffset commits
> ----------------------------------------------
>
> Key: SOLR-18408
> URL: https://issues.apache.org/jira/browse/SOLR-18408
> Project: Solr
> Issue Type: Bug
> Components: module - crossDC
> Affects Versions: 10.0, 9.10.1
> Reporter: Andrzej Bialecki
> Assignee: Andrzej Bialecki
> Priority: Major
>
> {{PartitionManager.getOffsetForPartition}} returns
> {{partitionRecords.get(partitionRecords.size() - 1).offset() + 1}}, which
> is the last offset of the _entire_ poll batch for that partition. Every
> {{WorkUnit}} created while iterating that batch therefore receives the same
> {{nextOffset}}. The first WorkUnit to complete commits an offset past all the
> later WorkUnits that are still in flight.
> Example failure scenario: a poll returns records at offsets 100–199 for one
> partition and the collapse logic splits them into three WorkUnits (e.g.
> because some records contained only additions and some others contained also
> deletes, which enforces a buffer flush + new WorkUnit). All three units now
> carry {{{}nextOffset = 200{}}}. The first completes and commits 200. The
> consumer crashes while the other two are still executing. On restart the
> consumer resumes at 200 and the records in the two unfinished WorkUnits are
> never replayed. Data is lost.
> Additionally, {{PartitionManager.checkForOffsetUpdates}} synchronizes on the
> method argument {{TopicPartition}} which is an "equal-but-distinct" object
> returned by Kafka from different threads, so the synchronization doesn't work.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]