KnowledgeHub
Questions
Tags
Users
Search
Alex Rivera
|
Logout
Edit Question
Title
Body
I am fairly new to python. I am using the multiprocessing module for reading lines of text on stdin, converting them in some way and writing them into a database. Here's a snippet of my code: batch = [] pool = multiprocessing.Pool(20) i = 0 for i, content in enumerate(sys.stdin): batch.append(content) if len(batch) >= 10000: pool.apply_async(insert, args=(batch,i+1)) batch = [] pool.apply_async(insert, args=(batch,i)) pool.close() pool.join() Now that all works fine, until I get to process huge input files (hundreds of millions of lines) that i pipe into my python program. At some point, when my database gets slower, I see the memory getting full. After some playing, it turned out that pool.apply_async as well as pool.map_async never ever block, so that the queue of the calls to be processed grows bigger and bigger. What is the correct approach to my problem? I would expect a parameter that I can set, that will block the pool.apply_async call, as soon as a certain queue length has been reached. AFAIR in Java one can give the ThreadPoolExecutor a BlockingQueue with a fixed length for that purpose. Thanks!
Tags (comma-separated)
Save Edits
Cancel