Le traduzioni sono generate tramite traduzione automatica. In caso di conflitto tra il contenuto di una traduzione e la versione originale in Inglese, quest'ultima prevarrà.
Configurazione dei controlli di gestione degli errori per le sorgenti di eventi Kafka
Puoi configurare il modo in cui Lambda gestisce gli errori e i tentativi per le mappature delle sorgenti degli eventi Kafka. Queste configurazioni ti aiutano a controllare il modo in cui Lambda elabora i record non riusciti e gestisce il comportamento dei nuovi tentativi.
Configurazioni di nuovo tentativo disponibili
Le seguenti configurazioni di ripetizione sono disponibili sia per Amazon MSK che per le sorgenti di eventi Kafka autogestite:
-
Numero massimo di tentativi di ripetizione: il numero massimo di tentativi di Lambda quando la funzione restituisce un errore. Questo non conta il tentativo di richiamo iniziale. L'impostazione predefinita è -1 (infinito). Quando si configurano sia tentativi infiniti che una destinazione in caso di errore, Lambda applica automaticamente un massimo di 10 tentativi di nuovo tentativo.
-
Età massima del record: l'età massima di un record che Lambda invia alla tua funzione. L'impostazione predefinita è -1 (infinito).
-
Dividi il batch in caso di errore: quando la funzione restituisce un errore, suddividi il batch in due batch più piccoli e riprova ciascuno separatamente. Questo aiuta a isolare i record problematici.
-
Risposta parziale in batch: consenti alla funzione di restituire informazioni su quali record di un batch non sono stati elaborati, in modo che Lambda possa riprovare solo i record non riusciti.
Configurazione dei controlli di gestione degli errori (console)
È possibile configurare il comportamento dei nuovi tentativi durante la creazione o l'aggiornamento di una mappatura della sorgente di eventi Kafka nella console Lambda.
Per configurare il comportamento dei tentativi di ripetizione per una sorgente di eventi Kafka (console)
-
Aprire la pagina Funzioni
della console Lambda. -
Scegli il nome della tua funzione.
-
Esegui una delle seguenti operazioni:
-
Per aggiungere un nuovo trigger Kafka, in Panoramica delle funzioni, scegli Aggiungi trigger.
-
Per modificare un trigger Kafka esistente, scegli il trigger e quindi scegli Modifica.
-
-
In Event poller configuration, seleziona la modalità provisioned per configurare i controlli di gestione degli errori:
-
Per i tentativi di riprova, inserisci il numero massimo di tentativi (0-10000 o -1 per infinito).
-
Per Età massima del record, inserisci l'età massima in secondi (60-604800 o -1 per infinito).
-
Per abilitare la suddivisione in batch in caso di errori, seleziona Dividi batch in caso di errore.
-
Per abilitare la risposta parziale in batch, selezionare ReportBatchItemFailures.
-
-
Scegli Aggiungi o Salva.
Configurazione del comportamento dei tentativi (AWS CLI)
Utilizzate i seguenti AWS CLI comandi per configurare il comportamento dei tentativi per le mappature delle sorgenti degli eventi Kafka.
Creazione di una mappatura della sorgente di eventi con configurazioni di ripetizione
L'esempio seguente crea una mappatura autogestita della sorgente di eventi Kafka con controlli di gestione degli errori:
aws lambda create-event-source-mapping \ --function-name my-kafka-function \ --topics my-kafka-topic \ --source-access-configuration Type=SASL_SCRAM_512_AUTH,URI=arn:aws:secretsmanager:us-east-1:111122223333:secret:MyBrokerSecretName \ --self-managed-event-source '{"Endpoints":{"KAFKA_BOOTSTRAP_SERVERS":["abc.xyz.com:9092"]}}' \ --starting-position LATEST \ --provisioned-poller-config MinimumPollers=1,MaximumPollers=1 \ --maximum-retry-attempts 3 \ --maximum-record-age-in-seconds 3600 \ --bisect-batch-on-function-error \ --function-response-types "ReportBatchItemFailures"
Per le sorgenti di eventi Amazon MSK:
aws lambda create-event-source-mapping \ --event-source-arn arn:aws:kafka:us-east-1:111122223333:cluster/my-cluster/fc2f5bdf-fd1b-45ad-85dd-15b4a5a6247e-2 \ --topics AWSMSKKafkaTopic \ --starting-position LATEST \ --function-name my-kafka-function \ --source-access-configurations '[{"Type": "SASL_SCRAM_512_AUTH","URI": "arn:aws:secretsmanager:us-east-1:111122223333:secret:my-secret"}]' \ --provisioned-poller-config MinimumPollers=1,MaximumPollers=1 \ --maximum-retry-attempts 3 \ --maximum-record-age-in-seconds 3600 \ --bisect-batch-on-function-error \ --function-response-types "ReportBatchItemFailures"
Aggiornamento delle configurazioni dei nuovi tentativi
Utilizzate il update-event-source-mapping comando per modificare le configurazioni dei nuovi tentativi per una mappatura della sorgente di eventi esistente:
aws lambda update-event-source-mapping \ --uuid 12345678-1234-1234-1234-123456789012 \ --maximum-retry-attempts 5 \ --maximum-record-age-in-seconds 7200 \ --bisect-batch-on-function-error \ --function-response-types "ReportBatchItemFailures"
PartialBatchResponse
La risposta parziale in batch, nota anche come ReportBatchItemFailures, è una funzionalità chiave per la gestione degli errori nell'integrazione di Lambda con i sorgenti Kafka. Senza questa funzionalità, quando si verifica un errore in uno degli elementi di un batch, comporta la rielaborazione di tutti i messaggi in quel batch. Con la risposta parziale in batch abilitata e implementata, il gestore restituisce gli identificatori solo per i messaggi non riusciti, consentendo a Lambda di riprovare solo quegli elementi specifici. Ciò fornisce un maggiore controllo sul modo in cui vengono elaborati i batch contenenti messaggi non riusciti.
Per segnalare errori batch, utilizzate questo schema JSON:
{ "batchItemFailures": [ { "itemIdentifier": { "partition": "topic-partition_number", "offset": 100 } }, ... ] }
Importante
Se restituisci un JSON o null valido vuoto, la mappatura dell'origine dell'evento considera un batch come elaborato correttamente. Qualsiasi topic-partition_number o offset restituito non valido che non era presente nell'evento richiamato viene considerato un errore e l'intero batch viene riprovato.
I seguenti esempi di codice mostrano come implementare una risposta batch parziale per le funzioni Lambda che ricevono eventi da sorgenti Kafka. La funzione riporta gli errori degli elementi batch nella risposta, segnalando a Lambda di riprovare tali messaggi in un secondo momento.
Ecco un'implementazione del gestore Python Lambda che mostra questo approccio:
import base64 from typing import Any, Dict, List def lambda_handler(event: Dict[str, Any], context: Any) -> Dict[str, List[Dict[str, Dict[str, Any]]]]: failures: List[Dict[str, Dict[str, Any]]] = [] records_dict = event.get("records", {}) for topic_partition, records_list in records_dict.items(): for record in records_list: topic = record.get("topic") partition = record.get("partition") offset = record.get("offset") value_b64 = record.get("value") try: data = base64.b64decode(value_b64).decode("utf-8") process_message(data) except Exception as exc: print(f"Failed to process record topic={topic} partition={partition} offset={offset}: {exc}") item_identifier: Dict[str, Any] = { "partition": f"{topic}-{partition}", "offset": int(offset) if offset is not None else None, } failures.append({"itemIdentifier": item_identifier}) return {"batchItemFailures": failures} def process_message(data: str) -> None: # Your business logic for a single message pass
Ecco una versione: Node.js
const { Buffer } = require("buffer"); const handler = async (event) => { const failures = []; for (let topicPartition in event.records) { const records = event.records[topicPartition]; for (const record of records) { const topic = record.topic; const partition = record.partition; const offset = record.offset; const valueBase64 = record.value; const data = Buffer.from(valueBase64, "base64").toString("utf8"); try { await processMessage(data); } catch (error) { console.error("Failed to process record", { topic, partition, offset, error }); const itemIdentifier = { "partition": `${topic}-${partition}`, "offset": Number(offset), }; failures.push({ itemIdentifier }); } } } return { batchItemFailures: failures }; }; async function processMessage(payload) { // Your business logic for a single message } module.exports = { handler };
Ecco una versione Java:
import com.amazonaws.services.lambda.runtime.Context; import com.amazonaws.services.lambda.runtime.RequestHandler; import java.util.ArrayList; import java.util.Base64; import java.util.HashMap; import java.util.List; import java.util.Map; public class KafkaBatchHandler implements RequestHandler<Map<String, Object>, Map<String, Object>> { @SuppressWarnings("unchecked") @Override public Map<String, Object> handleRequest(Map<String, Object> event, Context context) { List<Map<String, Object>> failures = new ArrayList<>(); Map<String, List<Map<String, Object>>> records = (Map<String, List<Map<String, Object>>>) event.getOrDefault("records", Map.of()); for (Map.Entry<String, List<Map<String, Object>>> entry : records.entrySet()) { for (Map<String, Object> record : entry.getValue()) { String topic = (String) record.get("topic"); Object partition = record.get("partition"); Object offset = record.get("offset"); String valueBase64 = (String) record.get("value"); try { String data = new String(Base64.getDecoder().decode(valueBase64), "UTF-8"); processMessage(data); } catch (Exception e) { System.err.printf("Failed to process record topic=%s partition=%s offset=%s: %s%n", topic, partition, offset, e.getMessage()); Map<String, Object> itemIdentifier = new HashMap<>(); itemIdentifier.put("partition", topic + "-" + partition); itemIdentifier.put("offset", offset instanceof Number ? ((Number) offset).longValue() : null); Map<String, Object> failure = new HashMap<>(); failure.put("itemIdentifier", itemIdentifier); failures.add(failure); } } } Map<String, Object> response = new HashMap<>(); response.put("batchItemFailures", failures); return response; } private void processMessage(String data) { // Your business logic for a single message } }