Commit 22427ae9 by chenyuanjie

kafka补偿机制-如果是爬虫中断不触发补偿

parent d32523eb
......@@ -338,7 +338,7 @@ class Templates(object):
def start_process_instance(self):
pass
def kafka_stream_stop(self):
def kafka_stream_stop(self, compensate=True):
if self._stopping:
return
self._stopping = True
......@@ -357,7 +357,11 @@ class Templates(object):
wait_seconds += 1
if self.query.isActive:
print(f"[停止] query.stop()后等待{wait_seconds}s仍处于active状态,继续执行后续流程")
if compensate:
self._compensate_missed_offsets() # 停止前核对topic末尾offset,若有遗漏区间则补偿消费
else:
# 爬虫中断场景(如让位优先抓日asin):数据量可能过大,且爬虫最终正常结束时会再触发一次补偿兜底,这里直接快速停止让出资源
print(f"[停止] 爬虫中断场景,跳过停止前补偿消费,优先快速让出资源")
if self.spark is not None:
self.spark.stop()
except Exception as e:
......@@ -520,7 +524,7 @@ class Templates(object):
else:
time.sleep(10)
if self.consumer_type == "latest":
self.kafka_stream_stop()
self.kafka_stream_stop(compensate=False)
else:
# 爬虫还在进行中
print(f"爬虫'{self.spider_type}'还在爬取中(spider_state={spider_state}, spider_is_ready={spider_is_ready}), 继续消费")
......
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