mjsax commented on a change in pull request #10407:
URL: https://github.com/apache/kafka/pull/10407#discussion_r602605779
##########
File path:
streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java
##########
@@ -509,38 +519,60 @@ void handleRevocation(final Collection<TopicPartition>
revokedPartitions) {
prepareCommitAndAddOffsetsToMap(commitNeededActiveTasks,
consumedOffsetsPerTask);
}
- // even if commit failed, we should still continue and complete
suspending those tasks,
- // so we would capture any exception and throw
+ // even if commit failed, we should still continue and complete
suspending those tasks, so we would capture
+ // any exception and rethrow it at the end
+ final Set<TaskId> corruptedTasks = new HashSet<>();
try {
commitOffsetsOrTransaction(consumedOffsetsPerTask);
+ } catch (final TaskCorruptedException e) {
+ log.warn("Some tasks were corrupted when trying to commit offsets,
these will be cleaned and revived: {}",
+ e.corruptedTasks());
+
+ // If we hit a TaskCorruptedException, we should just handle the
cleanup for those corrupted tasks right here
+ corruptedTasks.addAll(e.corruptedTasks());
+ final Map<Task, Collection<TopicPartition>>
corruptedTasksWithChangelogs = new HashMap<>();
+ for (final TaskId taskId : corruptedTasks) {
+ final Task task = tasks.task(taskId);
+ task.markChangelogAsCorrupted(task.changelogPartitions());
+ corruptedTasksWithChangelogs.put(task,
task.changelogPartitions());
+ }
+ closeAndRevive(corruptedTasksWithChangelogs);
} catch (final RuntimeException e) {
log.error("Exception caught while committing those revoked tasks "
+ revokedActiveTasks, e);
+
+ // TODO: KIP-572 need to handle TimeoutException, may be rethrown
from committing offsets under ALOS
Review comment:
I don't think we need to handle `TimeoutException` here? If a client
throws a `Timeout` it should be handled within `commitOffsetsOrTransaction` --
and `commitOffsetsOrTransaction` should only re-throw a timeout after
`task.timeout.ms` expired and for this case we should indeed let the thread die.
Or actually, it seems we don't want `commitOffsetsOrTransaction` to catch
client timeouts when we call it from here (we only want it to catch them during
regular processing). Thus, we might need to pass in a boolean flag... And for
this call, any client `TimeoutException` should be rethrown as
`TaskCorruptedException` ?
--
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.
For queries about this service, please contact Infrastructure at:
[email protected]