raghavan-arvind opened a new issue, #616:
URL: https://github.com/apache/pekko-connectors-kafka/issues/616

   Hi team, we noticed one issue when using CooperateStickyAssignor strategy 
with pekko connectors. Our understanding of the issue is as follows -- 
   
   Here is the rebalance handler:
   ```
         private var lastRevoked = Set.empty[TopicPartition]
   
         override def onRevoke(revokedTps: Set[TopicPartition], consumer: 
RestrictedConsumer): Unit =
           lastRevoked = revokedTps
   
         override def onAssign(assignedTps: Set[TopicPartition], consumer: 
RestrictedConsumer): Unit =
           for {
             tp <- lastRevoked -- assignedTps
             control <- subSources.get(tp)
           } control.filterRevokedPartitionsCB.invoke(Set(tp))
   ```
   
   and the callback for `filterRevokedPartitionsCB` which drops messages we 
have already polled for revoked topics:
   
   ```
     private def filterRevokedPartitions(topicPartitions: Set[TopicPartition]): 
Unit = {
       if (topicPartitions.nonEmpty) {
         log.debug("filtering out messages from revoked partitions {}", 
topicPartitions)
         // as buffer is an Iterator the filtering will be applied during `pump`
         buffer = buffer.filterNot { record =>
           val tp = new TopicPartition(record.topic, record.partition)
           topicPartitions.contains(tp)
         }
       }
     }
   ```
   
   With non-sticky rebalances, we always get:
   
   * `onRevoke` - revoke all partitions
   * `onAssign` - assign new partitions
   
   
   a full "stop-the-world" rebalance. In that case, the code is correct, 
`onRevoke` resets mutable variable `lastRevoked`, then the next `onAssign` 
assigns new partitions and drops any messages related to old partitions.
   
   But in cooperative sticky rebalancing, consumers will opt towards keeping 
their partitions, and it looks like we aren't getting an empty onRevoke 
callback if the consumer keeps its partitions, we just get:
   
   * `onAssign` - clears whatever the last revoked partition was, adds new 
topics
   
   
   Every successive `onAssign` callback, will drop all already-polled messages 
for the last revoked topics , since we never got a new callback for `onRevoke` 
so that mutable var never cleared. This causes us to drop some messages. We 
poll the next set of messages and commit afterwards, so we end up silently 
dropping messages.
   
   
   Would like to confirm whether this makes sense / is the 
`CooperativeStickyAssignor` generally supported or something we would like to 
add? This is something we would be potentially willing to implement support for 
if this is the direction you all would like to go in


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to