Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.
Menggunakan Lambda untuk memproses catatan dari Amazon Kinesis Data Streams
Anda dapat menggunakan fungsi Lambda untuk memproses catatan dalam aliran data Amazon Kinesis. Anda dapat memetakan fungsi Lambda ke konsumen throughput bersama Kinesis Data Streams (iterator standar), atau ke konsumen throughput khusus dengan fan-out yang ditingkatkan. https://docs.aws.amazon.com/kinesis/latest/dev/enhanced-consumers.html Untuk iterator standar, Lambda melakukan polling rekaman di setiap shard dalam aliran Kinesis Anda dengan menggunakan protokol HTTP. Pemetaan sumber kejadian berbagi throughput baca dengan konsumen lain dari shard tersebut.
Untuk perincian tentang aliran data Kinesis, lihat Data Pembacaan dari Amazon Kinesis Data Streams.
Untuk memulai, lihat Memproses rekaman Amazon Kinesis Data Streams dengan Lambda. Anda dapat mengonfigurasi parameter pemetaan sumber peristiwa, menggunakan jendela tumbling untuk agregasi, dan melihat contoh kode fungsi.
catatan
Kinesis mengenakan biaya untuk setiap shard dan, untuk keluaran yang ditingkatkan, pembacaan data dari aliran. Untuk perincian harga, lihat harga Amazon Kinesis
Aliran pemungutan suara dan batching
Lambda membaca rekaman dari aliran data dan memanggil fungsi Anda secara sinkron dengan kejadian yang berisi rekaman aliran. Lambda membaca rekaman dalam batch dan memanggil fungsi Anda untuk memproses rekaman dari batch. Setiap batch berisi catatan dari satu shard/data aliran.
Fungsi Lambda Anda adalah aplikasi konsumen untuk aliran data Anda. Itu memproses satu batch rekaman pada satu waktu dari setiap shard. Anda dapat memetakan fungsi Lambda ke konsumen dengan throughput bersama (iterator standar), atau ke konsumen dengan throughput khusus dengan keluaran yang ditingkatkan.
-
Iterator standar: Lambda mensurvei setiap pecahan dalam aliran Kinesis Anda untuk catatan pada tingkat dasar sekali per detik. Saat tersedia lebih banyak rekaman, Lambda melanjutkan pemrosesan batch sampai fungsi dapat menyusul aliran. Pemetaan sumber kejadian berbagi throughput baca dengan konsumen lain dari shard tersebut.
-
Peningkatan fan-out: Untuk meminimalkan latensi dan memaksimalkan throughput baca, buat konsumen aliran data dengan fan-out yang ditingkatkan. Konsumen dengan keluaran yang ditingkatkan mendapatkan koneksi khusus ke setiap shard yang tidak memengaruhi pembacaan aplikasi lain dari aliran tersebut. Konsumen streaming digunakan HTTP/2 untuk mengurangi latensi dengan mendorong catatan ke Lambda melalui koneksi berumur panjang dan dengan mengompresi header permintaan. Anda dapat membuat konsumen aliran dengan API Kinesis RegisterStreamConsumer.
aws kinesis register-stream-consumer \ --consumer-name con1 \ --stream-arn arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream
Anda akan melihat output berikut:
{ "Consumer": { "ConsumerName": "con1", "ConsumerARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream/consumer/con1:1540591608", "ConsumerStatus": "CREATING", "ConsumerCreationTimestamp": 1540591608.0 } }
Untuk meningkatkan kecepatan pemrosesan catatan fungsi Anda, tambahkan pecahan ke aliran
Apabila fungsi Anda tidak dapat menskalakan naik untuk menangani jumlah total batch yang bersamaan, mintalah kenaikan kuota atau cadangkan konkurensi untuk fungsi Anda.
Secara default, Lambda memanggil fungsi Anda segera setelah catatan tersedia. Jika batch yang dibaca Lambda dari sumber peristiwa hanya memiliki satu catatan di dalamnya, Lambda hanya mengirimkan satu catatan ke fungsi. Untuk menghindari pemanggilan fungsi dengan sejumlah kecil catatan, Anda dapat memberi tahu sumber peristiwa untuk menyangga catatan hingga 5 menit dengan mengonfigurasi jendela batching. Sebelum menjalankan fungsi, Lambda terus membaca catatan dari sumber peristiwa hingga mengumpulkan batch penuh, jendela batching berakhir, atau batch mencapai batas muatan 6 MB. Untuk informasi selengkapnya, lihat Perilaku batching.
Awas
Pemetaan sumber peristiwa lambda memproses setiap peristiwa setidaknya sekali, dan pemrosesan rekaman duplikat dapat terjadi. Untuk menghindari potensi masalah terkait peristiwa duplikat, kami sangat menyarankan agar Anda membuat kode fungsi Anda idempotent. Untuk mempelajari lebih lanjut, lihat Bagaimana cara membuat fungsi Lambda saya idempoten
Lambda tidak menunggu ekstensi yang dikonfigurasi selesai sebelum mengirim batch berikutnya untuk diproses. Dengan kata lain, ekstensi Anda dapat terus berjalan saat Lambda memproses kumpulan catatan berikutnya. Hal ini dapat menyebabkan masalah pelambatan jika Anda melanggar pengaturan atau batasan konkur ensi akun Anda. Untuk mendeteksi apakah ini merupakan masalah potensial, pantau fungsi Anda dan periksa apakah Anda melihat metrik konkurensi yang lebih tinggi dari yang diharapkan untuk pemetaan sumber acara Anda. Karena waktu yang singkat di antara pemanggilan, Lambda dapat secara singkat melaporkan penggunaan konkurensi yang lebih tinggi daripada jumlah pecahan. Ini bisa benar bahkan untuk fungsi Lambda tanpa ekstensi.
Konfigurasikan ParallelizationFactor pengaturan untuk memproses satu pecahan aliran data Kinesis dengan lebih dari satu pemanggilan Lambda secara bersamaan. Anda dapat menentukan jumlah batch bersamaan yang dipolling Lambda dari pecahan dengan menggunakan faktor paralelisasi dari 1 (default) hingga 10. Misalnya, ketika Anda menyetel ParallelizationFactor ke 2, Anda dapat memiliki maksimum 200 pemanggilan Lambda bersamaan untuk memproses 100 pecahan data Kinesis (meskipun dalam praktiknya, Anda mungkin melihat nilai yang berbeda untuk metrik). ConcurrentExecutions Hal ini membantu meningkatkan skala throughput pemrosesan ketika volume data tidak stabil dan IteratorAge tinggi. Ketika Anda meningkatkan jumlah batch bersamaan per shard, Lambda masih memastikan pemrosesan dalam urutan di tingkat kunci partisi.
Anda juga dapat menggunakan ParallelizationFactor dengan agregasi Kinesis. Perilaku pemetaan sumber peristiwa tergantung pada apakah Anda menggunakan fan-out yang disempurnakan:
-
Tanpa fan-out yang disem purnakan: Semua peristiwa di dalam acara gabungan harus memiliki kunci partisi yang sama. Kunci partisi juga harus cocok dengan peristiwa agregat. Jika peristiwa di dalam peristiwa agregat memiliki kunci partisi yang berbeda, Lambda tidak dapat menjamin pemrosesan acara secara berurutan dengan kunci partisi.
-
Dengan fan-out yang ditingkatkan: Pertama, Lambda memecahkan kode peristiwa agregat menjadi peristiwa individualnya. Acara agregat dapat memiliki kunci partisi yang berbeda dari peristiwa yang dikandungnya. Namun, peristiwa yang tidak sesuai dengan kunci partisi dijatu hkan dan hilang
. Lambda tidak memproses peristiwa ini, dan tidak mengirimnya ke tujuan kegagalan yang dikonfigurasi.
Contoh peristiwa
contoh
{ "Records": [ { "kinesis": { "kinesisSchemaVersion": "1.0", "partitionKey": "1", "sequenceNumber": "49590338271490256608559692538361571095921575989136588898", "data": "SGVsbG8sIHRoaXMgaXMgYSB0ZXN0Lg==", "approximateArrivalTimestamp": 1545084650.987 }, "eventSource": "aws:kinesis", "eventVersion": "1.0", "eventID": "shardId-000000000006:49590338271490256608559692538361571095921575989136588898", "eventName": "aws:kinesis:record", "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role", "awsRegion": "us-east-2", "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream" }, { "kinesis": { "kinesisSchemaVersion": "1.0", "partitionKey": "1", "sequenceNumber": "49590338271490256608559692540925702759324208523137515618", "data": "VGhpcyBpcyBvbmx5IGEgdGVzdC4=", "approximateArrivalTimestamp": 1545084711.166 }, "eventSource": "aws:kinesis", "eventVersion": "1.0", "eventID": "shardId-000000000006:49590338271490256608559692540925702759324208523137515618", "eventName": "aws:kinesis:record", "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role", "awsRegion": "us-east-2", "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream" } ] }