基于Celery+Redis的OCR未知文档类型识别多进程改造技术咨询
Great questions—let's tackle each one with your specific OCR application in mind, plus share actionable code snippets to refactor your calc_receipts() function.
1. Do I need to start a Celery Worker inside calc_receipts()?
Absolutely not. Celery Workers are long-running processes meant to be started independently (not from within your business logic function). Starting a Worker inside calc_receipts() would lead to duplicate Worker instances, resource leaks, and unpredictable behavior.
Instead, start your Worker(s) separately before running your application—either via the command line:
celery -A your_celery_app.tasks worker --loglevel=info --concurrency=4
Or via a Python script (for automation/containerized environments):
from celery import Celery app = Celery('ocr_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') if __name__ == '__main__': app.worker_main(['--loglevel=info', '--concurrency=4'])
Your calc_receipts() function only needs to submit tasks to the already-running Workers.
2. Can I start a Celery Client via Python code instead of the console?
Yes, absolutely. The Celery Client is just a Python object you initialize in your code, configured to talk to your Redis broker. Here's how to set it up for your workflow:
First, define your Celery task (move your calc_receipt logic into a task):
# tasks.py from celery import Celery from your_ocr_modules import OpenCV, Tesseract, DocumentType app = Celery('ocr_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') @app.task def validate_receipt_task(raw_img_path, dt): aligned_img_path = OpenCV.align_img( template_path=f"path/to/templates/{dt.name.lower()}.png", # Adjust path logic as needed raw_img_path=raw_img_path, result_img_path=f"path/to/aligned/{dt.name.lower()}_{raw_img_path.split('/')[-1]}", ) tesseract_result = Tesseract.read_from_img(img_path=aligned_img_path) if tesseract_result: return aligned_img_path, dt.value # Store enum value instead of enum instance for serialization return '', DocumentType.NONE.value
Then, in your main application code, initialize the Celery client and use it to submit tasks directly.
3. How to manage Redis results and wait for all 4 tasks to complete?
Use Celery's group primitive to submit all 4 validation tasks in parallel, then wait for all results to come back. You can then filter the results to find the valid document type.
Here's how to refactor your calc_receipts() function:
from celery import group from tasks import validate_receipt_task from your_ocr_modules import DocumentType, map_receipt_to_receipts def calc_receipts(self, raw_img_path, selected_doc_type): # Get all document types to validate doc_types = map_receipt_to_receipts[selected_doc_type] # Create a group of parallel tasks job = group( validate_receipt_task.s(raw_img_path, dt) for dt in doc_types ) # Execute the group and wait for all results results = job.apply_async() results.wait() # Blocks until all tasks finish # Process results to find the valid document type for result in results.get(): aligned_path, dt_value = result if dt_value != DocumentType.NONE.value: # Convert back to enum instance valid_dt = DocumentType(dt_value) return aligned_path, valid_dt # If no valid type found return '', DocumentType.NONE
Celery handles storing results in Redis automatically (thanks to the backend configuration), and cleans up old results based on your result_expires setting (default is 1 day). You don't need to manually clean Redis data unless you want to override the default expiration timeline.
4. General best practices for Celery + Redis in this scenario
- Configure result expiration: Set
result_expiresin your Celery app config to auto-clean old results and avoid Redis bloat:app.conf.result_expires = 3600 # Expire results after 1 hour - Make tasks idempotent: Ensure running
validate_receipt_taskmultiple times with the same inputs doesn't cause issues (e.g., use unique filenames for aligned images to prevent overwrites). - Set task timeouts: Prevent stuck tasks from hanging by adding a timeout to your task:
@app.task(time_limit=30, soft_time_limit=25) # Hard kill after 30s, warn after 25s def validate_receipt_task(raw_img_path, dt): # ... your logic ... - Scale Workers appropriately: Match the number of Worker concurrency slots to your available CPU cores (e.g.,
--concurrency=4for a 4-core machine) to maximize parallelism without overloading. - Avoid passing large objects: Don't pass raw image data in tasks—stick to file paths (like you're already doing) to keep task messages small and efficient.
- Monitor tasks: Use Celery's Flower tool to track task status, worker health, and performance—it's invaluable for debugging slow or failed tasks.
内容的提问来源于stack exchange,提问作者Slobjo

