Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

how to create per worker process variable #4021

starplanet opened this issue May 8, 2017 · 1 comment

how to create per worker process variable #4021

starplanet opened this issue May 8, 2017 · 1 comment


Copy link

@starplanet starplanet commented May 8, 2017


  • I have included the output of celery -A proj report in the issue.
    (if you are not able to do this, then at least specify the Celery
    version affected).
  • I have verified that the issue exists against the master branch of Celery.

Steps to reproduce

This problem comes from producer.flush make celery hang.

I paste the code in the following:

from celery import Celery
from kafka import KafkaProducer

app = Celery('test', broker='redis://')

producer = KafkaProducer(bootstrap_servers=['', ''])

def send_msg():
    # producer = KafkaProducer(bootstrap_servers=['', ''])
    for i in range(10):
        producer.send('test', b'this is the %dth test message' % i)

if __name__ == '__main__':

I want to create producer variable per worker process, and I think worker_process_init signal will help.

But I don't know how to declare producer variable per worker process which will be then used in task func.

Can someone help me? Thanks.


This comment has been minimized.

Copy link

@auvipy auvipy commented Jan 11, 2018

check the celery docs in detail. also I kafka is not supported now

@auvipy auvipy closed this Jan 11, 2018
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment
None yet
2 participants
You can’t perform that action at this time.