Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Mengkonfigurasi kontrol penanganan kesalahan untuk sumber peristiwa Kafka
Anda dapat mengonfigurasi cara Lambda menangani kesalahan dan percobaan ulang untuk pemetaan sumber peristiwa Kafka Anda. Konfigurasi ini membantu Anda mengontrol bagaimana Lambda memproses catatan yang gagal dan mengelola perilaku coba ulang.
Konfigurasi coba ulang yang tersedia
Konfigurasi coba ulang berikut tersedia untuk Amazon MSK dan sumber peristiwa Kafka yang dikelola sendiri:
-
Upaya coba ulang maksimum — Jumlah maksimum kali Lambda mencoba lagi saat fungsi Anda mengembalikan kesalahan. Ini tidak menghitung upaya pemanggilan awal. Defaultnya adalah -1 (tak terbatas). Saat Anda mengonfigurasi percobaan ulang tak terbatas dan tujuan pada kegagalan, Lambda secara otomatis menerapkan maksimal 10 upaya percobaan ulang.
-
Usia rekaman maksimum — Usia maksimum rekaman yang dikirim Lambda ke fungsi Anda. Defaultnya adalah -1 (tak terbatas).
-
Pisahkan batch pada kesalahan - Ketika fungsi Anda mengembalikan kesalahan, bagi batch menjadi dua batch yang lebih kecil dan coba lagi masing-masing secara terpisah. Ini membantu mengisolasi catatan bermasalah.
-
Respon batch par tial — Izinkan fungsi Anda mengembalikan informasi tentang catatan mana dalam batch yang gagal diproses, sehingga Lambda dapat mencoba lagi hanya catatan yang gagal.
Mengkonfigurasi kontrol penanganan kesalahan (konsol)
Anda dapat mengonfigurasi perilaku coba ulang saat membuat atau memperbarui pemetaan sumber peristiwa Kafka di konsol Lambda.
Untuk mengonfigurasi perilaku coba ulang untuk sumber peristiwa Kafka (konsol)
-
Buka halaman Fungsi
di konsol Lambda. -
Pilih nama fungsi Anda.
-
Lakukan salah satu tindakan berikut:
-
Untuk menambahkan pemicu Kafka baru, di bawah Ikhtis ar fungsi, pilih Tambahkan pemicu.
-
Untuk memodifikasi pemicu Kafka yang ada, pilih pemicu lalu pilih Edit.
-
-
Di bawah konfigurasi poller peristiwa, pilih mode yang disediakan untuk mengonfigurasi kontrol penanganan kesalahan:
-
Untuk Upaya Coba Ulang, masukkan jumlah maksimum percobaan ulang (0-10000, atau -1 untuk tak terbatas).
-
Untuk Usia rekor maksimum, masukkan usia maksimum dalam detik (60-604800, atau -1 untuk tak terbatas).
-
Untuk mengaktifkan pemisahan batch saat terjadi kesalahan, pilih Pisahkan batch pada kesalahan.
-
Untuk mengaktifkan respons batch parsional, pilih ReportBatchItemFailures.
-
-
Pilih Tam bah atau Simpan.
Mengkonfigurasi perilaku coba lagi (AWS CLI)
Gunakan AWS CLI perintah berikut untuk mengonfigurasi perilaku coba ulang untuk pemetaan sumber peristiwa Kafka Anda.
Membuat pemetaan sumber acara dengan konfigurasi coba ulang
Contoh berikut membuat pemetaan sumber peristiwa Kafka yang dikelola sendiri dengan kontrol penanganan kesalahan:
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"
Untuk sumber acara 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"
Memperbarui konfigurasi coba lagi
Gunakan update-event-source-mapping perintah untuk memodifikasi konfigurasi coba ulang untuk pemetaan sumber peristiwa yang ada:
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
Respons batch parsional, juga dikenal sebagai ReportBatchItemFailures, adalah fitur utama untuk penanganan kesalahan dalam integrasi Lambda dengan sumber Kafka. Tanpa fitur ini, ketika kesalahan terjadi di salah satu item dalam batch, itu menghasilkan pemrosesan ulang semua pesan dalam batch itu. Dengan respons batch parsional diaktifkan dan diimplementasikan, handler mengembalikan pengidentifikasi hanya untuk pesan yang gagal, memungkinkan Lambda untuk mencoba lagi hanya item tertentu. Ini memberikan kontrol yang lebih besar atas bagaimana batch yang berisi pesan gagal diproses.
Untuk melaporkan kesalahan batch, gunakan skema JSON ini:
{ "batchItemFailures": [ { "itemIdentifier": { "partition": "topic-partition_number", "offset": 100 } }, ... ] }
penting
Jika Anda mengembalikan JSON atau null yang valid kosong, pemetaan sumber peristiwa menganggap batch berhasil diproses. Setiap topic-partition_number atau offset yang tidak valid yang dikembalikan yang tidak ada dalam peristiwa yang dipanggil diperlakukan sebagai kegagalan dan seluruh batch dicoba ulang.
Contoh kode berikut menunjukkan bagaimana menerapkan respons batch parsional untuk fungsi Lambda yang menerima peristiwa dari sumber Kafka. Fungsi melaporkan kegagalan item batch dalam respons, memberi sinyal ke Lambda untuk mencoba lagi pesan tersebut nanti.
Berikut adalah implementasi handler Python Lambda yang menunjukkan pendekatan ini:
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
Berikut adalah ver Node.js sinya:
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 };
Berikut adalah versi 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 } }