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]