Skip to content

Add a concurrency model with ThreadPoolExecutor#5011

Closed
alfred-sa wants to merge 1 commit into
celery:masterfrom
alfred-sa:master
Closed

Add a concurrency model with ThreadPoolExecutor#5011
alfred-sa wants to merge 1 commit into
celery:masterfrom
alfred-sa:master

Conversation

@alfred-sa

@alfred-sa alfred-sa commented Aug 29, 2018

Copy link
Copy Markdown
Contributor

Hello,

I don't know if it would be useful to someone else, but I needed to implement a concurrency model based on thread (and not processes) in celery because I wanted to pass future objects between tasks and some coroutines and it is not pickable.

Of course it is not scale out and it has limitations (because of the GIL), but it can be useful for people who wants the share memory between tasks and another thread (for example the asyncio event loop) without blocking as the 'solo' concurrency model.

I am open to comment. Maybe this model won't be needed anymore in celery 5 as it could be replaced by an asyncio loop, which is not possible in celery 4.
And I would be happy to help on a massive asyncio refactoring for Celery 5.

Regards

@thedrow

thedrow commented Aug 29, 2018

Copy link
Copy Markdown
Contributor

I've been meaning to do this myself. I'll review this as soon as possible.

@thedrow thedrow left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This PR requires some adjustments before we can merge it. Some of them are outlined in the diff itself.
We also require:

  • Unit tests to ensure this works correctly.
  • A requirements file for Python 2.7 users in requirements/extra/ is necessary. Please ensure to provide the appropriate version markers.
  • Documentation adjustments.

signal_safe = False

def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

super() calls are only valid with Python 3.
Celery 4.x still supports Python 2.7 so we'll need to adjust that.


def on_stop(self):
self.executor.shutdown()
super().on_stop()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

super() calls are only valid with Python 3.
Celery 4.x still supports Python 2.7 so we'll need to adjust that.

def _get_info(self):
return {
'max-concurrency': self.limit,
'threads': len(self.executor._threads)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm not a big fan of using private APIs. Can we change this to something more sensible?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The actual number of threads in the executor is not accessible from a public API. So I can do 2 things: either modify concurrent.futures/thread.py to create one (modification to a core python lib) or remove de actual number of threads info in the concurrency model (ie remove line 39).

The former is not sure to be accepted by the PSF. I already have a pull request pending for Python 3.8, so I can try anyway ?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Seems like you are correct. There's no public API for that.
I think a TODO comment to change that once it lands is sufficient.

@alfred-sa

Copy link
Copy Markdown
Contributor Author

Thank you for the review, I will apply your recommandation in the next week.

@thedrow

thedrow commented Oct 8, 2018

Copy link
Copy Markdown
Contributor

@whuji Did you get a chance to work on this?

@alfred-sa

Copy link
Copy Markdown
Contributor Author

Not yet. I am doing it right now.

@alfred-sa

alfred-sa commented Oct 8, 2018

Copy link
Copy Markdown
Contributor Author

I have submitted a new RP : #5099

Duplicate of #5099

@alfred-sa alfred-sa closed this Oct 8, 2018
@alfred-sa
alfred-sa deleted the master branch October 8, 2018 15:25
@thedrow thedrow removed this from the v4.3 milestone Oct 9, 2018
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants