最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Python實(shí)現(xiàn)數(shù)據(jù)庫(kù)并行讀取和寫(xiě)入實(shí)例

 更新時(shí)間:2017年06月09日 11:12:41   作者:AsuraDong  
本篇文章主要介紹了Python實(shí)現(xiàn)數(shù)據(jù)庫(kù)并行讀取和寫(xiě)入實(shí)例,非常具有實(shí)用價(jià)值,需要的朋友可以參考下

這篇主要記錄一下如何實(shí)現(xiàn)對(duì)數(shù)據(jù)庫(kù)的并行運(yùn)算來(lái)節(jié)省代碼運(yùn)行時(shí)間。語(yǔ)言是Python,其他語(yǔ)言思路一樣。

前言

一共23w條數(shù)據(jù),是之前通過(guò)自然語(yǔ)言分析處理過(guò)的數(shù)據(jù),附一張截圖:


要實(shí)現(xiàn)對(duì)news主體的讀取,并且找到其中含有的股票名稱(chēng),只要發(fā)現(xiàn),就將這支股票和對(duì)應(yīng)的日期、score寫(xiě)入數(shù)據(jù)庫(kù)。

顯然,幾十萬(wàn)條數(shù)據(jù)要是一條條讀寫(xiě),然后在本機(jī)上操作,耗時(shí)太久,可行性極低。所以,如何有效并行的讀取內(nèi)容,并且進(jìn)行操作,最后再寫(xiě)入數(shù)據(jù)庫(kù)呢?

并行讀取和寫(xiě)入

并行讀?。簞?chuàng)建N*max_process個(gè)進(jìn)程,對(duì)數(shù)據(jù)庫(kù)進(jìn)行讀取。讀取的時(shí)候應(yīng)該注意:

  1. 每個(gè)進(jìn)程需要分配不同的connection和對(duì)應(yīng)的cursor,否則數(shù)據(jù)庫(kù)會(huì)報(bào)錯(cuò)。
  2. 數(shù)據(jù)庫(kù)必須能承受相應(yīng)的高并發(fā)訪問(wèn)(可以手動(dòng)更改)

實(shí)現(xiàn)的時(shí)候,如果不在進(jìn)程里面創(chuàng)建新的connection,就會(huì)發(fā)生沖突,每個(gè)進(jìn)程拿到權(quán)限后,會(huì)被下個(gè)進(jìn)程釋放,所以匯報(bào)出來(lái)NoneType Error的錯(cuò)誤。

  1. 并行寫(xiě)入:在對(duì)數(shù)據(jù)庫(kù)進(jìn)行更改的時(shí)候,不可以多進(jìn)程更改。所以,我們需要根據(jù)已有的表,創(chuàng)建max_process-1個(gè)同樣結(jié)構(gòu)的表用來(lái)寫(xiě)入。表的命名規(guī)則可以直接在原來(lái)基礎(chǔ)上加上1,2,3...數(shù)字可以通過(guò)對(duì)max_process取余得到。

此時(shí),對(duì)應(yīng)進(jìn)程里面先后出現(xiàn)讀入的conn(保存消息后關(guān)閉)和寫(xiě)入的conn。每個(gè)進(jìn)程對(duì)應(yīng)的表的index就是 主循環(huán)中的num對(duì)max_process取余(100->4,101->5),這樣每個(gè)進(jìn)程只對(duì)一個(gè)表進(jìn)行操作了。

部分代碼實(shí)現(xiàn)

max_process = 16 #最大進(jìn)程數(shù)

def read_SQL_write(r_host,r_port,r_user,r_passwd,r_db,r_charset,w_host,w_port,w_user,w_passwd,w_db,w_charset,cmd,index=None):
  #得到tem字典保存著信息
  try:
    conn = pymysql.Connect(host=r_host, port=r_port, user=r_user, passwd =r_passwd, db =r_db, charset =r_charset)
    cursor = conn.cursor()
    cursor.execute(cmd)
  except Exception as e:
    error = "[-][-]%d fail to connect SQL for reading" % index
    log_error('error.log',error)
    return 
  else:
    tem = cursor.fetchone()
    print('[+][+]%d succeed to connect SQL for reading' % index)
  finally:
    cursor.close()
    conn.close()
  
  try:
    conn = pymysql.Connect(host=w_host, port=w_port, user=w_user, passwd =w_passwd, db =w_db, charset =w_charset)
    cursor = conn.cursor()
    cursor.execute(cmd)
  except Exception as e:
    error = "[-][-]%d fail to connect SQL for writing" % index
    log_error('error.log',error)
    return 
  else:
    print('[+][+]%d succeed to connect SQL for writing' % index)
  
  
  r_dict = dict()
  r_dict['id'] = tem[0]
  r_dict['content_id'] = tem[1]
  r_dict['pub_date'] = tem[2]
  r_dict['title'] = cht_to_chs(tem[3])
  r_dict['title_score'] =tem[4]![](http://images2015.cnblogs.com/blog/1172464/201706/1172464-20170609000900309-1810357590.png)

  r_dict['news_content'] = cht_to_chs(tem[5])
  r_dict['content_score'] = tem[6]
  
  for key in stock_dict.keys():
    #能找到對(duì)應(yīng)的股票
    if stock_dict[key][1] and ( r_dict['title'].find(stock_dict[key][1])!=-1 or r_dict['news_content'].find(stock_dict[key][1])!=-1 ):
      w_dict=dict()
      w_dict['code'] = key
      w_dict['english_name'] = stock_dict[key][0]
      w_dict['cn_name'] = stock_dict[key][1]
      #得到分?jǐn)?shù)
      if r_dict['title_score']:
        w_dict['score']=r_dict['title_score']
      else:
        w_dict['score']=r_dict['content_score']
      
      #開(kāi)始寫(xiě)入
      try:
        global max_process
        cmd = "INSERT INTO dyx_stock_score%d VALUES ('%s', '%s' , %d , '%s' , '%s' , %.2f );" % \
          (index%max_process ,r_dict['content_id'] ,r_dict['pub_date'] ,w_dict['code'] ,w_dict['english_name'] ,w_dict['cn_name'] ,w_dict['score'])
        cursor.execute(cmd)
        conn.commit()
      except Exception as e:
        error = "  [-]%d fail to write to SQL" % index
        cursor.rollback()
        log_error('error.log',error)
      else:
        print("  [+]%d succeed to write to SQL" % index)

  cursor.close()
  conn.close()
def main():
  num = 238143#數(shù)據(jù)庫(kù)查詢(xún)拿到的總數(shù)
  p = None
  for index in range(1,num+1):
    if index%max_process==1:
      if p:
        p.close()
        p.join()
      p = multiprocessing.Pool(max_process)
    r_cmd = ('select id,content_id,pub_date,title,title_score,news_content,content_score from dyx_emotion_analysis where id = %d;' % (index))
    p.apply_async(func = read_SQL_write,args=(r_host,r_port,r_user,r_passwd,r_db,r_charset,w_host,w_port,w_user,w_passwd,w_db,w_charset,r_cmd,index,))

  if p:
    p.close()
    p.join()

以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。

相關(guān)文章

最新評(píng)論

邵东县| 乐昌市| 郎溪县| 沂南县| 罗江县| 佳木斯市| 临江市| 和平区| 长寿区| 莱州市| 江川县| 罗城| 枣庄市| 十堰市| 灌阳县| 阿克陶县| 苏尼特左旗| 泸西县| 习水县| 离岛区| 左权县| 吴忠市| 大名县| 邵阳县| 邻水| 河间市| 菏泽市| 成安县| 平度市| 兴安盟| 邵阳市| 巴马| 黑龙江省| 旬阳县| 新营市| 三河市| 安丘市| 襄樊市| 德令哈市| 黄骅市| 太原市|