【问题标题】:Python - Mult-Threading Help - Reading Multiple Files - ETL Into SQL ServerPython - 多线程帮助 - 读取多个文件 - ETL 到 SQL Server
【发布时间】:2017-06-21 04:14:44
【问题描述】:

我正在开发一个从本地驱动器读取 DBF 文件并将数据加载到 sql server 表中的程序。我对 Python 非常熟悉,并且我发现了一些关于多线程的细节,其中大部分都令人困惑。读取和插入的性能很慢,看看我的 CPU 使用率,我有足够的容量。我也在运行 SSD。

此代码将扩展为从大约 400 个 zip 文件中的大约 20 个 DBF 文件中读取。所以我们总共讨论了 8000 个 DBF 文件。

我很难做到这一点。可以指点一下吗?

这是我的代码(有点乱,但我稍后会清理它),

import os, pyodbc, datetime, shutil
from dbfread import DBF
from zipfile import ZipFile

# SQL Server Connection Test
cnxn = pyodbc.connect('DRIVER={SQL Server};SERVER=localhost\test;DATABASE=TEST_DBFIMPORT;UID=test;PWD=test')
cursor = cnxn.cursor()

dr = 'e:\\Backups\\dbf\\'
work = 'e:\\Backups\\work\\'
archive = 'e:\\Backups\\archive\\'


for r in os.listdir(dr):

    curdate = datetime.datetime.now()
    filepath = dr + r
    process = work + r
    arc = archive + r

    pth = r.replace(".sss","")
    zipfolder = work + pth
    filedateunix = os.path.getctime(filepath)
    filedateconverted=datetime.datetime.fromtimestamp(int(filedateunix)
                                                  ).strftime('%Y-%m-%d %H:%M:%S')
    shutil.move(filepath,process)
    with ZipFile(process) as zf:
        zf.extractall(zipfolder)


    cursor.execute(
        "insert into tblBackups(backupname, filedate, dateadded) values(?,?,?)",
    pth, filedateconverted, curdate)
    cnxn.commit()

    for dirpath, subdirs, files in os.walk (zipfolder):

        for file in files:
            dateadded = datetime.datetime.now()

            if file.endswith(('.dbf','.DBF')):
                dbflocation = os.path.abspath(os.path.join(dirpath,file)).lower()

                if dbflocation.__contains__("\\bk.dbf"):
                    table = DBF(dbflocation, lowernames=True, char_decode_errors='ignore')
                    for record in table.records:
                        rec1 = str(record['code'])
                        rec2 = str(record['name'])
                        rec3 = str(record['addr1'])
                        rec4 = str(record['addr2'])
                        rec5 = str(record['city'])
                        rec6 = str(record['state'])
                        rec7 = str(record['zip'])
                        rec8 = str(record['tel'])
                        rec9 = str(record['fax'])
                        cursor.execute(
                       "insert into tblbk(code,name,addr1,addr2,city,state,zip,tel,fax) values(?,?,?,?,?,?,?,?,?)",
                        rec1, rec2, rec3, rec4, rec5, rec6, rec7, rec8, rec9, rec10, rec11, rec12, rec13)
                cnxn.commit()


                if dbflocation.__contains__("\\cr.dbf"):
                    table = DBF(dbflocation, lowernames=True, char_decode_errors='ignore')
                    for record in table.records:
                        rec2 = str(record['cal_desc'])
                        rec3 = str(record['b_date'])
                        rec4 = str(record['b_time'])
                        rec5 = str(record['e_time'])
                        rec6 = str(record['with_desc'])
                        rec7 = str(record['recuruntil'])
                        rec8 = record['notes']
                        rec9 = dateadded
                        cursor.execute(
                        "insert into tblcalendar(cal_desc,b_date,b_time,e_time,with_desc,recuruntil,notes,dateadded) values(?,?,?,?,?,?,?,?)",
                        rec2, rec3, rec4, rec5, rec6, rec7, rec8, rec9)
                cnxn.commit() 

    shutil.move(process, archive)
    shutil.rmtree(zipfolder)

【问题讨论】:

  • 我想要的另一个选择是多处理,它可能更简单。

标签: python multithreading python-2.7 multiprocessing dbf


【解决方案1】:

tl;dr:先测量,后修复!


请注意,在最常见的 Python 实现 (CPython) 中,一次只能有一个线程执行 Python 字节码。 因此线程不是提高 CPU 密集型性能的好方法。如果工作是 I/O 密集型的,它们可以很好地工作。

但你首先应该做的是测量。这一点怎么强调都不过分。如果您不知道导致性能下降的原因,则无法修复它!

编写完成这项工作的单线程代码,并在分析器下运行它。先试试内置的cProfile。如果这不能给你足够的信息,请尝试例如line profiler。

分析应该告诉您哪些步骤消耗的时间最多。一旦你知道了这一点,你就可以开始改进了。

例如,使用multiprocessing 读取 DBF 文件是没有意义的,如果这是将数据填充到 SQL 服务器的操作最耗时!这甚至可能会减慢速度,因为随后有几个进程正在争夺 SQL 服务器的注意力。

如果 SQL 服务器不是瓶颈,并且它可以处理多个连接,我会使用multiprocessing,可能是Pool.map() 来并行读取 DBF 并将数据填充到 SQL 服务器中。在这种情况下,您应该 Pool.map 覆盖 DBF 文件名列表,以便在工作进程中打开这些文件。

【讨论】:

  • 谢谢 Roland,你说的对。除了我的代码之外没有其他瓶颈。我能够使用多处理模块并一次加载 15 个进程,延迟为 1 秒。我这样做是因为我注意到有时进程会尝试获取同一个文件。现在我遇到了 SQL Server CPU 瓶颈,这正是我想要看到的。 SQL Server CPU 使用率从大约 8% 变为 50%,处理 8k 个文件的时间从 10 多个小时变为 45 分钟!
  • @HMan06 快 10 倍以上?不错!
【解决方案2】:

您可以尝试executemany() 方法而不是循环中的单个插入。下面是一些 ETL 脚本中插入函数的示例:

def sql_insert(table_name, fields, rows, truncate_table = True):
    if len(rows) == 0:
        return

    cursor = mdwh_connection.cursor()
    cursor.fast_executemany = True
    values_sql = ('?, ' * (fields.count(',') + 1))[:-2]

    if truncate_table:
        sql_truncate(table_name, cursor)
    
    insert_sql = 'insert {0} ({1}) values ({2});'.format(table_name, fields, values_sql)
    current_row = 0
    batch_size = 50000

    while current_row < len(rows):
        cursor.executemany(insert_sql, rows[current_row:current_row + batch_size])
        mdwh_connection.commit()
        current_row += batch_size
        logging.info(
            '{} more records inserted. Total: {}'.format(
                min(batch_size,len(rows)-current_row+batch_size),
                min(current_row, len(rows))
            )
        )    

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-03-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多