Commit f853f1af by chenyuanjie

spark消费停止前,把kafka已有数据全部消费完再停止

parent f0a3b8c7
...@@ -977,7 +977,7 @@ class KafkaFlowAsinDetail(Templates): ...@@ -977,7 +977,7 @@ class KafkaFlowAsinDetail(Templates):
self.df_self_brand = F.broadcast(SparkUtil.read_jdbc_query( self.df_self_brand = F.broadcast(SparkUtil.read_jdbc_query(
session=self.spark, url=mysql_con['url'], pwd=mysql_con['pwd'], session=self.spark, url=mysql_con['url'], pwd=mysql_con['pwd'],
username=mysql_con['username'], query=mysql_brand_sql username=mysql_con['username'], query=mysql_brand_sql
)) ).persist(StorageLevel.MEMORY_ONLY))
self.df_self_brand.show(10, truncate=False) self.df_self_brand.show(10, truncate=False)
print("10. 读取隐藏分类(流量选品模块,id_path前缀匹配),用于asin_type计算") print("10. 读取隐藏分类(流量选品模块,id_path前缀匹配),用于asin_type计算")
# category_full_name + category_disable_config 按 id_path 前缀匹配(对齐 dwt.handle_asin_is_hide) # category_full_name + category_disable_config 按 id_path 前缀匹配(对齐 dwt.handle_asin_is_hide)
...@@ -992,7 +992,7 @@ class KafkaFlowAsinDetail(Templates): ...@@ -992,7 +992,7 @@ class KafkaFlowAsinDetail(Templates):
self.df_hide_category = F.broadcast(SparkUtil.read_jdbc_query( self.df_hide_category = F.broadcast(SparkUtil.read_jdbc_query(
session=self.spark, url=mysql_con['url'], pwd=mysql_con['pwd'], session=self.spark, url=mysql_con['url'], pwd=mysql_con['pwd'],
username=mysql_con['username'], query=sql username=mysql_con['username'], query=sql
)) ).persist(StorageLevel.MEMORY_ONLY))
self.df_hide_category.show(10, truncate=False) self.df_hide_category.show(10, truncate=False)
# 字段处理逻辑综合 # 字段处理逻辑综合
......
...@@ -348,6 +348,9 @@ class Templates(object): ...@@ -348,6 +348,9 @@ class Templates(object):
def _do_stop(): def _do_stop():
try: try:
if self.query is not None: if self.query is not None:
# 先把调用时刻Kafka里已经存在但还没被消费的积压全部跑完,再停,
# 避免爬虫刚标记完成、但还有一截尾巴数据没被下一次触发扫到就被stop掉
self.query.processAllAvailable()
self.query.stop() # 在子线程中调用,避免 foreachBatch 回调内死锁 self.query.stop() # 在子线程中调用,避免 foreachBatch 回调内死锁
if self.spark is not None: if self.spark is not None:
self.spark.stop() self.spark.stop()
......
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