Python调用Firebase Firestore onSnapshot监听集合出现重复回调问题求解
问题修复方案
核心原因
- 你在
while True循环中重复调用on_snapshot()注册监听器,旧监听器不会自动销毁,运行时间越长注册的监听器数量越多,同一份集合变更会触发所有已注册的监听器回调,直接导致打印内容重复、Generating...执行多次。 on_snapshot的首次回调是异步执行的,所以第一轮循环打印完while 0后,回调还没触发就进入了下一轮循环,才会出现while 1打印在首次用户输出前面的问题。
修复方案
方案1:改用定时轮询逻辑(最符合预期输出)
你当前的逻辑更接近定时拉取数据处理的场景,不需要使用Firestore的实时监听能力,直接每次循环调用get()拉取最新数据即可,完全避免监听器重复问题:
简化测试代码修复版
import time # 此处省略Firebase初始化代码 while_counter=0 while True: print('while %s ----------' %(while_counter)) while_counter += 1 time.sleep(1) # 单次拉取当前用户集合的所有文档 doc_snapshot = db.collection(u'users').get() user_counter = 0 for doc in doc_snapshot: print('user %s' %(user_counter)) user_counter += 1
业务代码修复版
import time # 此处省略Firebase初始化代码 # 记录已处理的子文档ID,防止边界情况下重复处理 processed_subdoc_ids = set() while_counter=0 while True: print('while_counter: %s' %(while_counter)) while_counter += 1 time.sleep(1) # 拉取所有用户文档 user_docs = db.collection(u'users').get() user_counter = 0 for doc in user_docs: print('user_counter: %s' %(user_counter)) user_counter += 1 # 拉取当前用户下flag为True的历史记录 sub_docs = db.collection(u'users').document(doc.id).collection(u'history').where(u'flag', u'==', True).stream() subdoc_counter = 0 for sub_doc in sub_docs: # 已处理过的文档直接跳过 if sub_doc.id in processed_subdoc_ids: continue print('subdoc_counter: %s' %(subdoc_counter)) subdoc_counter += 1 print('Generating...') sub_doc_dict = sub_doc.to_dict() usr_txt = sub_doc_dict['inputdata'] # 此处保留你原来的业务处理逻辑 # 更新flag为False db.collection(u'users').document(doc.id).collection(u'history').document(sub_doc.id).set({ u'outputdata': usr_txt, u'flag': False }, merge=True) # 标记为已处理 processed_subdoc_ids.add(sub_doc.id)
方案2:保留实时监听能力
如果你需要使用Firestore的实时变更推送能力,只需要把监听器注册逻辑移到循环外面,全局仅注册一次即可:
import time # 此处省略Firebase初始化代码 processed_subdoc_ids = set() def on_snapshot(doc_snapshot, changes, read_time): # 仅处理新增/修改的变更,减少不必要的重复计算 for change in changes: if change.type.name not in ('ADDED', 'MODIFIED'): continue doc = change.document # 以下逻辑和原业务逻辑一致 sub_docs = db.collection(u'users').document(doc.id).collection(u'history').where(u'flag', u'==', True).stream() subdoc_counter = 0 for sub_doc in sub_docs: if sub_doc.id in processed_subdoc_ids: continue print('subdoc_counter: %s' %(subdoc_counter)) subdoc_counter += 1 print('Generating...') sub_doc_dict = sub_doc.to_dict() usr_txt = sub_doc_dict['inputdata'] # 此处保留你原来的业务处理逻辑 db.collection(u'users').document(doc.id).collection(u'history').document(sub_doc.id).set({ u'outputdata': usr_txt, u'flag': False }, merge=True) processed_subdoc_ids.add(sub_doc.id) # 全局仅注册一次监听器 doc_ref = db.collection(u'users') doc_watch = doc_ref.on_snapshot(on_snapshot) # 循环仅负责打印计数 while_counter=0 while True: print('while_counter: %s' %(while_counter)) while_counter += 1 time.sleep(1)
内容的提问来源于stack exchange,提问作者fhashikawa
相关产品推荐
相关产品推荐

