Procesar trabajos de cola sin efectos duplicados

AWSBeginner
Practicar Ahora

Introducción

Los trabajos repetidos para el mismo pedido deben conservar el resultado de negocio completado. Añadirás una protección de escritura atómica a un consumidor proporcionado, probarás pedidos repetidos e independientes y confirmarás que los mensajes se confirman correctamente.

Completa primero Aislar trabajos fallidos con una cola de mensajes fallidos, Evitar pedidos duplicados con escrituras condicionales y Configurar y diagnosticar una función Lambda. Esta VM independiente proporciona el procesador inicial, el rol y una tabla de pedidos vacía.

Relación con las certificaciones

Este laboratorio ofrece práctica para los siguientes temas de examen.

Conectar una cola vacía al procesador

En este paso, crea una cola y conecta el consumidor proporcionado antes de enviar trabajos de negocio.

Usa AWS View junto a Terminal para comparar las colas actuales, los resultados del procesador y los pedidos almacenados. Conserva los datos de referencia ajenos al trabajo.

cd /home/labex/project

Crea la cola Standard. La sustitución de comandos, $(...), guarda la URL de cola devuelta para operaciones posteriores:

QUEUE_URL=$(aws sqs create-queue --queue-name labex-q05-jobs --attributes VisibilityTimeout=30 --query QueueUrl --output text)

Selecciona su ARN y crea una asignación de origen de eventos con tamaño de lote uno:

QUEUE_ARN=$(aws sqs get-queue-attributes --queue-url "$QUEUE_URL" --attribute-names QueueArn --query Attributes.QueueArn --output text)
MAPPING_ID=$(aws lambda create-event-source-mapping --function-name labex-q05-worker --event-source-arn "$QUEUE_ARN" --batch-size 1 --enabled --query UUID --output text)

Inspecciona la conexión:

aws lambda get-event-source-mapping --uuid "$MAPPING_ID" --query '{Source:EventSourceArn,Function:FunctionArn,Batch:BatchSize,State:State}'

Espera el origen labex-q05-jobs, labex-q05-worker, el tamaño de lote 1 y el estado Enabled. El rol proporcionado solo puede sondear y eliminar mensajes de esta cola y escribir en la tabla de pedidos, además de registrar ejecuciones. AWS View muestra una cola vacía, sin pedidos ni decisiones de escritura. No envíes un trabajo hasta que despliegues la protección en el siguiente paso.

Desplegar una protección atómica por clave de negocio

En este paso, sustituye la escritura incondicional del procesador y demuestra que el primer trabajo protegido se completa correctamente.

ID de mensaje y clave de negocio

Distintos mensajes de cola pueden contener el mismo id de negocio. El procesador confirma un duplicado coincidente sin volver a escribir el pedido.

Un consumidor idempotente conserva el mismo resultado de negocio cuando llega trabajo repetido. DynamoDB evalúa ConditionExpression='attribute_not_exists(id)' en la misma operación de escritura que PutItem. Esto evita una condición de carrera entre una lectura y una escritura separadas. El id de negocio es la clave; usar el messageId de SQS permitiría que una copia del mismo pedido enviada de nuevo por el productor volviera a escribir.

Cuando falla una condición, ReturnValuesOnConditionCheckFailure='ALL_OLD' devuelve el elemento existente. El código captura únicamente ConditionalCheckFailedException y compara ese elemento con el elemento solicitado. Los detalles idénticos indican un duplicado ya completado; los detalles en conflicto u otro error siguen provocando un fallo, en lugar de confirmarse silenciosamente.

Escribe el procesador completo siguiente. cat > app.py sustituye el archivo y el documento de entrada proporciona su contenido hasta el terminador PY. Las comillas de 'PY' impiden la expansión del shell dentro del código Python:

cat > app.py <<'PY'
import json
import os
import boto3
from botocore.exceptions import ClientError


def handler(event, context):
    print('EVENT '+json.dumps(event,sort_keys=True))
    if set(event)!={'Records'} or len(event['Records'])!=1:
        raise ValueError('Expected one SQS record')
    job=json.loads(event['Records'][0]['body'])
    print('JOB '+json.dumps(job,sort_keys=True))
    if set(job)!={'id','quantity'} or not isinstance(job['id'],str):
        raise ValueError('Use an id and quantity')
    quantity=job['quantity']
    if isinstance(quantity,bool) or not isinstance(quantity,int) or not 1<=quantity<=10:
        raise ValueError('Quantity must be an integer from 1 to 10')
    database=boto3.client('dynamodb',endpoint_url='http://127.0.0.1:5000',region_name='us-east-1')
    item={'id':{'S':job['id']},'quantity':{'N':str(quantity)},'total_cents':{'N':str(quantity*250+100)}}
    try:
        database.put_item(TableName=os.environ['TABLE_NAME'],Item=item,
                          ConditionExpression='attribute_not_exists(id)',
                          ReturnValuesOnConditionCheckFailure='ALL_OLD')
        result={'id':job['id'],'quantity':quantity,'processed':True,'duplicate':False}
    except ClientError as error:
        if error.response['Error']['Code']!='ConditionalCheckFailedException':
            raise
        if error.response.get('Item')!=item:
            raise ValueError('Order ID already exists with different details') from error
        result={'id':job['id'],'quantity':quantity,'processed':False,'duplicate':True}
    print('RESULT '+json.dumps(result,sort_keys=True))
    return result
PY

El análisis del evento y las comprobaciones de cantidad son contexto proporcionado. El cambio clave es el único put_item condicional, el tratamiento limitado de su fallo y el resultado processed/duplicate. Un duplicado devuelve un resultado correcto para que el consumidor pueda confirmar su copia SQS después de que el resultado de negocio ya se haya conservado.

Empaqueta app.py en la raíz del ZIP. El controlador de la función existente es app.handler, por lo que importa el nombre del archivo del módulo:

zip -q function.zip app.py

Carga el ZIP binario con fileb://:

aws lambda update-function-code --function-name labex-q05-worker --zip-file fileb://function.zip --query CodeSha256 --output text

El hash devuelto identifica el código desplegado; el procesamiento real demostrará su comportamiento. Envía el primer pedido:

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'

Observa AWS View hasta que la cola esté vacía, el procesador informe processed: true y aparezca el pedido. Después léelo:

aws dynamodb get-item --table-name labex-q05-orders --key '{"id":{"S":"dedup-order"}}' --consistent-read --query Item

Espera la cantidad 4 y el total 1100. La decisión de escritura muestra la condición, Accepted, ningún elemento previo y este elemento almacenado después. Una condición configurada o un archivo cargado por sí solos no pueden demostrar que el trabajo se haya realizado correctamente.

Consumir trabajos repetidos sin repetir la escritura

En este paso, envía dos copias nuevas del mismo pedido de negocio y demuestra que otro pedido sigue completándose de forma independiente.

Envía la misma carga útil de negocio dos veces. Cada respuesta SendMessage tiene un ID de mensaje SQS nuevo, pero la clave del pedido sigue siendo dedup-order:

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'
aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'

Envía una clave de negocio independiente:

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"other-order","quantity":1}'

Observa AWS View hasta que hayan terminado cuatro ejecuciones reales y la cola esté vacía. Las dos ejecuciones repetidas informan processed: false, duplicate: true. DynamoDB rechaza ambas escrituras condicionales con ConditionalCheckFailedException; cada elemento antes y después permanece sin cambios. El primer pedido y other-order tienen una escritura aceptada cada uno. El consumidor confirma las copias repetidas en lugar de reintentar un trabajo ya completado.

Lee los resultados de negocio:

aws dynamodb scan --table-name labex-q05-orders --query Items

Espera exactamente dedup-order/4/1100 y other-order/1/350; el orden de los elementos puede variar. Inspecciona los resultados reales del procesador. Un evento de registro puede contener varias líneas; esta canalización envía la salida JSON nativa a jq -r, divide cada evento en líneas y selecciona únicamente las líneas RESULT, dejando los identificadores de recepción fuera de los resultados mostrados:

aws logs filter-log-events --log-group-name /aws/lambda/labex-q05-worker --filter-pattern '"RESULT"' --output json | jq -r '.events[].message | split("\n")[] | select(startswith("RESULT "))'

Comprueba la confirmación:

aws sqs get-queue-attributes --queue-url "$QUEUE_URL" --attribute-names ApproximateNumberOfMessages ApproximateNumberOfMessagesNotVisible

Ambos recuentos son cero. Hubo cuatro entregas reales de mensajes y ejecuciones de la función, dos escrituras de negocio aceptadas y dos rechazos condicionales nativos. Contar solo los elementos de la tabla no revelaría las sobrescrituras incondicionales; usa también las decisiones reales y los resultados correspondientes del consumidor.

La protección cubre un elemento de negocio. No convierte la escritura de DynamoDB y la confirmación SQS en una sola operación atómica. El registro del pedido es la protección duradera contra duplicados. Eliminarlo permite que la misma clave vuelva a crear un elemento, por lo que la retención y el diseño de claves de negocio forman parte de la política de un sistema real. Estas repeticiones sintéticas demuestran un tratamiento seguro del mismo trabajo de negocio; no establecen una garantía de concurrencia de producción ni de exactamente una vez de extremo a extremo.

Ejemplo de AWS View: cuatro ejecuciones reales producen dos pedidos y dos rechazos condicionales que conservan el elemento existente, con la cola vacía.

Eliminar la conexión del consumidor y los resultados de negocio

En este paso, elimina la conexión y los recursos propios después de confirmar el tratamiento de duplicados.

Detén tu asignación de origen de eventos antes de eliminar su cola de origen:

aws lambda delete-event-source-mapping --uuid "$MAPPING_ID" --query UUID --output text
aws sqs delete-queue --queue-url "$QUEUE_URL"

Elimina ambos elementos de negocio sintéticos y los registros de ejecución de este laboratorio:

aws dynamodb delete-item --table-name labex-q05-orders --key '{"id":{"S":"dedup-order"}}'
aws dynamodb delete-item --table-name labex-q05-orders --key '{"id":{"S":"other-order"}}'
aws logs delete-log-group --log-group-name /aws/lambda/labex-q05-worker

Confirma la ausencia y conserva los datos de referencia ajenos al trabajo:

aws lambda list-event-source-mappings --function-name labex-q05-worker --query EventSourceMappings
aws sqs list-queues
aws dynamodb scan --table-name labex-q05-orders --query Items
aws dynamodb scan --table-name labex-q05-reference --query Items

Las asignaciones y los pedidos están vacíos, no queda ninguna URL de cola y el elemento de referencia todavía indica keep unchanged. La función proporcionada y las estructuras de tablas permanecen en el entorno. Un error de red o de autenticación no puede demostrar la eliminación. AWS View conserva las decisiones históricas de escritura mientras muestra los recursos actuales vacíos.

Elimina tus archivos ordinarios de código y archivo comprimido:

rm -f app.py function.zip

Ejecuta la comprobación de limpieza antes de finalizar el entorno.

Resumen

Desplegaste una condición atómica de DynamoDB por clave de negocio en un consumidor real de cola. El primer pedido se completó; dos copias enviadas por separado se consumieron y confirmaron después de que el rechazo condicional nativo conservara el elemento original. Una clave distinta siguió produciendo su propio resultado. Las escrituras reales, los resultados del procesador y el estado vacío de la cola demostraron el resultado antes de la limpieza de recursos propios.

El reto aplica el aislamiento de fallos con límite y el proyecto sin servidor combina después la recuperación desde DLQ con esta protección por clave de negocio.