【问题标题】:How to use multiprocessing in current Python application?如何在当前 Python 应用程序中使用多处理?
【发布时间】:2021-09-14 07:23:44
【问题描述】:

我有一个从不同目录读取数千个文件的应用程序,它读取它们,对它们进行一些处理,然后将数据发送到数据库。我有 1 个问题,大约需要。 1 小时完成 1 个目录中的所有文件,我有 19 个目录(将来可能会更多)。现在它一个接一个地做,我想并行运行所有东西,所以我加快了速度。

这是我的代码:

import mysql.connector
import csv
import os
import time
from datetime import datetime
import ntpath
import configparser

config = configparser.ConfigParser()                         
config.read('C:\Desktop\Energy\file_cfg.ini')

source = config['PATHS']['source']
archive = config['PATHS']['archive']

mydb = mysql.connector.connect(
        host= config['DB']['host'],
        user = config['DB']['user'],
        passwd = config['DB']['passwd'],
        database= config['DB']['database']
    )

cursor = mydb.cursor()
select_antenna = "SELECT * FROM `antenna`"
cursor.execute(select_antenna)

mp_mysql = [i[0] for i in cursor.fetchall()]
mp_server = os.listdir(source)

# microbeats clean.
cursor.execute("TRUNCATE TABLE microbeats")

for mp in mp_mysql:
    if mp in mp_server:
        subdir_paths = os.path.join(source, mp)

        for file in os.listdir(subdir_paths):
            file_paths = os.path.join(subdir_paths, file)
                        
            cr_time_s = os.path.getctime(file_paths)
            cr_time = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(cr_time_s))
                        
        all_file_paths = [os.path.join(subdir_paths, f) for f in os.listdir(subdir_paths)]
        full_file_paths = [p for p in all_file_paths if os.path.getsize(p) > 0] #<----- Control empty files.
                
        if full_file_paths != []:
            newest_file_paths = max(full_file_paths, key=os.path.getctime)
                        
            for file in all_file_paths:
                if file == newest_file_paths and os.path.getctime(newest_file_paths) < time.time() - 120: 
                    with open(file, 'rt') as f:
                        reader = csv.reader(f, delimiter ='\t')
                        line_data0 = list()
                        col = next(reader)
                        
                        for line in reader:
                            line.insert(0, mp)
                            line.insert(1, cr_time)
                            if line != []: #<----- Control empty directories.
                                line_data0.append(line)
                            
                                                                                                                                              
                        q1 = ("INSERT INTO microbeats"
                                     "(`antenna`,`datetime`,`system`,`item`,`event`, `status`, `accident`)"
                                     "VALUES (%s, %s, %s,%s, %s, %s, %s)")

                        for line in line_data0:
                             cursor.execute(q1, line)  

【问题讨论】:

  • 这里有 CPU-bound 和 I/O-bound 的代码。是的,多处理会加快速度,但这里有一些设计问题。尝试将 for 循环 for file in all_file_paths 转换为函数,然后使用 concurrent.futures.ProcessPoolExecutor.map 执行该函数。如果这速度足够快,那就收工吧。否则,您需要通过预处理查询(可能具有并发性)然后与线程同时执行查询来解决您的设计问题。
  • 首先,您可能并不真正需要多处理,因为多线程对于这种类型的应用程序来说很好。为了提高效率,您应该查看 execute_many() 而不是执行大量离散插入。使用连接池,因为 MySQL 连接不是线程安全的。不要忘记提交您的更改!
  • 其实不用multiprocessing,这个问题更适合Code Review。请试一试,发布您的代码、预期结果和实际结果,或转发至Code Review
  • 你们能用代码告诉我你的意思吗?这也让我有机会接受你的回答。

标签: python multithreading multiprocessing config mysql-python


【解决方案1】:

我正在使用多处理,其中每个进程都有自己的数据库连接。我已经对您的代码进行了最小的更改,以尝试并行处理目录。但是,我不确定 subdir_paths 这样的变量是否正确命名,因为其名称末尾的“s”表示它包含多个路径名。

有人建议这个问题更适合 Code Review 的原因是因为大概您有一个已经在运行的程序,并且您只是在寻找性能改进(当然,这适用于很大一部分的在 SO 上发布的带有multiprocessing 标记的问题)。这类问题应该发到https://codereview.stackexchange.com/

import mysql.connector
import csv
import os
import time
from datetime import datetime
import ntpath
import configparser
from multiprocessing import Pool, cpu_count

config = configparser.ConfigParser()
config.read('C:\Desktop\Energy\file_cfg.ini')

source = config['PATHS']['source']
archive = config['PATHS']['archive']

def get_connnection():
    mydb = mysql.connector.connect(
            host= config['DB']['host'],
            user = config['DB']['user'],
            passwd = config['DB']['passwd'],
            database= config['DB']['database']
        )
    return mydb

def get_mp_list():
    select_antenna = "SELECT * FROM `antenna`"
    mydb = get_connection()
    cursor = mydb.cursor()
    cursor.execute(select_antenna)

    mp_mysql = [i[0] for i in cursor.fetchall()]
    mp_server = os.listdir(source)

    # microbeats clean.
    cursor.execute("TRUNCATE TABLE microbeats")
    mydb.commit()
    mydb.close()

    mp_list = [mp for mp in mp_mysql if mp in mp_server]
    return mp_list

def process_mp(mp):
    subdir_paths = os.path.join(source, mp)
    for file in os.listdir(subdir_paths):
        file_paths = os.path.join(subdir_paths, file)

        cr_time_s = os.path.getctime(file_paths)
        cr_time = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(cr_time_s))

    all_file_paths = [os.path.join(subdir_paths, f) for f in os.listdir(subdir_paths)]
    full_file_paths = [p for p in all_file_paths if os.path.getsize(p) > 0] #<----- Control empty files.

    if full_file_paths != []:
        newest_file_paths = max(full_file_paths, key=os.path.getctime)

        mydb = get_connection()
        cursor = mydb.cursor()
        did_insert = False
        q1 = ("INSERT INTO microbeats"
                     "(`antenna`,`datetime`,`system`,`item`,`event`, `status`, `accident`)"
                     "VALUES (%s, %s, %s,%s, %s, %s, %s)")
        for file in all_file_paths:
            if file == newest_file_paths and os.path.getctime(newest_file_paths) < time.time() - 120:
                with open(file, 'rt') as f:
                    reader = csv.reader(f, delimiter ='\t')
                    line_data0 = list()
                    col = next(reader)

                    for line in reader:
                        line.insert(0, mp)
                        line.insert(1, cr_time)
                        if line != []: #<----- Control empty directories.
                            line_data0.append(line)

                    if line_data0:
                        cursor.executemany(q1, line_data0)
                        did_insert = True
        if did_insert:
            mydb.commit()
        mydb.close()

def main():
    mp_list = get_mp_list()
    pool = Pool(min(cpu_count(), len(mp_list)))
    results = pool.imap_unordered(process_mp, mp_list)
    while True:
        try:
            result = next(results)
        except StopIteration:
            break
        except BaseException as e:
            print(e)

if __name__ == '__main__':
    main()

【讨论】:

  • ReferenceError: weakly-referenced object no longer exists
  • @Mediterráneo 错误在哪里?由于显而易见的原因,我无法自己运行代码。
  • Traceback (most recent call last): File "c:/data/Energy/Desktop/test_booboo.py", line 86, in &lt;module&gt; main() File "c:/data/Energy/Desktop/test_booboo.py", line 79, in main all_subdir_paths = get_all_subdir_paths() File "c:/data/Energy/Desktop/test_booboo.py", line 28, in get_all_subdir_paths cursor.execute(select_antennas) File "C:\ProgramData\Anaconda3\lib\site-packages\mysql\connector\cursor_cext.py", line 232, in execute if not self._cnx: ReferenceError: weakly-referenced object no longer exists
  • 尝试更新的代码,尽管我可能已将错误移至别处。如果是这样,我还有一件事要尝试。
  • 我已经更新了代码,试图摆脱 mp 未定义错误。
猜你喜欢
  • 2016-11-21
  • 1970-01-01
  • 2016-04-13
  • 2018-07-31
  • 2017-06-26
  • 2019-05-20
  • 2020-09-21
  • 2023-03-11
  • 2019-12-10
相关资源
最近更新 更多