celery-beat
beat如何动态生成task的入参数据?
from celery.beat import BeatLazyFunc
beat_schedule = {
'test-every-5-minutes': {
'task': 'test',
'schedule': 300,
'kwargs': {
"current": BeatCallBack(datetime.datetime.now)
}
}
}
集成gRPC
众所周知, grpc的channel不能跨进程使用, 而worker启动时, 会从worker主进程启动n个worker子进程, 然后分别在子进程 执行task任务。所以在celery作为客户端长连接grpc服务端时,channel应该在worker_process_init信号中连接,即分别在n个子进程 创建和连接channel。 不要在worker_init中定义, 会导致只有1个进程能启动工作,其它进程报错退出。
从代码启动worker进程
https://docs.celeryq.dev/en/latest/userguide/application.html#main-name
chain应用
单元测试
celery有提供pytest的插件协助完成celery的单元测试
pytest_plugins = ('celery.contrib.pytest', )
https://docs.celeryq.dev/en/latest/userguide/testing.html#celery-app-celery-app-used-for-testing
有BUG,有机会参与开源贡献!
https://github.com/celery/celery/issues/7750
我认为应该改进文档提供的样例代码:
def test_create_task(celery_worker):
@celery_worker.app.task
def mul(x, y):
return x * y
assert mul.delay(4, 4).get(timeout=10) == 16
开源贡献
在node2机器上已clone
设置消息过期时间
参考资料:
https://www.rabbitmq.com/ttl.html#per-message-ttl-in-publishers
https://stackoverflow.com/questions/26990438/how-to-set-per-message-expiration-ttl-in-celery
第一种方法:
设置expiration属性, 例如下面设置42秒
my_awesome_task.apply_async(args=(11,), expiration=42)
第二种方法:
定义celery task时传入参数
@shared_task(expires=20)
def import_contacts():
"""expires: 20秒过期时间"""
from datetime import datetime
print("Now: " + str(datetime.now()))
print("Do: import contacts")
除此之外,还可以定义队列属性设置过期时间,以后有空再研究。
常见疑问
数据库重启后, celery beat能否正常工作
经实验, 不影响。当停止mysql数据库时, beat的控制台输出:
django.db.utils.OperationalError: (2013, "Lost connection to server at 'handshake: reading initial communication packet', system error: 0")
django.db.utils.OperationalError: (2002, "Can't connect to server on '127.0.0.1' (10061)")
重新启动mysql后,beat恢复发送任务