Kiến trúc và triển khai Celery cho hệ thống Python

Tổng quan về Celery

Celery là giải pháp hàng đầu để xây dựng các ứng dụng cần xử lý công việc nền hiệu quả. Đây là một hệ thống phân tán được thiết kế nhằm xử lý quy trình truyền tải thông tin giữa các dịch vụ một cách ổn định. Trọng tâm chính của nó là quản lý hàng đợi tác vụ (task queue) theo cơ chế bất đồng bộ, đồng thời hỗ trợ mạnh mẽ tính năng lên lịch thực thi tự động.

Hình dưới đây minh họa luồng dữ liệu cơ bản trong kiến trúc của Celery:

Celery Architecture Diagram

Trường hợp sử dụng phổ biến

  • Thực thi tác vụ không đồng bộ: Chuyển giao các thao tác tốn thời gian sang nền để không làm gián đoạn luồng chính của ứng dụng. Ví dụ điển hình bao gồm gửi thông báo qua SMS/email, đẩy cảnh báo đến người dùng, hoặc chuyển đổi định dạng media.
  • Lập lịch trình công việc: Thực hiện các tác vụ định kỳ tương tự như Cronjob trên Linux, chẳng hạn như thu thập thống kê dữ liệu cuối ngày hoặc dọn dẹp bảng ghi tạm thời.

Thành phần trung gian (Message Broker)

Celery cần một môi trường trung gian để lưu trữ và truyền tiếp yêu cầu từ người gọi đến worker. Hai lựa chọn phổ biến nhất hiện nay là RabbitMQ hoặc Redis.

Cài đặt và Khởi tạo

Để bắt đầu, ta cần cài đặt gói Celery kèm theo driver Redis thông qua pip:

pip install "celery[redis]"

Xử lý tác vụ bất đồng bộ

Để hiểu rõ sự khác biệt, hãy xem xét một hàm xử lý đơn giản có độ trễ giả lập.

Phiên bản trực tiếp (Blocking):

# -*- coding: utf-8 -*-
import time
import sys

def process_simple_task(a, b):
    print('Đang chạy tác vụ')
    # Giả lập xử lý nặng trong 4 giây
    time.sleep(4)
    return a + b

if __name__ == '__main__':
    print('Bắt đầu chương trình')
    result = process_simple_task(2, 8)
    print('Kết thúc chờ')
    print(result)

Kết quả chạy ra sẽ hiển thị sự chờ đợi dài hơi trước khi in ra kết quả cuối cùng.

Phiên bản tối ưu hóa bằng Celery:

Tạo file worker_tasks.py chứa định nghĩa tác vụ:

# -*- coding: utf-8 -*-
import time
from celery import Celery

# Cấu hình kết nối tới broker và nơi lưu kết quả
broker_url = 'redis://127.0.0.1:6379/0'
result_backend = 'redis://127.0.0.1:6379/1'

celery_app = Celery('project_tasks', broker=broker_url, backend=result_backend)

@celery_app.task(bind=True)
def complex_calculation(x, y):
    print('Đang xử lý trong background')
    time.sleep(5)  # Thời gian xử lý giả lập
    return x * y

Tạo file client_main.py để gọi tác vụ:

# -*- coding: utf-8 -*-
from worker_tasks import complex_calculation

if __name__ == '__main__':
    print('Gửi lệnh đi...')
    # Phương thức delay đẩy tác vụ vào hàng đợi ngay lập tức
    async_result = complex_calculation.delay(10, 5)
    
    print('Đã trả về ID nhiệm vụ')
    print(f'Task ID: {async_result.id}')

Lúc này, chương trình chính không bị treo. Nó nhận về một ID duy nhất ngay lập tức. Tác vụ thực tế sẽ nằm im trong Redis chờ worker xử lý.

Vận hành Worker

Để các tác vụ được thực thi, bạn cần khởi động quy trình worker tại thư mục chứa mã nguồn:

celery -A worker_tasks worker --loglevel=INFO
  • -A: Chỉ định module chứa ứng dụng Celery.
  • --loglevel=INFO: Thiết lập mức độ chi tiết của nhật ký hệ thống.

Nếu phiên bản quá mới, có thể gặp lỗi tương thích. Khi chạy thành công, terminal sẽ hiển thị danh sách tác vụ đã đăng ký và cấu hình kết nối Redis.

Cấu trúc dự án chuẩn

Để dễ bảo trì, nên tách biệt cấu hình, logic tác vụ và điểm vào của ứng dụng.

1. File cấu hình (celery_config.py)

# -*- coding: utf-8 -*-
class Config:
    # Địa chỉ broker
    CELERY_BROKER_URL = 'redis://localhost:6379/3'
    # Nơi lưu kết quả sau khi hoàn thành
    CELERY_RESULT_BACKEND = 'redis://localhost:6379/4'
    # Cài đặt múi giờ mặc định
    CELERY_TIMEZONE = 'Asia/Ho_Chi_Minh'
    
    # Import các module chứa task để tự động đăng ký
    CELERY_IMPORTS = ('tasks.module_processing',)

    # Cấu hình lịch trình (Beat Scheduler)
    CELERY_BEAT_SCHEDULE = {
        'daily-stats-report': {
            'task': 'tasks.module_processing.generate_report',
            'schedule': 60.0,  # Chạy mỗi 60 giây
        },
    }

2. Module Tác vụ (tasks/module_processing.py)

# -*- coding: utf-8 -*-
from my_project.celery_init import app

@app.task
def generate_report():
    # Logic thống kê
    return {'status': 'success'}

3. Khởi tạo ứng dụng (my_project/celery_init.py)

# -*- coding: utf-8 -*-
from celery import Celery

app = Celery('my_project')

# Tải cấu hình từ file riêng
app.config_from_object('celery_config.Config')

# Tự động quét các package chứa task
app.autodiscover_tasks()

4. Điểm kích hoạt (run.py)

# -*- coding: utf-8 -*-
from tasks.module_processing import generate_report

if __name__ == '__main__':
    # Kích hoạt tác vụ không đồng bộ
    generate_report.apply_async()
    print('Yêu cầu đã gửi vào hàng đợi')

Bằng cách chia nhỏ các tệp như trên, việc mở rộng thêm chức năng hay thay đổi cấu hình kết nối sẽ trở nên đơn giản hơn mà không ảnh hưởng đến logic nghiệp vụ cốt lõi.

Thẻ: python celery Redis distributed-tasks message-queue

Đăng vào ngày 22 tháng 8 lúc 07:42