【发布时间】:2021-09-04 16:27:05
【问题描述】:
我需要创建一个 celery 组任务,我想等待它完成,但我不清楚文档是如何实现的:
这是我目前的状态:
def import_media(request):
keys = []
for obj in s3_resource.Bucket(env.str('S3_BUCKET')).objects.all():
if obj.key.endswith(('.m4v', '.mp4', '.m4a', '.mp3')):
keys.append(obj.key)
for key in keys:
url = s3_client.generate_presigned_url(
ClientMethod='get_object',
Params={'Bucket': env.str('S3_BUCKET'), 'Key': key},
ExpiresIn=86400,
)
if not Files.objects.filter(descriptor=strip_descriptor_url_scheme(url)).exists():
extract_descriptor.apply_async(kwargs={"descriptor": str(url)})
return None
现在我需要在组内为我拥有的每个 URL 创建一个新任务,我该怎么做?
我现在设法让我的流程像这样工作:
@require_http_methods(("GET"))
def import_media(request):
keys = []
urls = []
for obj in s3_resource.Bucket(env.str('S3_BUCKET')).objects.all():
if obj.key.endswith(('.m4v', '.mp4', '.m4a', '.mp3')):
keys.append(obj.key)
for key in keys:
url = s3_client.generate_presigned_url(
ClientMethod='get_object',
Params={'Bucket': env.str('S3_BUCKET'), 'Key': key},
ExpiresIn=86400,
)
if not Files.objects.filter(descriptor=strip_descriptor_url_scheme(url)).exists():
new_file = Files.objects.create(descriptor=strip_descriptor_url_scheme(url))
new_file.save()
urls.append(url)
workflow = (
group([extract_descriptor.s(url) for url in urls]).delay()
)
workflow.get(timeout=None, interval=0.5)
print("hello - Further processing here")
return None
有什么优化的建议吗?至少现在它工作得很好!
提前致谢
【问题讨论】:
标签: python django celery django-celery