Skip to content

Commit 3f2c913

Browse files
Merge upstream master
1 parent 91c0ba7 commit 3f2c913

1 file changed

Lines changed: 2 additions & 1 deletion

File tree

streamz/sources.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -515,7 +515,8 @@ def checkpoint_emit(_part):
515515
continue
516516
if 'auto.offset.reset' in self.consumer_params.keys():
517517
if self.consumer_params['auto.offset.reset'] == 'latest' and \
518-
self.positions == [-1001] * self.npartitions:
518+
(self.positions == [-1001] * self.npartitions
519+
or self.positions == [0] * self.npartitions):
519520
self.positions[partition] = high
520521
current_position = self.positions[partition]
521522
lowest = max(current_position, low)

0 commit comments

Comments
 (0)