@@ -45,6 +45,9 @@ def __init__(self, redis_key, task_table=None):
4545
4646 self ._table_item = setting .TAB_ITEM
4747 self ._table_request = setting .TAB_REQUSETS .format (redis_key = redis_key )
48+ self ._table_failed_items = setting .TAB_FAILED_ITEMS .format (
49+ redis_key = redis_key
50+ )
4851
4952 self ._item_tables = {
5053 # 'xxx_item': {'tab_item': 'xxx:xxx_item'} # 记录item名与redis中item名对应关系
@@ -246,39 +249,37 @@ def __pick_items(self, items, is_update_item=False):
246249
247250 return datas_dict
248251
249- def __export_to_db (self , tab_item , datas , is_update = False , update_keys = ()):
250- to_table = tools .get_info (tab_item , ":s_(.*?)_item$" , fetch_one = True )
251-
252+ def __export_to_db (self , table , datas , is_update = False , update_keys = ()):
252253 # 打点 校验
253- self .check_datas (table = to_table , datas = datas )
254+ self .check_datas (table = table , datas = datas )
254255
255256 for pipeline in self ._pipelines :
256257 if is_update :
257- if to_table == self ._task_table and not isinstance (
258+ if table == self ._task_table and not isinstance (
258259 pipeline , MysqlPipeline
259260 ):
260261 continue
261262
262- if not pipeline .update_items (to_table , datas , update_keys = update_keys ):
263+ if not pipeline .update_items (table , datas , update_keys = update_keys ):
263264 log .error (
264- f"{ pipeline .__class__ .__name__ } 更新数据失败. table: { to_table } items: { datas } "
265+ f"{ pipeline .__class__ .__name__ } 更新数据失败. table: { table } items: { datas } "
265266 )
266267 return False
267268
268269 else :
269- if not pipeline .save_items (to_table , datas ):
270+ if not pipeline .save_items (table , datas ):
270271 log .error (
271- f"{ pipeline .__class__ .__name__ } 保存数据失败. table: { to_table } items: { datas } "
272+ f"{ pipeline .__class__ .__name__ } 保存数据失败. table: { table } items: { datas } "
272273 )
273274 return False
274275
275276 # 若是任务表, 且上面的pipeline里没mysql,则需调用mysql更新任务
276- if not self ._have_mysql_pipeline and is_update and to_table == self ._task_table :
277+ if not self ._have_mysql_pipeline and is_update and table == self ._task_table :
277278 if not self .mysql_pipeline .update_items (
278- to_table , datas , update_keys = update_keys
279+ table , datas , update_keys = update_keys
279280 ):
280281 log .error (
281- f"{ pipeline . __class__ .__name__ } 更新数据失败. table: { to_table } items: { datas } "
282+ f"{ self . mysql_pipeline . __class__ .__name__ } 更新数据失败. table: { table } items: { datas } "
282283 )
283284 return False
284285
@@ -299,36 +300,46 @@ def __add_item_to_db(
299300 update_items_dict = self .__pick_items (update_items , is_update_item = True )
300301
301302 # item批量入库
303+ failed_items = {}
302304 while items_dict :
303305 tab_item , datas = items_dict .popitem ()
306+ table = tools .get_info (tab_item , ":s_(.*?)_item$" , fetch_one = True )
304307
305308 log .debug (
306309 """
307310 -------------- item 批量入库 --------------
308311 表名: %s
309312 datas: %s
310313 """
311- % (tab_item , tools .dumps_json (datas , indent = 16 ))
314+ % (table , tools .dumps_json (datas , indent = 16 ))
312315 )
313316
314- export_success = self .__export_to_db (tab_item , datas )
317+ export_success = self .__export_to_db (table , datas )
318+ if not export_success :
319+ failed_items ["add" ] = {"table" : table , "datas" : datas }
320+ break
315321
316322 # 执行批量update
317323 while update_items_dict :
318324 tab_item , datas = update_items_dict .popitem ()
325+ table = tools .get_info (tab_item , ":s_(.*?)_item$" , fetch_one = True )
326+
319327 log .debug (
320328 """
321329 -------------- item 批量更新 --------------
322330 表名: %s
323331 datas: %s
324332 """
325- % (tab_item , tools .dumps_json (datas , indent = 16 ))
333+ % (table , tools .dumps_json (datas , indent = 16 ))
326334 )
327335
328336 update_keys = self ._item_update_keys .get (tab_item )
329337 export_success = self .__export_to_db (
330- tab_item , datas , is_update = True , update_keys = update_keys
338+ table , datas , is_update = True , update_keys = update_keys
331339 )
340+ if not export_success :
341+ failed_items ["update" ] = {"table" : table , "datas" : datas }
342+ break
332343
333344 if export_success :
334345 # 执行回调
@@ -361,10 +372,17 @@ def __add_item_to_db(
361372 )
362373
363374 if self .export_falied_times > setting .EXPORT_DATA_MAX_RETRY_TIMES :
375+ failed_items ["requests" ] = requests
376+ self .redis_db .sadd (self ._table_failed_items , failed_items )
377+
364378 # 删除做过的request
365379 if requests :
366380 self .redis_db .zrem (self ._table_request , requests )
367- log .error ("入库超过最大重试次数,不再重试" )
381+ log .error (
382+ "入库超过最大重试次数,不再重试, 数据记录到redis, items:\n {}" .format (
383+ tools .dumps_json (failed_items )
384+ )
385+ )
368386 else :
369387 tip = ["入库不成功" ]
370388 if callbacks :
0 commit comments