[ 
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]

Reply via email to