prefetch_multiplier_tuning
Ajuste fino del prefetch multiplier de los workers Celery para equilibrar throughput y latencia segun el tipo de tarea del pipeline de verificacion de identidad. Las tareas GPU-bound (face_match, liveness) requieren configuracion diferente a las CPU-bound (OCR, doc_processing) para evitar cuellos de botella y maximizar la utilizacion de recursos.
When to use
Usa esta skill cuando trabajes con el worker_pool_agent y necesites optimizar el rendimiento de los workers Celery en el pipeline KYC. Aplica cuando observes alta latencia en tareas GPU-bound, subutilizacion de workers CPU-bound, o desbalance de carga entre colas de verificacion.
Instructions
Identificar y clasificar las tareas del pipeline por tipo de recurso consumido:
# backend/modules/worker_pool/task_classification.py
TASK_RESOURCE_MAP = {
"face_match": "gpu_bound", # ArcFace inference
"liveness_detection": "gpu_bound", # Anti-spoofing model
"ocr_extraction": "cpu_bound", # PaddleOCR / EasyOCR
"doc_processing": "cpu_bound", # OpenCV processing
"antifraud_analysis": "cpu_bound", # ELA, metadata analysis
"decision_engine": "io_bound", # DB queries, scoring
}
Configurar prefetch multiplier diferenciado por tipo de worker:
# backend/modules/worker_pool/celery_config.py
# GPU-bound workers: prefetch=1 para evitar acaparar tareas mientras GPU esta ocupada
GPU_WORKER_CONFIG = {
"worker_prefetch_multiplier": 1,
"task_acks_late": True,
"worker_concurrency": 1, # Una tarea GPU a la vez
}
# CPU-bound workers: prefetch=4 para mantener pipeline lleno
CPU_WORKER_CONFIG = {
"worker_prefetch_multiplier": 4,
"task_acks_late": True,
"worker_concurrency": 4, # Multiples tareas CPU en paralelo
}
# IO-bound workers: prefetch=8 para alto throughput
IO_WORKER_CONFIG = {
"worker_prefetch_multiplier": 8,
"task_acks_late": False,
"worker_concurrency": 8,
}
Crear scripts de lanzamiento de workers con configuracion especifica por cola:
# GPU workers - face_match y liveness
celery -A worker_pool worker \
--queues=face_match,liveness \
--concurrency=1 \
--prefetch-multiplier=1 \
--hostname=gpu-worker@%h
# CPU workers - OCR y document processing
celery -A worker_pool worker \
--queues=ocr,doc_processing,antifraud \
--concurrency=4 \
--prefetch-multiplier=4 \
--hostname=cpu-worker@%h
Implementar metricas para medir el impacto del prefetch multiplier:
# backend/modules/worker_pool/metrics.py
from prometheus_client import Histogram, Gauge
task_wait_time = Histogram(
"kyc_task_wait_seconds",
"Tiempo de espera en cola antes de ejecucion",
["task_name", "worker_type"],
)
prefetch_utilization = Gauge(
"kyc_prefetch_utilization",
"Ratio de tareas prefetched vs ejecutadas",
["worker_type"],
)
Implementar ajuste dinamico del prefetch multiplier basado en carga:
# backend/modules/worker_pool/dynamic_prefetch.py
from celery.signals import worker_ready
import psutil
@worker_ready.connect
def adjust_prefetch(sender, **kwargs):
cpu_percent = psutil.cpu_percent(interval=1)
current_multiplier = sender.app.conf.worker_prefetch_multiplier
if cpu_percent > 85 and current_multiplier > 1:
sender.app.conf.worker_prefetch_multiplier = max(1, current_multiplier - 1)
elif cpu_percent < 40 and current_multiplier < 8:
sender.app.conf.worker_prefetch_multiplier = current_multiplier + 1
Configurar task_acks_late junto con el prefetch para garantizar que tareas KYC no se pierdan:
# Combinacion critica para tareas de verificacion
app.conf.update(
task_acks_late=True, # ACK solo despues de completar
task_reject_on_worker_lost=True, # Reencolar si worker muere
worker_prefetch_multiplier=1, # Default conservador
)
Crear tests de carga para validar la configuracion de prefetch:
# backend/tests/test_prefetch_tuning.py
def test_gpu_worker_no_starvation():
"""Verificar que GPU workers no acaparan tareas."""
results = []
for _ in range(10):
r = face_match_task.delay(test_session_id)
results.append(r)
wait_times = [r.result["wait_time"] for r in results]
assert max(wait_times) < 8.0 # Dentro del SLA de 8 segundos
Notes
- Un prefetch_multiplier=1 es obligatorio para tareas GPU-bound del pipeline KYC (face_match, liveness); valores mayores causan que un worker acapare tareas mientras la GPU esta ocupada, aumentando la latencia global.
- Siempre combinar prefetch_multiplier con task_acks_late=True para tareas criticas de verificacion; esto evita perder tareas si un worker muere durante el procesamiento biometrico.
- Revisar periodicamente las metricas de wait_time por cola para detectar desajustes; el objetivo es mantener el tiempo total de verificacion por debajo de 8 segundos segun los SLA del sistema.
1---2name: prefetch-multiplier-tuning3description: Ajuste del prefetch multiplier de Celery workers para optimizar throughput vs latencia en tareas KYC4---56# prefetch_multiplier_tuning78Ajuste fino del prefetch multiplier de los workers Celery para equilibrar throughput y latencia segun el tipo de tarea del pipeline de verificacion de identidad. Las tareas GPU-bound (face_match, liveness) requieren configuracion diferente a las CPU-bound (OCR, doc_processing) para evitar cuellos de botella y maximizar la utilizacion de recursos.910## When to use1112Usa esta skill cuando trabajes con el **worker_pool_agent** y necesites optimizar el rendimiento de los workers Celery en el pipeline KYC. Aplica cuando observes alta latencia en tareas GPU-bound, subutilizacion de workers CPU-bound, o desbalance de carga entre colas de verificacion.1314## Instructions15161. Identificar y clasificar las tareas del pipeline por tipo de recurso consumido:17 ```python18 # backend/modules/worker_pool/task_classification.py19 TASK_RESOURCE_MAP = {20 "face_match": "gpu_bound", # ArcFace inference21 "liveness_detection": "gpu_bound", # Anti-spoofing model22 "ocr_extraction": "cpu_bound", # PaddleOCR / EasyOCR23 "doc_processing": "cpu_bound", # OpenCV processing24 "antifraud_analysis": "cpu_bound", # ELA, metadata analysis25 "decision_engine": "io_bound", # DB queries, scoring26 }27 ```28292. Configurar prefetch multiplier diferenciado por tipo de worker:30 ```python31 # backend/modules/worker_pool/celery_config.py32 # GPU-bound workers: prefetch=1 para evitar acaparar tareas mientras GPU esta ocupada33 GPU_WORKER_CONFIG = {34 "worker_prefetch_multiplier": 1,35 "task_acks_late": True,36 "worker_concurrency": 1, # Una tarea GPU a la vez37 }3839 # CPU-bound workers: prefetch=4 para mantener pipeline lleno40 CPU_WORKER_CONFIG = {41 "worker_prefetch_multiplier": 4,42 "task_acks_late": True,43 "worker_concurrency": 4, # Multiples tareas CPU en paralelo44 }4546 # IO-bound workers: prefetch=8 para alto throughput47 IO_WORKER_CONFIG = {48 "worker_prefetch_multiplier": 8,49 "task_acks_late": False,50 "worker_concurrency": 8,51 }52 ```53543. Crear scripts de lanzamiento de workers con configuracion especifica por cola:55 ```bash56 # GPU workers - face_match y liveness57 celery -A worker_pool worker \58 --queues=face_match,liveness \59 --concurrency=1 \60 --prefetch-multiplier=1 \61 --hostname=gpu-worker@%h6263 # CPU workers - OCR y document processing64 celery -A worker_pool worker \65 --queues=ocr,doc_processing,antifraud \66 --concurrency=4 \67 --prefetch-multiplier=4 \68 --hostname=cpu-worker@%h69 ```70714. Implementar metricas para medir el impacto del prefetch multiplier:72 ```python73 # backend/modules/worker_pool/metrics.py74 from prometheus_client import Histogram, Gauge7576 task_wait_time = Histogram(77 "kyc_task_wait_seconds",78 "Tiempo de espera en cola antes de ejecucion",79 ["task_name", "worker_type"],80 )81 prefetch_utilization = Gauge(82 "kyc_prefetch_utilization",83 "Ratio de tareas prefetched vs ejecutadas",84 ["worker_type"],85 )86 ```87885. Implementar ajuste dinamico del prefetch multiplier basado en carga:89 ```python90 # backend/modules/worker_pool/dynamic_prefetch.py91 from celery.signals import worker_ready92 import psutil9394 @worker_ready.connect95 def adjust_prefetch(sender, **kwargs):96 cpu_percent = psutil.cpu_percent(interval=1)97 current_multiplier = sender.app.conf.worker_prefetch_multiplier9899 if cpu_percent > 85 and current_multiplier > 1:100 sender.app.conf.worker_prefetch_multiplier = max(1, current_multiplier - 1)101 elif cpu_percent < 40 and current_multiplier < 8:102 sender.app.conf.worker_prefetch_multiplier = current_multiplier + 1103 ```1041056. Configurar task_acks_late junto con el prefetch para garantizar que tareas KYC no se pierdan:106 ```python107 # Combinacion critica para tareas de verificacion108 app.conf.update(109 task_acks_late=True, # ACK solo despues de completar110 task_reject_on_worker_lost=True, # Reencolar si worker muere111 worker_prefetch_multiplier=1, # Default conservador112 )113 ```1141157. Crear tests de carga para validar la configuracion de prefetch:116 ```python117 # backend/tests/test_prefetch_tuning.py118 def test_gpu_worker_no_starvation():119 """Verificar que GPU workers no acaparan tareas."""120 results = []121 for _ in range(10):122 r = face_match_task.delay(test_session_id)123 results.append(r)124 wait_times = [r.result["wait_time"] for r in results]125 assert max(wait_times) < 8.0 # Dentro del SLA de 8 segundos126 ```127128## Notes129130- Un prefetch_multiplier=1 es obligatorio para tareas GPU-bound del pipeline KYC (face_match, liveness); valores mayores causan que un worker acapare tareas mientras la GPU esta ocupada, aumentando la latencia global.131- Siempre combinar prefetch_multiplier con task_acks_late=True para tareas criticas de verificacion; esto evita perder tareas si un worker muere durante el procesamiento biometrico.132- Revisar periodicamente las metricas de wait_time por cola para detectar desajustes; el objetivo es mantener el tiempo total de verificacion por debajo de 8 segundos segun los SLA del sistema.