preprocessing_db.py 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137
  1. # -*- coding: utf-8 -*-
  2. """遗留入库与批量处理脚本。
  3. 按 ARCHITECTURE.md Phase 4 要求,从 ``interface/Preprocessing.py`` 迁出。
  4. 类型:LEGACY(非运行时,仅用于历史数据迁移或调试)。
  5. 原位置:``interface/Preprocessing.py`` 中以下函数:
  6. - ``persistenceData`` — 将中间结果保存到数据库(线上不执行)
  7. - ``persistenceData1`` — 将实体/句子中间结果保存到数据库(线上不执行)
  8. - ``_handle`` — 批量表格解析的子任务处理函数
  9. - ``getPredictTable`` — 批量表格预测脚本入口
  10. ``interface/Preprocessing.py`` 仍 re-export 以上全部名称,老 import 不受影响。
  11. """
  12. from __future__ import absolute_import
  13. import json
  14. from bs4 import BeautifulSoup
  15. from BiddingKG.dl.preprocess.table_parser import tableToText
  16. __all__ = [
  17. "persistenceData",
  18. "persistenceData1",
  19. "_handle",
  20. "getPredictTable",
  21. ]
  22. def persistenceData(data):
  23. '''
  24. @summary:将中间结果保存到数据库-线上生产的时候不需要执行
  25. '''
  26. # Phase 1: PG 连接走 infra/db(原硬编码 host=192.168.2.101, password=postgres 已移除)
  27. from BiddingKG.dl.infra.db import get_connection
  28. conn = get_connection("BiddingKG")
  29. cursor = conn.cursor()
  30. for item_index in range(len(data)):
  31. item = data[item_index]
  32. doc_id = item[0]
  33. dic = item[1]
  34. code = dic['code']
  35. name = dic['name']
  36. prem = dic['prem']
  37. if len(code)==0:
  38. code_insert = ""
  39. else:
  40. code_insert = ";".join(code)
  41. prem_insert = ""
  42. for item in prem:
  43. for x in item:
  44. if isinstance(x, list):
  45. if len(x)>0:
  46. for x1 in x:
  47. prem_insert+="/".join(x1)+","
  48. prem_insert+="$"
  49. else:
  50. prem_insert+=str(x)+"$"
  51. prem_insert+=";"
  52. sql = " insert into predict_validation(doc_id,code,name,prem) values('"+doc_id+"','"+code_insert+"','"+name+"','"+prem_insert+"')"
  53. cursor.execute(sql)
  54. conn.commit()
  55. conn.close()
  56. def persistenceData1(list_entitys,list_sentences):
  57. '''
  58. @summary:将中间结果保存到数据库-线上生产的时候不需要执行
  59. '''
  60. # Phase 1: PG 连接走 infra/db
  61. from BiddingKG.dl.infra.db import get_connection
  62. conn = get_connection("BiddingKG")
  63. cursor = conn.cursor()
  64. for list_entity in list_entitys:
  65. for entity in list_entity:
  66. if entity.values is not None:
  67. sql = " insert into predict_entity(entity_id,entity_text,entity_type,doc_id,sentence_index,begin_index,end_index,label,values) values('"+str(entity.entity_id)+"','"+str(entity.entity_text)+"','"+str(entity.entity_type)+"','"+str(entity.doc_id)+"',"+str(entity.sentence_index)+","+str(entity.begin_index)+","+str(entity.end_index)+","+str(entity.label)+",array"+str(entity.values)+")"
  68. else:
  69. sql = " insert into predict_entity(entity_id,entity_text,entity_type,doc_id,sentence_index,begin_index,end_index) values('"+str(entity.entity_id)+"','"+str(entity.entity_text)+"','"+str(entity.entity_type)+"','"+str(entity.doc_id)+"',"+str(entity.sentence_index)+","+str(entity.begin_index)+","+str(entity.end_index)+")"
  70. cursor.execute(sql)
  71. for list_sentence in list_sentences:
  72. for sentence in list_sentence:
  73. str_tokens = "["
  74. for item in sentence.tokens:
  75. str_tokens += "'"
  76. if item=="'":
  77. str_tokens += "''"
  78. else:
  79. str_tokens += item
  80. str_tokens += "',"
  81. str_tokens = str_tokens[:-1]+"]"
  82. sql = " insert into predict_sentences(doc_id,sentence_index,tokens) values('"+sentence.doc_id+"',"+str(sentence.sentence_index)+",array"+str_tokens+")"
  83. cursor.execute(sql)
  84. conn.commit()
  85. conn.close()
  86. def _handle(item,result_queue):
  87. dochtml = item["dochtml"]
  88. docid = item["docid"]
  89. list_innerTable = tableToText(BeautifulSoup(dochtml,"lxml"))
  90. flag = False
  91. if list_innerTable:
  92. flag = True
  93. for table in list_innerTable:
  94. result_queue.put({"docid":docid,"json_table":json.dumps(table,ensure_ascii=False)})
  95. def getPredictTable():
  96. filename = "D:\Workspace2016\DataExport\data\websouce_doc.csv"
  97. import pandas as pd
  98. import json
  99. from BiddingKG.dl.common.MultiHandler import MultiHandler,Queue
  100. df = pd.read_csv(filename)
  101. df_data = {"json_table":[],"docid":[]}
  102. _count = 0
  103. _sum = len(df["docid"])
  104. task_queue = Queue()
  105. result_queue = Queue()
  106. _index = 0
  107. for dochtml,docid in zip(df["dochtmlcon"],df["docid"]):
  108. task_queue.put({"docid":docid,"dochtml":dochtml,"json_table":None})
  109. _index += 1
  110. mh = MultiHandler(task_queue=task_queue,task_handler=_handle,result_queue=result_queue,process_count=5,thread_count=1)
  111. mh.run()
  112. while True:
  113. try:
  114. item = result_queue.get(block=True,timeout=1)
  115. df_data["docid"].append(item["docid"])
  116. df_data["json_table"].append(item["json_table"])
  117. except Exception as e:
  118. print(e)
  119. break
  120. df_1 = pd.DataFrame(df_data)
  121. df_1.to_csv("../form/websource_67000_table.csv",columns=["docid","json_table"])