summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--kafka/consumer/fetcher.py5
1 files changed, 5 insertions, 0 deletions
diff --git a/kafka/consumer/fetcher.py b/kafka/consumer/fetcher.py
index c1f98eb..375090a 100644
--- a/kafka/consumer/fetcher.py
+++ b/kafka/consumer/fetcher.py
@@ -338,6 +338,8 @@ class Fetcher(six.Iterator):
for record in self._unpack_message_set(tp, messages):
# Fetched compressed messages may include additional records
if record.offset < fetch_offset:
+ log.debug("Skipping message offset: %s (expecting %s)",
+ record.offset, fetch_offset)
continue
drained[tp].append(record)
else:
@@ -419,6 +421,9 @@ class Fetcher(six.Iterator):
# Compressed messagesets may include earlier messages
# It is also possible that the user called seek()
elif msg.offset != self._subscriptions.assignment[tp].position:
+ log.debug("Skipping message offset: %s (expecting %s)",
+ msg.offset,
+ self._subscriptions.assignment[tp].position)
continue
self._subscriptions.assignment[tp].position = msg.offset + 1