in python: child processes going defunct while others are not, unsure why

multiprocessing, parallel-processing, python

Solution

There's not enough information to be sure, but the problem is very likely to be that `Slave.do_work` is raising an unhandled exception. (There are many lines of your code that could do that in various different conditions.)

When you do that, the child process will just exit.

On POSIX systems… well, the full details are a bit complicated, but in the simple case (what you have here), a child process that exits will stick around as a `<defunct>` process until it gets reaped (because the parent either `wait`s on it, or exits). Since your parent code doesn't wait on the children until the queue is finished, that's exactly what happens.

So, there's a simple duct-tape fix:

def do_work(self):
    self.log(str(os.getpid()))
    while True:
        try:
            # the rest of your code
        except Exception as e:
            self.log("something appropriate {}".format(e))
            # you may also want to post a reply back to the parent

You might also want to break the massive `try` up into different ones, so you can distinguish between all the different stages where things could go wrong (especially if some of them mean you need a reply, and some mean you don't).

However, it looks like what you're attempting to do is duplicate exactly the behavior of `multiprocessing.Pool`, but have missed the bar in a couple places. Which raises the question: why not just use `Pool` in the first place? You could then simplify/optimize things ever further by using one of the `map` family methods. For example, your entire `Master.run` could be reduced to:

self.init()
pool = multiprocessing.Pool(Master.SLAVE_COUNT, initializer=slave_setup)
pool.map(slave_job, tables)
pool.join()

And this will handle exceptions for you, and allow you to return values/exceptions if you later need that, and let you use the built-in `logging` library instead of trying to build your own, and so on. And it should only take about a dozens lines of minor code changes to `Slave`, and then you're done.

If you want to submit new jobs from within jobs, the easiest way to do this is probably with a `Future`-based API (which turns things around, making the future result the focus and the pool/executor the dumb thing that provides them, instead of making the pool the focus and the result the dumb thing it gives back), but there are multiple ways to do it with `Pool` as well. For example, right now, you're not returning anything from each job, so, you can just return a list of `tables` to execute. Here's a simple example that shows how to do it:

import multiprocessing

def foo(x):
    print(x, x**2)
    return list(range(x))

if __name__ == '__main__':
    pool = multiprocessing.Pool(2)
    jobs = [5]
    while jobs:
        jobs, oldjobs = [], jobs
        for job in oldjobs:
            jobs.extend(pool.apply(foo, [job]))
    pool.close()
    pool.join()

Obviously you can condense this a bit by replacing the whole loop with, e.g., a list comprehension fed into `itertools.chain`, and you can make it a lot cleaner-looking by passing "a submitter" object to each job and adding to that instead of returning a list of new jobs, and so on. But I wanted to make it as explicit as possible to show how little there is to it.

At any rate, if you think the explicit queue is easier to understand and manage, go for it. Just look at the source for `multiprocessing.worker` and/or `concurrent.futures.ProcessPoolExecutor` to see what you need to do yourself. It's not that hard, but there are enough things you could get wrong (personally, I always forget at least one edge case when I try to do something like this myself) that it's work looking at code that gets it right.

Alternatively, it seems like the only reason you can't use `concurrent.futures.ProcessPoolExecutor` here is that you need to initialize some per-process state (the `boto.s3.key.Key`, `MySqlWrap`, etc.), for what are probably very good caching reasons. (If this involves a web-service query, a database connect, etc., you certainly don't want to do that once per task!) But there are a few different ways around that.

But you can subclass `ProcessPoolExecutor` and override the undocumented function `_adjust_process_count` (see the source for how simple it is) to pass your setup function, and… that's all you have to do.

Or you can mix and match. Wrap the `Future` from `concurrent.futures` around the `AsyncResult` from `multiprocessing`.

Problem

edit: the answer was that the os was axing processes because i was consuming all the memory i am spawning enough subprocesses to keep the load average 1:1 with cores, however at some point within the hour, this script could run for days, 3 of the processes go : ``` tipu 14804 0.0 0.0 328776 428 pts/1 Sl 00:20 0:00 python run.py tipu 14808 64.4 24.1 2163796 1848156 pts/1 Rl 00:20 44:41 python run.py tipu 14809 8.2 0.0 0 0 pts/1 Z 00:20 5:43 [python] <defunct> tipu 14810 60.3 24.3 2180308 1864664 pts/1 Rl 00:20 41:49 python run.py tipu 14811 20.2 0.0 0 0 pts/1 Z 00:20 14:04 [python] <defunct> tipu 14812 22.0 0.0 0 0 pts/1 Z 00:20 15:18 [python] <defunct> tipu 15358 0.0 0.0 103292 872 pts/1 S+ 01:30 0:00 grep python ``` i have no idea why this is happening, attached is the master and slave. i can attach the mysql/pg wrappers if needed as well, any suggestions? `slave.py`: ``` from boto.s3.key import Key import multiprocessing import gzip import os from mysql_wrapper import MySQLWrap from pgsql_wrapper import PGSQLWrap import boto import re class Slave: CHUNKS = 250000 BUCKET_NAME = "bucket" AWS_ACCESS_KEY = "" AWS_ACCESS_SECRET = "" KEY = Key(boto.connect_s3(AWS_ACCESS_KEY, AWS_ACCESS_SECRET).get_bucket(BUCKET_NAME)) S3_ROOT = "redshift_data_imports" COLUMN_CACHE = {} DEFAULT_COLUMN_VALUES = {} def __init__(self, job_queue): self.log_handler = open("logs/%s" % str(multiprocessing.current_process().name), "a"); self.mysql = MySQLWrap(self.log_handler) self.pg = PGSQLWrap(self.log_handler) self.job_queue = job_queue def do_work(self): self.log(str(os.getpid())) while True: #sample job in the abstract: mysql_db.table_with_date-iteration job = self.job_queue.get() #queue is empty if job is None: self.log_handler.close() self.pg.close() self.mysql.close() print("good bye and good day from %d" % (os.getpid())) self.job_queue.task_done() break #curtail iteration table = job.split('-')[0] #strip redshift table from job name redshift_table = re.sub(r"(_[1-9].*)", "", table.split(".")[1]) iteration = int(job.split("-")[1]) offset = (iteration - 1) * self.CHUNKS #columns redshift is expecting #bad tables will slip through and error out, so we catch it try: colnames = self.COLUMN_CACHE[redshift_table] except KeyError: self.job_queue.task_done() continue #mysql fields to use in SELECT statement fields = self.get_fields(table) #list subtraction determining which columns redshift has that mysql does not delta = (list(set(colnames) - set(fields.keys()))) #subtract columns that have a default value and so do not need padding if delta: delta = list(set(delta) - set(self.DEFAULT_COLUMN_VALUES[redshift_table])) #concatinate columns with padded \N select_fields = ",".join(fields.values()) + (",\\N" * len(delta)) query = "SELECT %s FROM %s LIMIT %d, %d" % (select_fields, table, offset, self.CHUNKS) rows = self.mysql.execute(query) self.log("%s: %s\n" % (table, len(rows))) if not rows: self.job_queue.task_done() continue #if there is more data potentially, add it to the queue if len(rows) == self.CHUNKS: self.log("putting %s-%s" % (table, (iteration+1))) self.job_queue.put("%s-%s" % (table, (iteration+1))) #various characters need escaping clean_rows = [] redshift_escape_chars = set( ["\\", "|", "\t", "\r", "\n"] ) in_chars = "" for row in rows: new_row = [] for value in row: if value is not None: in_chars = str(value) else: in_chars = "" #escape any naughty characters new_row.append("".join(["\\" + c if c in redshift_escape_chars else c for c in in_chars])) new_row = "\t".join(new_row) clean_rows.append(new_row) rows = ",".join(fields.keys() + delta) rows += "\n" + "\n".join(clean_rows) offset = offset + self.CHUNKS filename = "%s-%s.gz" % (table, iteration) self.move_file_to_s3(filename, rows) self.begin_data_import(job, redshift_table, ",".join(fields.keys() + delta)) self.job_queue.task_done() def move_file_to_s3(self, uri, contents): tmp_file = "/dev/shm/%s" % str(os.getpid()) self.KEY.key = "%s/%s" % (self.S3_ROOT, uri) self.log("key is %s" % self.KEY.key ) f = gzip.open(tmp_file, "wb") f.write(contents) f.close() #local saving allows for debugging when copy commands fail #text_file = open("tsv/%s" % uri, "w") #text_file.write(contents) #text_file.close() self.KEY.set_contents_from_filename(tmp_file, replace=True) def get_fields(self, table): """ Returns a dict used as: {"column_name": "altered_column_name"} Currently only the debug column gets altered """ exclude_fields = ["_qproc_id", "_mob_id", "_gw_id", "_batch_id", "Field"] query = "show columns from %s" % (table) fields = self.mysql.execute(query) #key raw field, value mysql formatted field new_fields = {} #for field in fields: for field in [val[0] for val in fields]: if field in exclude_fields: continue old_field = field if "debug_mode" == field.strip(): field = "IFNULL(debug_mode, 0)" new_fields[old_field] = field return new_fields def log(self, text): self.log_handler.write("\n%s" % text) def begin_data_import(self, table, redshift_table, fields): query = "copy %s (%s) from 's3://bucket/redshift_data_imports/%s' \ credentials 'aws_access_key_id=%s;aws_secret_access_key=%s' delimiter '\\t' \ gzip NULL AS '' COMPUPDATE ON ESCAPE IGNOREHEADER 1;" \ % (redshift_table, fields, table, self.AWS_ACCESS_KEY, self.AWS_ACCESS_SECRET) self.pg.execute(query) ``` `master.py`: ``` from slave import Slave as Slave import multiprocessing from mysql_wrapper import MySQLWrap as MySQLWrap from pgsql_wrapper import PGSQLWrap as PGSQLWrap class Master: SLAVE_COUNT = 5 def __init__(self): self.mysql = MySQLWrap() self.pg = PGSQLWrap() def do_work(table): pass def get_table_listings(self): """Gathers a list of MySQL log tables needed to be imported""" query = 'show databases' result = self.mysql.execute(query) #turns list[tuple] into a flat list databases = list(sum(result, ())) #overriding during development databases = ['db1', 'db2', 'db3']] exclude = ('mysql', 'Database', 'information_schema') scannable_tables = [] for database in databases: if database in exclude: continue query = "show tables from %s" % database result = self.mysql.execute(query) #turns list[tuple] into a flat list tables = list(sum(result, ())) for table in tables: exclude = ("Tables_in_%s" % database, "(", "201303", "detailed", "ltv") #exclude any of the unfavorables if any(s in table for s in exclude): continue scannable_tables.append("%s.%s-1" % (database, table)) return scannable_tables def init(self): #fetch redshift columns once and cache #get columns from redshift so we can pad the mysql column delta with nulls tables = ('table1', 'table2', 'table3') for table in tables: #cache columns query = "SELECT column_name FROM information_schema.columns WHERE \ table_name = '%s'" % (table) result = self.pg.execute(query, async=False, ret=True) Slave.COLUMN_CACHE[table] = list(sum(result, ())) #cache default values query = "SELECT column_name FROM information_schema.columns WHERE \ table_name = '%s' and column_default is not \ null" % (table) result = self.pg.execute(query, async=False, ret=True) #turns list[tuple] into a flat list result = list(sum(result, ())) Slave.DEFAULT_COLUMN_VALUES[table] = result def run(self): self.init() job_queue = multiprocessing.JoinableQueue() tables = self.get_table_listings() for table in tables: job_queue.put(table) processes = [] for i in range(Master.SLAVE_COUNT): process = multiprocessing.Process(target=slave_runner, args=(job_queue,)) process.daemon = True process.start() processes.append(process) #blocks this process until queue reaches 0 job_queue.join() #signal each child process to GTFO for i in range(Master.SLAVE_COUNT): job_queue.put(None) #blocks this process until queue reaches 0 job_queue.join() job_queue.close() #do not end this process until child processes close out for process in processes: process.join() #toodles ! print("this is master saying goodbye") def slave_runner(queue): slave = Slave(queue) slave.do_work() ```

Original source