Is there a workaround for the blocking that happens with Firebase Python SDK? Like adding a completion callback?

Viewed 950

Recently, I have moved my REST server code in express.js to using FastAPI. So far, I've been successful in the transition until recently. I've noticed based on the firebase python admin sdk documention, unlike node.js, the python sdk is blocking. The documentation says here:

In Python and Go Admin SDKs, all write methods are blocking. That is, the write methods do not return until the writes are committed to the database.

I think this feature is having a certain effect on my code. It also could be how I've structured my code as well. Some code from one of my files is below:

from app.services.new_service import nService
from firebase_admin import db
import json
import redis

class TryNewService:
  async def tryNew_func(self, request):
    # I've already initialized everything in another file for firebase
    ref = db.reference()
    r = redis.Redis()
    holdingData = await nService().dialogflow_session(request)
    fulfillmentText = json.dumps(holdingData[-1])
    body = await request.json()

    if ("user_prelimInfo_address" in holdingData):
      holdingData.append("session")
      holdingData.append(body["session"])
      print(holdingData)
      return(holdingData)
    else:
      if (("Default Welcome Intent" in holdingData)):
        pass
      else:
        UserVal = r.hget(name='{}'.format(body["session"]), key="userId").decode("utf-8")
        ref.child("users/{}".format(UserVal)).child("c_data").set({holdingData[0]:holdingData[1]})
        print(holdingData)
      return(fulfillmentText)

Is there any workaround for the blocking effect of usingref.set() line in my code? Kinda like adding a callback in node.js? I'm new to the asyncio world of python 3.


Update as of 06/13/2020: So I added following code and am now getting a RuntimeError: Task attached to a different loop. In my second else statement I do the following:

loop = asyncio.new_event_loop()
UserVal = r.hget(name='{}'.format(body["session"]), key="userId").decode("utf-8")
with concurrent.futures.ThreadPoolExecutor(max_workers=20) as pool:
  result = await loop.run_in_executor(pool, ref.child("users/{}".format(UserVal)).child("c_data").set({holdingData[0]:holdingData[1]}))
  print("custom thread pool:{}".format(result))

With this new RuntimeError, I would appreciate some help in figuring out.

2 Answers

If you want to run synchronous code inside an async coroutine, then the steps are:

  1. loop = get_event_loop()
    Note: Get and not new. Get provides current event_loop, and new_even_loop returns a new one
  2. await loop.run_in_executor(None, sync_method)
    • First parameter = None -> use default executor instance
    • Second parameter (sync_method) is the synchronous code to be called.

Remember that resources used by sync_method need to be properly synchronized:

  • a) either using asyncio.Lock
  • b) or using asyncio.run_coroutine_threadsafe function(see an example below)

Forget for this case about ThreadPoolExecutor (that provides a way to I/O parallelism, versus concurrency provided by asyncio).

You can try following code:

loop = asyncio.get_event_loop()
UserVal = r.hget(name='{}'.format(body["session"]), key="userId").decode("utf-8")
result = await loop.run_in_executor(None, sync_method, ref, UserVal, holdingData)
print("custom thread pool:{}".format(result))

With a new function:

def sync_method(ref, UserVal, holdingData):
    result = ref.child("users/{}".format(UserVal)).child("c_data").set({holdingData[0]:holdingData[1]}))
    return result

Please let me know your feedback

Note: previous code it's untested. I have only tested next minimum example (using pytest & pytest-asyncio):

import asyncio
import time

import pytest


@pytest.mark.asyncio
async def test_1():
    loop = asyncio.get_event_loop()
    delay = 3.0
    result = await loop.run_in_executor(None, sync_method, delay)
    print(f"Result = {result}")

def sync_method(delay):
    time.sleep(delay)
    print(f"dddd {delay}")
    return "OK"

Answer @jeff-ridgeway comment:

Let's try to change previous answer to clarify how to use run_coroutine_threadsafe, to execute from a sync worker thread a coroutine that gather these shared resources:

  1. Add loop as additional parameter in run_in_executor
  2. Move all shared resources from sync_method to a new async_method, that is executed with run_coroutine_threadsafe
loop = asyncio.get_event_loop()
UserVal = r.hget(name='{}'.format(body["session"]), key="userId").decode("utf-8")
result = await loop.run_in_executor(None, sync_method, ref, UserVal, holdingData, loop)
print("custom thread pool:{}".format(result))

def sync_method(ref, UserVal, holdingData, loop):
    coro = async_method(ref, UserVal, holdingData)
    future = asyncio.run_coroutine_threadsafe(coro, loop)
    future.result()
async def async_method(ref, UserVal, holdingData)
    result = ref.child("users/{}".format(UserVal)).child("c_data").set({holdingData[0]:holdingData[1]}))
    return result

Note: previous code is untested. And now my tested minimum example updated:

@pytest.mark.asyncio
async def test_1():
    loop = asyncio.get_event_loop()
    delay = 3.0
    result = await loop.run_in_executor(None, sync_method, delay, loop)
    print(f"Result = {result}")

def sync_method(delay, loop):
    coro = async_method(delay)
    future = asyncio.run_coroutine_threadsafe(coro, loop)
    return future.result()

async def async_method(delay):
    time.sleep(delay)
    print(f"dddd {delay}")
    return "OK"

I hope this can be helpful

Related