惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Recorded Future
Recorded Future
S
Secure Thoughts
D
Docker
The GitHub Blog
The GitHub Blog
Y
Y Combinator Blog
雷峰网
雷峰网
T
The Blog of Author Tim Ferriss
D
DataBreaches.Net
F
Fortinet All Blogs
月光博客
月光博客
S
Security @ Cisco Blogs
AI
AI
IT之家
IT之家
P
Proofpoint News Feed
Stack Overflow Blog
Stack Overflow Blog
大猫的无限游戏
大猫的无限游戏
N
News and Events Feed by Topic
Hacker News: Ask HN
Hacker News: Ask HN
人人都是产品经理
人人都是产品经理
Recent Announcements
Recent Announcements
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
Google Online Security Blog
Google Online Security Blog
K
Kaspersky official blog
Application and Cybersecurity Blog
Application and Cybersecurity Blog
P
Proofpoint News Feed
Latest news
Latest news
Webroot Blog
Webroot Blog
GbyAI
GbyAI
L
Lohrmann on Cybersecurity
T
Threat Research - Cisco Blogs
S
Schneier on Security
G
Google Developers Blog
N
News | PayPal Newsroom
C
Check Point Blog
T
Tenable Blog
L
LINUX DO - 最新话题
C
CERT Recently Published Vulnerability Notes
T
Troy Hunt's Blog
Know Your Adversary
Know Your Adversary
SecWiki News
SecWiki News
TaoSecurity Blog
TaoSecurity Blog
Recent Commits to openclaw:main
Recent Commits to openclaw:main
爱范儿
爱范儿
W
WeLiveSecurity
Project Zero
Project Zero
M
MIT News - Artificial intelligence
H
Hacker News: Front Page
Apple Machine Learning Research
Apple Machine Learning Research
V
V2EX
C
Cyber Attacks, Cyber Crime and Cyber Security

郑文峰的博客

使用dify对接飞书多维表格 使用n8n对接飞书多维表格 服务启动时出现 OOM 一次服务升级时pg表DDL执行超时失败 Go语言高效IO缓冲技术详解 Go语言延迟初始化(Lazy Initialization)最佳实践 Go语言字符串拼接性能对比与优化指南 Go语言结构体内存对齐完全指南 Go语言空结构体:零内存消耗的高效编程 Go语言堆栈分配与逃逸分析深度解析 Go语言原子操作完全指南 Go语言内存预分配完全指南 Go语言不可变数据共享:无锁并发编程实践 Go语言零拷贝技术完全指南 Go语言遍历性能深度解析:从原理到优化实践 Go语言Interface Boxing原理与性能优化指南 Go协程池深度解析:原理、实现与最佳实践 使用etcd分布式锁导致的协程泄露与死锁问题 基于pre-commit的Python代码规范落地实践 初识 MCP Server pulsar阻塞导致logstash无法接入日志 django-prometheus使用及源码分析 kube-proxy源码分析 kubernetes service如何通过iptables转发 tcp缓存引起的日志丢失 django-apschedule定时任务异常停止 理解calico容器网络通信方案原理 理解flannel的三种容器网络方案原理 理解Linux IPIP隧道 理解VXLAN网络 理解Linux TunTap设备 快速了解iptables kafka中listener和advertised.listeners的作用 django rest_framework 分页 django后端服务、logstash和flink接入VictoriaMetrics指标监控 python中import原理 docker容器单机网络 手动实现docker容器bridge网络模型 mysql之MVCC原理 mysql之日志 使用java开发logstash的filter插件 使用python实现单例模式的三种方式 redis之缓存 redis之分片集群 redis之哨兵机制 redis之主从库同步 redis之持久化 redis之五种基本数据类型 go中如何处理error pod中将代码与运行环境分离 ddt源码分析 python装饰器的使用方法 读书笔记:如何阅读一本书 使用ddt实现unittest的参数化测试 使用kubeadm安装k8s 优化gin表单的错误提示信息 gin中validator模块的源码分析 go简单使用grpc python简单使用grpc k8s之PV、PVC和StorageClass k8s之StatefulSet k8s之DaemonSet k8s之Job和CronJob k8s之ConfigMap和Secret k8s之Service k8s之Pod k8s之Deployment 容器的本质 docker容器 python迭代器与生成器 python元编程 python垃圾回收机制 python上下文管理器 django rest_framework使用jwt django rest_framework异常处理 django rest_framework 自定义文档 django压缩文件下载 django rest_framework使用pytest单元测试 django restframework choice 自定义输出数据 django Filtering 使用 django viewset 和 Router 配合使用时报的错 django model的序列化 django中使用AbStractUser django.core.exceptions.ImproperlyConfigured Application labels aren't unique, duplicates users django 中 media配置 django 外键引用自身和on_delete参数 django 警告 while time zone support is active Flask使用flask_socketio实现websocket flask结合mongo tornado 文件上传 tornado 使用jwt完成用户异步认证 tornado 用户密码 bcrypt加密 tornado 结合wtforms使用表单操作 tornado finish和write区别 tornado 使用peewee-async 完成异步orm数据库操作 pyspark streaming简介 和 消费 kafka示例 使用hue创建ozzie的pyspark action workflow count的性能优化 django rest_framework Authentication 网站
django celery 结合使用
zhengwenfeng · 2022-08-10 · via 郑文峰的博客

# 简介

本文主要介绍django和celery结合使用的案例。

celery 是一个异步任务的调度工具,可以完成一些异步任务和定时任务。

本文使用djcelery来完成django和celery的结合使用。

该案例在github中django_celery_demo (opens new window)

# 流程

任务发布者(Producer)将任务丢到消息队列(Broker)中,任务消费者(worker)从消息代理中获取任务执行,然后将保存存储结果(backend)。

# 消息分发与任务调度的实现机制

# celery-beat

celery 有个定时功能,通过定时去将task丢到broker中,然后worker去执行任务。但是有个确定是,该定时任务必须硬编写到代码中,不可在程序运行中动态增加任务。使用djcelery可以将定时任务写入到数据库中,然后通过操作数据库操作定时任务。

# 案例1

访问接口,异步调用程序中task

# 配置celery

安装**djcelery**

pip install django_celery

在settings中设置celery配置

代码: django_celery_demo/settings.py

import djcelery
djcelery.setup_loader() # 加载djcelery


# 允许的格式
CELERY_ACCEPT_CONTENT = ['pickle', 'json', 'yaml']

BROKER_URL = 'redis://localhost:6379/1' # redis作为中间件
BROKER_TRANSPORT = 'redis'

CELERYBEAT_SCHEDULER = 'djcelery.schedulers.DatabaseScheduler' # 定时任务使用数据库来操作
CELERY_RESULT_BACKEND = 'djcelery.backends.database:DatabaseBackend'  # 结果存储到数据库中

# worker 并发数
CELERY_CONCURRENCY = 2

# 指定导入task任务
CELERY_IMPORTS = {
'tasks.tasks'
}

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20

celery app 配置

代码: django_celery_demo/celery.py

import os

import django
from celery import Celery, shared_task
from celery.schedules import crontab
from celery.signals import task_success
from django.conf import settings
from django.utils import timezone

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'django_celery_demo.settings')
django.setup()

app = Celery('django_celery_demo')
app.config_from_object('django.conf:settings') # celery app 加载 settings中的配置

app.now = timezone.now # 设置时间时区和django一样

# 加载每个django app下的tasks.py中的task任务
app.autodiscover_tasks(lambda: settings.INSTALLED_APPS)

# 这个一个task
@app.task(bind=True)
def debug_task(self):
    print('Request: {0!r}'.format(self.request))
    
# 异步执行这个task
debug_task.delay()

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27

创建djcelery中的表

会自动创建djcelery中的表。里面有保存定时记录、结果记录等等表。

python manage.py migrate

# 在view中异步执行task

在app中创建**add**task

代码: demo/tasks.py

from celery import shared_task

@shared_task(name="add")
def add(a, b):
    return int(a) + int(b)

1
2
3
4
5

创建view去异步执行该task

代码: demo/views.py

from django.http import HttpResponse
from demo.tasks import add as add_task

def add(request):
    a = request.GET["a"]
    b = request.GET["b"]
    add_task.delay(a, b)
    
    return HttpResponse("success")

1
2
3
4
5
6
7
8
9

url中配置view

from demo.views import add

urlpatterns = [
    path('add', add),
]

1
2
3
4
5

运行celery worker

celery -A django_celery_demo worker -l info

运行项目

python manage.py runserver 0:8888

访问接口

http://127.0.0.1:8888/add?a=1&b=2 (opens new window)

结果: 返回success,在worker中可以看到add任务被调用,并且结果是3

# 案例2

定时调用异步任务

# 定时任务简介

有两种定时任务方式,这里使用的是常见的crontab,与linux一毛一样,方便很多。

# 配置

配置和案例1中一样。

# 定时任务

硬编码中创建定时任务

每分钟调用一次add task

代码: django_celery_demo/celery.py

# 这个是硬编码的定时任务
app.conf.beat_schedule = {
    'aa': {
        'task': 'add',
        'schedule': crontab(minute="*/1"),
        'args': (2, 4)
    },
}

1
2
3
4
5
6
7
8

开启celery beat

celery beat -A django_celery_demo -l info

这个服务会将数据库中的定时任务丢到broker

# 案例三-路由

将不同的任务放到不同的队列中,放到不同的worker中。

图: 消息分发与任务调度的实现机制

default = Exchange('default', type="direct")
frequent = Exchange('frequent', type="direct")

CELERY_QUEUES = {
    Queue('default', default, routing_key="default"),
    Queue('frequent', frequent, routing_key="frequent")
}

app.conf.task_default_queue = 'default'

task_routes = {
    'apps.periodic.tasks.oozie_workflow_task': {'queue': 'default'},
    'apps.periodic.tasks.oozie_workflow_status': {'queue': 'custom'}
}

app.conf.beat_schedule = {

    # 每分钟检查oozie运行中的任务状态
    'oozie_workflow_status': {
        'task': 'oozie_workflow_status',
        'schedule': crontab(),
        'args': (2, 1)
    }
}

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24

# 管理worker 进程

使用supervisor来管理worker进程。

# 基本操作

安装

pip install supervisor

生成默认配置文件

echo_supervisord_conf > /etc/supervisor/supervisord.conf

# 命令

supervisorctl

命令 描述
status
reread 读取配置文件
update 加载最新的进程
stop 进程名
start 进程名
reload 重新加载配置

# 配合celery使用

supervisord.conf中添加下面的配置。

[include]
; files = relative/directory/*.ini
files = /home/jim/conf/supervisor/supervisord.conf.d/*.conf

1
2
3

创建配置文件/home/jim/conf/supervisor/supervisord.conf.d/celeryd_worker.conf,添加下面配置

[program:celeryworker]
command=celery -A datahub_poster worker -l info
directory=/home/hadoop/jim/projs/datahub_poster
stdout_logfile=/yun/jim/log/supervisor/celeryworker.log
;stderr_logfile=/yun/jim/log/supervisor/celeryworker_err.log
redirect_stderr=true
autorestart=true
autostart=true
numprocs=1
startsecs=10
stopwaitsecs = 600
priority=15

[program:celerybeat]
command=celery -A datahub_poster beat -l info
directory=/home/hadoop/jim/projs/datahub_poster
stdout_logfile=/yun/jim/log/supervisor/celerybeat.log
;stderr_logfile=/yun/jim/log/supervisor/celerybeat_err.log
redirect_stderr=true
autorestart=true
autostart=true
numprocs=1
startsecs=10
stopwaitsecs = 600
priority=15

[program:celery_flower]
command=celery -A datahub_poster flower --port=5555
directory=/home/hadoop/jim/projs/datahub_poster
stdout_logfile=/yun/jim/log/supervisor/celery_flower.log
;stderr_logfile=/yun/jim/log/supervisor/celery_flower_err.log
redirect_stderr=true
autorestart=true
autostart=true
numprocs=1
startsecs=10
stopwaitsecs = 600
priority=15

[inet_http_server]
port=127.0.0.1:9001

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41

使用配置文件启动supervisor

supervisord -c /etc/supervisor/supervisord.conf

# 问题

  1. 在supervisorctl status时,出现http://localhost:9001 refused connection错误。

解决办法:

在配置文件supervisord.conf中添加

[inet_http_server]
port=127.0.0.1:9001

1
2

然后再update或reload以下。

# 使用flower监控celery

可以通过flower监控celery中的worker、task等等。

安装flower

pip install flower

运行

celery flower --broker=redis://localhost:6379/0

# 持久化

问题: 每次重启flower之后发现,以前的task运行记录清空。

解决: 启动flower时添加 --persistent=True,可以持久化task

# 时区问题

flower会读取celery的时区配置,在项目中配置下面参数即可。

TIME_ZONE = 'Asia/Shanghai'
CELERY_TIMEZONE = TIME_ZONE

1
2