Commit 1425f4b3 by chenyuanjie

fix

parent 964d93bf
...@@ -401,6 +401,7 @@ class Templates(object): ...@@ -401,6 +401,7 @@ class Templates(object):
partition_data = self.get_kafka_partitions_data(consumer=consumer, topic_name=self.topic_name) partition_data = self.get_kafka_partitions_data(consumer=consumer, topic_name=self.topic_name)
consumer.close() consumer.close()
# 覆盖topic全部分区(无遗漏分区start=end)
current_offsets, target_end_offsets = {}, {} current_offsets, target_end_offsets = {}, {}
gap_total = 0 gap_total = 0
for pid, info in partition_data.items(): for pid, info in partition_data.items():
...@@ -408,9 +409,9 @@ class Templates(object): ...@@ -408,9 +409,9 @@ class Templates(object):
begin_offset = info['beginning_offsets'] begin_offset = info['beginning_offsets']
if committed is not None and str(pid) in committed: if committed is not None and str(pid) in committed:
begin_offset = max(begin_offset, int(committed[str(pid)])) begin_offset = max(begin_offset, int(committed[str(pid)]))
if end_offset > begin_offset:
current_offsets[pid] = begin_offset current_offsets[pid] = begin_offset
target_end_offsets[pid] = end_offset target_end_offsets[pid] = end_offset
if end_offset > begin_offset:
gap_total += end_offset - begin_offset gap_total += end_offset - begin_offset
if gap_total == 0: if gap_total == 0:
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment