【问题标题】:writing to csv file using map_async使用 map_async 写入 csv 文件
【发布时间】:2015-09-07 10:37:07
【问题描述】:

我有以下代码:

#!/usr/bin/env python

def do_job(row):
  # COMPUTING INTENSIVE OPERATION
  sleep(1)
  row.append(int(row[0])**2)

  # WRITING TO FILE - ATOMICITY ENSURED
  semaphore.acquire()
  print "Inside semaphore before writing to file: (%s,%s,%s)" % (row[0], row[1], row[2])
  csvWriter.writerow(row)
  print "Inside semaphore after writing to file"
  semaphore.release()

  # RETURNING VALUE
  return row

def parallel_csv_processing(inputFile, header=["Default", "header", "please", "change"], separator=",", skipRows = 0, cpuCount = 1):
  # OPEN FH FOR READING INPUT FILE
  inputFH   = open(inputFile,  "rb")
  csvReader = csv.reader(inputFH, delimiter=separator)

  # SKIP HEADERS
  for skip in xrange(skipRows):
    csvReader.next()

  # WRITE HEADER TO OUTPUT FILE
  csvWriter.writerow(header)

  # COMPUTING INTENSIVE OPERATIONS
  try:
    p = Pool(processes = cpuCount)
    # results = p.map(do_job, csvReader, chunksize = 10)
    results = p.map_async(do_job, csvReader, chunksize = 10)

  except KeyboardInterrupt:
    p.close()
    p.terminate()
    p.join()

  # WAIT FOR RESULTS
  # results.get()
  p.close()
  p.join()

  # CLOSE FH FOR READING INPUT
  inputFH.close()

if __name__ == '__main__':
  import csv
  from time import sleep
  from multiprocessing import Pool
  from multiprocessing import cpu_count
  from multiprocessing import current_process
  from multiprocessing import Semaphore
  from pprint import pprint as pp
  import calendar
  import time

  SCRIPT_START_TIME  = calendar.timegm(time.gmtime())
  inputFile  = "input.csv"
  outputFile = "output.csv"
  semaphore = Semaphore(1)

  # OPEN FH FOR WRITING OUTPUT FILE
  outputFH  = open(outputFile, "wt")
  csvWriter = csv.writer(outputFH, lineterminator='\n')
  csvWriter.writerow(["before","calling","multiprocessing"])
  parallel_csv_processing(inputFile, cpuCount = cpu_count())
  csvWriter.writerow(["after","calling","multiprocessing"])

  # CLOSE FH FOR WRITING OUTPUT
  outputFH.close()

  SCRIPT_STOP_TIME   = calendar.timegm(time.gmtime())
  SCRIPT_DURATION    = SCRIPT_STOP_TIME - SCRIPT_START_TIME
  print "Script duration:    %s seconds" % SCRIPT_DURATION

在终端运行后输出如下:

Inside semaphore before writing to file: (0,0,0)
Inside semaphore after writing to file
Inside semaphore before writing to file: (1,3,1)
Inside semaphore after writing to file
Inside semaphore before writing to file: (2,6,4)
Inside semaphore after writing to file
Inside semaphore before writing to file: (3,9,9)
Inside semaphore after writing to file
Inside semaphore before writing to file: (4,12,16)
Inside semaphore after writing to file
Inside semaphore before writing to file: (5,15,25)
Inside semaphore after writing to file
Inside semaphore before writing to file: (6,18,36)
Inside semaphore after writing to file
Inside semaphore before writing to file: (7,21,49)
Inside semaphore after writing to file
Inside semaphore before writing to file: (8,24,64)
Inside semaphore after writing to file
Inside semaphore before writing to file: (9,27,81)
Inside semaphore after writing to file
Script duration:    10 seconds

input.csv的内容如下:

0,0
1,3
2,6
3,9
4,12
5,15
6,18
7,21
8,24
9,27

output.csv 创建的内容如下:

before,calling,multiprocessing
Default,header,please,change
after,calling,multiprocessing

为什么没有从parallel_csv_processing 分别写给output.csvdo_job 方法?

【问题讨论】:

    标签: python csv semaphore python-multiprocessing


    【解决方案1】:

    您的进程正在静默失败并出现异常 - 具体而言,在生成的进程中,脚本没有 csvWriter 的值,因为它们每个都在单独的 python 解释器中,并且没有运行 main() - 这是故意的,您不希望子进程运行主进程。 do_job() 函数只能访问您在 map_async() 调用中显式传递给它的值,并且您没有传递 csvWriter。即使您不确定它是否会工作,也不知道文件句柄是否在主进程和多处理创建的进程之间共享。

    在 do_job 中的代码周围放置一个 try/except ,您将看到异常。

    def do_job(row):
      try:
          # COMPUTING INTENSIVE OPERATION
          sleep(1)
          row.append(int(row[0])**2)
    
          # WRITING TO FILE - ATOMICITY ENSURED
          semaphore.acquire()
          print "Inside semaphore before writing to file: (%s,%s,%s)" % (row[0], row[1], row[2])
          csvWriter.writerow(row)
          print "Inside semaphore after writing to file"
          semaphore.release()
    
          # RETURNING VALUE
          return row
      except:
          print "exception"
    

    显然,在实际代码中应该正确处理异常,但是如果您运行它,您现在会看到每次调用 do_job 时都会打印异常。

    查看多处理文档以获得更多指导 - 在 Python 2.7 标准库文档中的标题“16.6.1.4。在进程之间共享状态”下。

    【讨论】:

    • 感谢您的回复。所以你是说没有直接传递给map_async() 的对象在do_job() 中是不可见的,对吧?那么semaphore 怎么可能是可见的呢?此外,当我在__main__ 中声明global_variable = "global variable" 时,我可以使用print global_variabledo_job 打印它。这不是同一个场景吗?
    • 我想问为什么这个异常是静默的?根据我以前在引发异常时的经验,我立即得到了通知。是否可以明确地说不要默默地抑制异常?我可以使用tryexcept,但我怎么知道哪部分代码导致了问题?是否有一些用于此目的的全局设置?
    • 没有 cPython 实现将异常从不同进程返回到 main() 代码的方法 - 您必须自己实现,可能会看到 stackoverflow.com/questions/6728236/…stackoverflow.com/questions/16943404/…
    • 谢谢,我会检查链接。您能否也对我提到的传递变量做出反应?一种可能的解释是文件处理程序不允许在进程之间共享,而其他类型是。这个对吗?谢谢
    • 多处理创建的进程在完全独立的内存空间中运行(例如,它们可能是不同的 CPU 内核) - 如果您想要单个解释器/内存空间,请使用具有非常相似 API 但没有的线程t 使用操作系统进程,因此必须在单个 CPU 内核中运行。阅读有关多处理的 Python 标准库文档,尝试一下。
    猜你喜欢
    • 2011-08-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-01-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多