When using airflow I get "UserWarning: MongoClient opened before fork" + best practice tips to instantiate mongo client in python

Viewed 155

I have created a python flow that executes various steps. I also use MongoDB which for now is used for config purposes but will further use it for persistence. When I execute my code from Pycharm all is good but when I executed it through Airflow I get the fork warning "UserWarning: MongoClient opened before fork. Create MongoClient only after forking.".

At first I thought that the problem is that I open many instances of the Mongo client so I used the singleton pattern so that the connection is instantiated only once -please see code below-. The only way to get rid of the warning is by adding connect=False parameter as you may see commented in my code, but it seems to be a workaround and not a steady solution --according to documentation-.

So what is the problem here? It has something to do with Airflow maybe -since I am not using the mongo_hook-? In addition please let me know if using the singleton pattern to instantiate once the Mongo Client is a good practice? The client is getting called from various modules.

Note that Mongo as well as Airflow are running in separate dockers.

def singleton(class_): instances = {}

def get_instance(*args, **kwargs):
    if class_ not in instances:
        instances[class_] = class_(*args, **kwargs)
    return instances[class_]

return get_instance

@singleton class MongoPersistence(Persistence):

def __init__(self, driver, host, user, password, port, db):
    self._uri = '{driver}://{host}:{port}'.format(driver=driver, host=host, port=port)        
    self._client = MongoClient(self._uri
                         , serverSelectionTimeoutMS=3000  # 3 second timeout
                         , username=self.user
                         , password=self.password
                         # ,connect=False
                               )

    def find_one(self, **kwargs):
    return self._client[kwargs.get('db_name')][kwargs.get('collection_name')].find_one(kwargs.get('query'))

Thank you, Dina

0 Answers
Related