View a markdown version of this page

Tutorial: Menggunakan pemetaan sumber peristiwa Amazon MSK untuk memanggil fungsi Lambda - AWS Lambda

Terjemahan disediakan oleh mesin penerjemah. Jika konten terjemahan yang diberikan bertentangan dengan versi bahasa Inggris aslinya, utamakan versi bahasa Inggris.

Tutorial: Menggunakan pemetaan sumber peristiwa Amazon MSK untuk memanggil fungsi Lambda

Dalam tutorial ini, Anda melakukan hal berikut:

  • Buat fungsi Lambda di AWS akun yang sama dengan cluster Amazon MSK yang ada.

  • Konfigurasikan jaringan dan otentikasi untuk Lambda untuk berkomunikasi dengan Amazon MSK.

  • Siapkan pemetaan sumber peristiwa Lambda Amazon MSK, yang menjalankan fungsi Lambda Anda saat peristiwa muncul di topik.

Setelah Anda selesai dengan langkah-langkah ini, ketika acara dikirim ke Amazon MSK, Anda dapat mengatur fungsi Lambda untuk memproses peristiwa tersebut secara otomatis dengan kode Lambda kustom Anda sendiri.

Apa yang dapat Anda lakukan dengan fitur ini?

Contoh solusi: Gunakan pemetaan sumber acara MSK untuk mengirimkan skor langsung kepada pelanggan Anda.

Pertimbangkan skenario berikut: Perusahaan Anda menghosting aplikasi web tempat pelanggan Anda dapat melihat informasi tentang acara langsung, seperti permainan olahraga. Pembaruan informasi dari game diberikan kepada tim Anda melalui topik Kafka di Amazon MSK. Anda ingin merancang solusi yang menggunakan pembaruan dari topik MSK untuk memberikan tampilan terbaru dari acara langsung kepada pelanggan di dalam aplikasi yang Anda kembangkan. Anda telah memutuskan pendekatan desain berikut: Aplikasi klien Anda berkomunikasi dengan backend tanpa server yang dihosting di. AWS Klien terhubung melalui sesi websocket menggunakan API Gateway WebSocket API Amazon.

Dalam solusi ini, Anda memerlukan komponen yang membaca peristiwa MSK, melakukan beberapa logika khusus untuk mempersiapkan peristiwa tersebut untuk lapisan aplikasi dan kemudian meneruskan informasi tersebut ke API Gateway API. Anda dapat mengimplementasikan komponen ini dengan AWS Lambda, dengan menyediakan logika kustom Anda dalam fungsi Lambda, lalu memanggilnya dengan pemetaan sumber peristiwa AWS Lambda Amazon MSK.

Untuk informasi selengkapnya tentang mengimplementasikan solusi menggunakan API Gateway WebSocket API Amazon, lihat tutorial WebSocket API di dokumentasi API Gateway.

Prasyarat

AWS Akun dengan sumber daya yang telah dikonfigurasi sebelumnya berikut:

Untuk memenuhi prasyarat ini, sebaiknya ikuti Mem ulai menggunakan Amazon MSK di dokumentasi Amazon MSK.

  • Cluster Amazon MSK. Lihat Membuat cluster Amazon MSK di Mem ulai menggunakan Amazon M SK.

  • Konfigurasi berikut:

    • Pastikan otentikasi berbasis peran IAM Di aktifkan di pengaturan keamanan cluster Anda. Ini meningkatkan keamanan Anda dengan membatasi fungsi Lambda Anda untuk hanya mengakses sumber daya Amazon MSK yang diperlukan. Ini diaktifkan secara default pada cluster Amazon MSK baru.

    • Pastikan akses publik dinonaktifkan di pengaturan jaringan cluster Anda. Membatasi akses cluster Amazon MSK Anda ke internet meningkatkan keamanan Anda dengan membatasi berapa banyak perantara yang menangani data Anda. Ini diaktifkan secara default pada cluster Amazon MSK baru.

  • Topik Kafka di cluster Amazon MSK Anda untuk digunakan untuk solusi ini. Lihat Membuat topik di Mem ulai menggunakan Amazon MSK.

  • Host admin Kafka yang disiapkan untuk mengambil informasi dari klaster Kafka Anda dan mengirim peristiwa Kafka ke topik Anda untuk pengujian, seperti instans Amazon EC2 dengan CLI admin Kafka dan pustaka IAM Amazon MSK diinstal. Lihat Membuat mesin klien di Mem ulai menggunakan Amazon MSK.

Setelah Anda menyiapkan sumber daya ini, kumpulkan informasi berikut dari AWS akun Anda untuk mengonfirmasi bahwa Anda siap untuk melanjutkan.

  • Nama cluster Amazon MSK Anda. Anda dapat menemukan informasi ini di konsol Amazon MSK.

  • UUID cluster, bagian dari ARN untuk cluster Amazon MSK Anda, yang dapat Anda temukan di konsol Amazon MSK. Ikuti prosedur di Daftar cluster di dokumentasi Amazon MSK untuk menemukan informasi ini.

  • Grup keamanan yang terkait dengan cluster Amazon MSK Anda. Anda dapat menemukan informasi ini di konsol Amazon MSK. Dalam langkah-langkah berikut, rujuk ini sebagai milik AndaclusterSecurityGroups.

  • Id Amazon VPC yang berisi cluster Amazon MSK Anda. Anda dapat menemukan informasi ini dengan mengidentifikasi subnet yang terkait dengan cluster Amazon MSK Anda di konsol Amazon MSK, lalu mengidentifikasi Amazon VPC yang terkait dengan subnet di Amazon VPC Console.

  • Nama topik Kafka yang digunakan dalam solusi Anda. Anda dapat menemukan informasi ini dengan memanggil cluster Amazon MSK Anda dengan Kafka topics CLI dari host admin Kafka Anda. Untuk informasi selengkapnya tentang topik CLI, lihat Men ambahkan dan menghapus topik dalam dokumentasi Kafka.

  • Nama grup konsumen untuk topik Kafka Anda, cocok untuk digunakan oleh fungsi Lambda Anda. Grup ini dapat dibuat secara otomatis oleh Lambda, jadi Anda tidak perlu membuatnya dengan Kafka CLI. Jika Anda perlu mengelola grup konsumen Anda, untuk mempelajari lebih lanjut tentang CLI grup konsumen, lihat Mengelola Grup Konsumen di dokumentasi Kafka.

Izin berikut di AWS akun Anda:

  • Izin untuk membuat dan mengelola fungsi Lambda.

  • Izin untuk membuat kebijakan IAM dan mengaitkannya dengan fungsi Lambda Anda.

  • Izin untuk membuat titik akhir Amazon VPC dan mengubah konfigurasi jaringan di Amazon VPC yang menghosting cluster Amazon MSK Anda.

Jika Anda belum menginstal AWS Command Line Interface, ikuti langkah-langkah di Meng instal atau memperbarui versi terbaru AWS CLI untuk menginstalnya.

Tutorial ini membutuhkan terminal baris perintah atau shell untuk menjalankan perintah. Di Linux dan macOS, gunakan shell dan manajer paket pilihan Anda.

catatan

Di Windows, beberapa perintah Bash CLI yang biasa Anda gunakan dengan Lambda (sepertizip) tidak didukung oleh terminal bawaan sistem operasi. Untuk mendapatkan Windows-integrated versi Ubuntu dan Bash, instal Subsistem Windows untuk Linux.

Konfigurasikan konektivitas jaringan agar Lambda berkomunikasi dengan Amazon MSK

Gunakan AWS PrivateLink untuk menghubungkan Lambda dan Amazon MSK. Anda dapat melakukannya dengan membuat antarmuka titik akhir Amazon VPC di konsol Amazon VPC. Untuk informasi selengkapnya tentang konfigurasi jaringan, lihatMengkonfigurasi cluster Amazon MSK dan jaringan Amazon VPC Anda untuk Lambda.

Ketika pemetaan sumber peristiwa Amazon MSK berjalan atas nama fungsi Lambda, pemetaan tersebut mengasumsikan peran eksekusi fungsi Lambda. Peran IAM ini mengotorisasi pemetaan untuk mengakses sumber daya yang diamankan oleh IAM, seperti cluster Amazon MSK Anda. Meskipun komponen berbagi peran eksekusi, pemetaan Amazon MSK dan fungsi Lambda Anda memiliki persyaratan konektivitas terpisah untuk tugas masing-masing, seperti yang ditunjukkan pada diagram berikut.

Fungsi Lambda mensurvei cluster dan berkomunikasi dengan Lambda menggunakan. AWS STS

Pemetaan sumber acara milik grup keamanan cluster Amazon MSK Anda. Pada langkah jaringan ini, buat titik akhir Amazon VPC dari VPC cluster Amazon MSK Anda untuk menghubungkan pemetaan sumber peristiwa ke layanan Lambda dan STS. Amankan titik akhir ini untuk menerima lalu lintas dari grup keamanan cluster Amazon MSK Anda. Kemudian, sesuaikan grup keamanan cluster Amazon MSK untuk memungkinkan pemetaan sumber peristiwa berkomunikasi dengan cluster Amazon MSK.

Anda dapat mengonfigurasi langkah-langkah berikut menggunakan Konsol Manajemen AWS.

Untuk mengkonfigurasi antarmuka titik akhir Amazon VPC untuk menghubungkan Lambda dan Amazon MSK
  1. Buat grup keamanan untuk antarmuka Anda titik akhir Amazon VPC,endpointSecurityGroup, yang memungkinkan lalu lintas TCP masuk pada 443 dari. clusterSecurityGroups Ikuti prosedur di Buat grup keamanan di dokumentasi Amazon EC2 untuk membuat grup keamanan. Kemudian, ikuti prosedur di Tambahkan aturan ke grup keamanan di dokumentasi Amazon EC2 untuk menambahkan aturan yang sesuai.

    Buat grup keamanan dengan informasi berikut:

    Saat menambahkan aturan masuk, buat aturan untuk setiap grup keamanan diclusterSecurityGroups. Untuk setiap aturan:

    • Untuk Jenis, pilih HTTPS.

    • Untuk Sumber, pilih salah satunyaclusterSecurityGroups.

  2. Buat titik akhir yang menghubungkan layanan Lambda ke Amazon VPC yang berisi cluster Amazon MSK Anda. Ikuti prosedur di Buat titik akhir antarmuka.

    Buat titik akhir antarmuka dengan informasi berikut:

    • Untuk Nama Layanan, pilihcom.amazonaws.regionName.lambda, di mana menghosting regionName fungsi Lambda Anda.

    • Untuk VPC, pilih Amazon VPC yang berisi cluster Amazon MSK Anda.

    • Untuk grup keamanan, pilihendpointSecurityGroup, yang Anda buat sebelumnya.

    • Untuk Subnet, pilih subnet yang menghosting cluster Amazon MSK Anda.

    • Untuk Kebijakan, berikan dokumen kebijakan berikut, yang mengamankan titik akhir untuk digunakan oleh prinsipal layanan Lambda untuk tindakan tersebutlambda:InvokeFunction.

      { "Statement": [ { "Action": "lambda:InvokeFunction", "Effect": "Allow", "Principal": { "Service": [ "lambda.amazonaws.com" ] }, "Resource": "*" } ] }
    • Pastikan En able DNS name tetap disetel.

  3. Buat titik akhir yang menghubungkan AWS STS layanan ke Amazon VPC yang berisi cluster Amazon MSK Anda. Ikuti prosedur di Buat titik akhir antarmuka.

    Buat titik akhir antarmuka dengan informasi berikut:

    • Untuk Nama Layanan, pilih AWS STS.

    • Untuk VPC, pilih Amazon VPC yang berisi cluster Amazon MSK Anda.

    • Untuk Grup keamanan, pilihendpointSecurityGroup.

    • Untuk Subnet, pilih subnet yang menghosting cluster Amazon MSK Anda.

    • Untuk Kebijakan, berikan dokumen kebijakan berikut, yang mengamankan titik akhir untuk digunakan oleh prinsipal layanan Lambda untuk tindakan tersebutsts:AssumeRole.

      { "Statement": [ { "Action": "sts:AssumeRole", "Effect": "Allow", "Principal": { "Service": [ "lambda.amazonaws.com" ] }, "Resource": "*" } ] }
    • Pastikan En able DNS name tetap disetel.

  4. Untuk setiap grup keamanan yang terkait dengan cluster Amazon MSK Anda, yaitu, diclusterSecurityGroups, izinkan hal berikut:

    • Izinkan semua lalu lintas TCP masuk dan keluar pada 9098 ke semuaclusterSecurityGroups, termasuk di dalam dirinya sendiri.

    • Izinkan semua lalu lintas TCP keluar pada 443.

    Beberapa lalu lintas ini diizinkan oleh aturan grup keamanan default, jadi jika cluster Anda dilampirkan ke grup keamanan tunggal, dan grup tersebut memiliki aturan default, aturan tambahan tidak diperlukan. Untuk menyesuaikan aturan grup keamanan, ikuti prosedur di Men ambahkan aturan ke grup keamanan di dokumentasi Amazon EC2.

    Tambahkan aturan ke grup keamanan Anda dengan informasi berikut:

    • Untuk setiap aturan masuk atau aturan keluar untuk port 9098, berikan

      • Untuk Jenis, pilih TCP Kustom.

      • Untuk rentang Port, berikan 9098.

      • Untuk Sumber, berikan salah satunyaclusterSecurityGroups.

    • Untuk setiap aturan masuk untuk port 443, untuk Jenis, pilih HTTPS.

Buat peran IAM untuk Lambda untuk dibaca dari topik Amazon MSK Anda

Identifikasi persyaratan autentikasi untuk Lambda untuk dibaca dari topik Amazon MSK Anda, lalu tentukan persyaratan tersebut dalam kebijakan. Buat peran,lambdaAuthRole, yang mengizinkan Lambda untuk menggunakan izin tersebut. Otorisasi tindakan pada cluster Amazon MSK menggunakan tindakan kafka-cluster IAM. Kemudian, otorisasi Lambda untuk melakukan tindakan Amazon MSK kafka dan Amazon EC2 yang diperlukan untuk menemukan dan terhubung ke cluster Amazon MSK Anda, serta CloudWatch tindakan sehingga Lambda dapat mencatat apa yang telah dilakukannya.

Untuk menjelaskan persyaratan autentikasi agar Lambda dibaca dari Amazon MSK
  1. Tulis dokumen kebijakan IAM (dokumen JSON),clusterAuthPolicy, yang memungkinkan Lambda membaca dari topik Kafka Anda di cluster Amazon MSK Anda menggunakan grup konsumen Kafka Anda. Lambda membutuhkan grup konsumen Kafka untuk diatur saat membaca.

    Ubah template berikut agar sesuai dengan prasyarat Anda:

    JSON
    { "Version":"2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "kafka-cluster:Connect", "kafka-cluster:DescribeGroup", "kafka-cluster:AlterGroup", "kafka-cluster:DescribeTopic", "kafka-cluster:ReadData", "kafka-cluster:DescribeClusterDynamicConfiguration" ], "Resource": [ "arn:aws:kafka:us-east-1:111122223333:cluster/mskClusterName/cluster-uuid", "arn:aws:kafka:us-east-1:111122223333:topic/mskClusterName/cluster-uuid/mskTopicName", "arn:aws:kafka:us-east-1:111122223333:group/mskClusterName/cluster-uuid/mskGroupName" ] } ] }

    Untuk informasi lebih lanjut, konsultasikanMengkonfigurasi izin Lambda untuk pemetaan sumber acara Amazon MSK. Saat menulis kebijakan Anda:

    • Ganti us-east-1 dan 111122223333 dengan Wilayah AWS dan Akun AWS dari cluster Amazon MSK Anda.

    • UntukmskClusterName, berikan nama cluster Amazon MSK Anda.

    • Untukcluster-uuid, berikan UUID di ARN untuk cluster Amazon MSK Anda.

    • UntukmskTopicName, berikan nama topik Kafka Anda.

    • UntukmskGroupName, berikan nama grup konsumen Kafka Anda.

  2. Identifikasi Amazon MSK, Amazon EC2, dan CloudWatch izin yang diperlukan Lambda untuk menemukan dan menghubungkan cluster Amazon MSK Anda, dan catat peristiwa tersebut.

    Kebijakan ter AWSLambdaMSKExecutionRole kelola secara permisif mendefinisikan izin yang diperlukan. Gunakan dalam langkah-langkah berikut.

    Di lingkungan produksi, lakukan penilaian AWSLambdaMSKExecutionRole untuk membatasi kebijakan peran eksekusi berdasarkan prinsip hak istimewa terkecil, lalu tulis kebijakan untuk peran Anda yang menggantikan kebijakan terkelola ini.

Untuk detail tentang bahasa kebijakan IAM, lihat dokumentasi IAM.

Sekarang setelah Anda menulis dokumen kebijakan Anda, buat kebijakan IAM sehingga Anda dapat melampirkannya ke peran Anda. Anda dapat melakukan ini menggunakan konsol dengan prosedur berikut.

Untuk membuat kebijakan IAM dari dokumen kebijakan Anda
  1. Masuk ke Konsol Manajemen AWS dan buka konsol IAM di https://console.aws.amazon.com/iam/.

  2. Di panel navigasi sebelah kiri, pilih Kebijakan.

  3. Pilih Buat kebijakan.

  4. Di bagian Editor kebijakan, pilih opsi JSON.

  5. TempelclusterAuthPolicy.

  6. Setelah selesai menambahkan izin ke kebijakan, pilih Berikutnya.

  7. Pada halaman Tinjau dan buat, ketik Nama Kebijakan dan Deskripsi (opsional) untuk kebijakan yang Anda buat. Tinjau Izin yang ditentukan dalam kebijakan ini untuk melihat izin yang diberikan oleh kebijakan Anda.

  8. Pilih Buat kebijakan untuk menyimpan kebijakan baru Anda.

Untuk informasi selengkapnya, lihat Membuat kebijakan IAM di dokumentasi IAM.

Sekarang setelah Anda memiliki kebijakan IAM yang sesuai, buat peran dan lampirkan ke dalamnya. Anda dapat melakukan ini menggunakan konsol dengan prosedur berikut.

Untuk membuat peran eksekusi di konsol IAM
  1. Buka Halaman peran di konsol IAM.

  2. Pilih Buat peran.

  3. Di bawah Jenis entitas tepercaya, pilih AWS layanan.

  4. Di bawah Kasus penggunaan, pilih Lambda.

  5. Pilih Berikutnya.

  6. Pilih kebijakan berikut:

    • clusterAuthPolicy

    • AWSLambdaMSKExecutionRole

  7. Pilih Berikutnya.

  8. Untuk Nama peran, lambdaAuthRole masukkan lalu pilih Buat peran.

Untuk informasi selengkapnya, lihat Mendefinisikan izin fungsi Lambda dengan peran pelaksanaan.

Buat fungsi Lambda untuk membaca dari topik Amazon MSK Anda

Buat fungsi Lambda yang dikonfigurasi untuk menggunakan peran IAM Anda. Anda dapat membuat fungsi Lambda menggunakan konsol.

Untuk membuat fungsi Lambda menggunakan konfigurasi autentikasi Anda
  1. Buka konsol Lambda dan pilih Buat fungsi dari header.

  2. Pilih Penulis dari awal.

  3. Untuk Nama fungsi, berikan nama yang sesuai pilihan Anda.

  4. Untuk Runtime, pilih versi terbaru yang didukung Node.js untuk menggunakan kode yang disediakan dalam tutorial ini.

  5. Pilih Ubah peran eksekusi default.

  6. Pilih Gunakan peran yang ada.

  7. Untuk Peran yang ada, pilihlambdaAuthRole.

Dalam lingkungan produksi, Anda biasanya perlu menambahkan kebijakan lebih lanjut ke peran eksekusi untuk fungsi Lambda Anda untuk memproses peristiwa Amazon MSK Anda secara bermakna. Untuk informasi selengkapnya tentang menambahkan kebijakan ke peran Anda, lihat Menambahkan atau menghapus izin identitas dalam dokumentasi IAM.

Buat pemetaan sumber peristiwa ke fungsi Lambda Anda

Pemetaan sumber peristiwa Amazon MSK Anda menyediakan layanan Lambda informasi yang diperlukan untuk memanggil Lambda Anda ketika peristiwa Amazon MSK yang sesuai terjadi. Anda dapat membuat pemetaan Amazon MSK menggunakan konsol. Buat pemicu Lambda, maka pemetaan sumber peristiwa diatur secara otomatis.

Untuk membuat pemicu Lambda (dan pemetaan sumber peristiwa)
  1. Arahkan ke halaman ikhtisar fungsi Lambda Anda.

  2. Di bagian ikhtisar fungsi, pilih Tambahkan pemicu di kiri bawah.

  3. Di menu tar ik-turun Pilih sumber, pilih Amazon M SK.

  4. Jangan mengatur otentikasi.

  5. Untuk cluster MSK, pilih nama cluster Anda.

  6. Untuk Ukuran Batch, masukkan 1. Langkah ini membuat fitur ini lebih mudah untuk diuji, dan bukan nilai ideal dalam produksi.

  7. Untuk Nama Top ik, berikan nama topik Kafka Anda.

  8. Untuk ID grup Konsumen, berikan id grup konsumen Kafka Anda.

Perbarui fungsi Lambda Anda untuk membaca data streaming Anda

Lambda memberikan informasi tentang peristiwa Kafka melalui parameter metode peristiwa. Untuk contoh struktur acara Amazon MSK, lihatContoh peristiwa. Setelah memahami cara menafsirkan peristiwa Amazon MSK yang diteruskan Lambda, Anda dapat mengubah kode fungsi Lambda untuk menggunakan informasi yang mereka berikan.

Berikan kode berikut ke fungsi Lambda Anda untuk mencatat konten peristiwa Lambda Amazon MSK untuk tujuan pengujian:

.NET
SDK for .NET
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan.NET.

using System.Text; using Amazon.Lambda.Core; using Amazon.Lambda.KafkaEvents; // Assembly attribute to enable the Lambda function's JSON input to be converted into a .NET class. [assembly: LambdaSerializer(typeof(Amazon.Lambda.Serialization.SystemTextJson.DefaultLambdaJsonSerializer))] namespace MSKLambda; public class Function { /// <param name="input">The event for the Lambda function handler to process.</param> /// <param name="context">The ILambdaContext that provides methods for logging and describing the Lambda environment.</param> /// <returns></returns> public void FunctionHandler(KafkaEvent evnt, ILambdaContext context) { foreach (var record in evnt.Records) { Console.WriteLine("Key:" + record.Key); foreach (var eventRecord in record.Value) { var valueBytes = eventRecord.Value.ToArray(); var valueText = Encoding.UTF8.GetString(valueBytes); Console.WriteLine("Message:" + valueText); } } } }
Go
SDK untuk Go V2
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan Go.

package main import ( "encoding/base64" "fmt" "github.com/aws/aws-lambda-go/events" "github.com/aws/aws-lambda-go/lambda" ) func handler(event events.KafkaEvent) { for key, records := range event.Records { fmt.Println("Key:", key) for _, record := range records { fmt.Println("Record:", record) decodedValue, _ := base64.StdEncoding.DecodeString(record.Value) message := string(decodedValue) fmt.Println("Message:", message) } } } func main() { lambda.Start(handler) }
Java
SDK untuk Java 2.x
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan Java.

import com.amazonaws.services.lambda.runtime.Context; import com.amazonaws.services.lambda.runtime.RequestHandler; import com.amazonaws.services.lambda.runtime.events.KafkaEvent; import com.amazonaws.services.lambda.runtime.events.KafkaEvent.KafkaEventRecord; import java.util.Base64; import java.util.Map; public class Example implements RequestHandler<KafkaEvent, Void> { @Override public Void handleRequest(KafkaEvent event, Context context) { for (Map.Entry<String, java.util.List<KafkaEventRecord>> entry : event.getRecords().entrySet()) { String key = entry.getKey(); System.out.println("Key: " + key); for (KafkaEventRecord record : entry.getValue()) { System.out.println("Record: " + record); byte[] value = Base64.getDecoder().decode(record.getValue()); String message = new String(value); System.out.println("Message: " + message); } } return null; } }
JavaScript
SDK untuk JavaScript (v3)
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan menggunakan Lambda. JavaScript

exports.handler = async (event) => { // Iterate through keys for (let key in event.records) { console.log('Key: ', key) // Iterate through records event.records[key].map((record) => { console.log('Record: ', record) // Decode base64 const msg = Buffer.from(record.value, 'base64').toString() console.log('Message:', msg) }) } }

Mengkonsumsi acara Amazon MSK dengan menggunakan Lambda. TypeScript

import { MSKEvent, Context } from "aws-lambda"; import { Buffer } from "buffer"; import { Logger } from "@aws-lambda-powertools/logger"; const logger = new Logger({ logLevel: "INFO", serviceName: "msk-handler-sample", }); export const handler = async ( event: MSKEvent, context: Context ): Promise<void> => { for (const [topic, topicRecords] of Object.entries(event.records)) { logger.info(`Processing key: ${topic}`); // Process each record in the partition for (const record of topicRecords) { try { // Decode the message value from base64 const decodedMessage = Buffer.from(record.value, 'base64').toString(); logger.info({ message: decodedMessage }); } catch (error) { logger.error('Error processing event', { error }); throw error; } }; } }
PHP
SDK untuk PHP
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan PHP.

<?php // Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. // SPDX-License-Identifier: Apache-2.0 // using bref/bref and bref/logger for simplicity use Bref\Context\Context; use Bref\Event\Kafka\KafkaEvent; use Bref\Event\Handler as StdHandler; use Bref\Logger\StderrLogger; require __DIR__ . '/vendor/autoload.php'; class Handler implements StdHandler { private StderrLogger $logger; public function __construct(StderrLogger $logger) { $this->logger = $logger; } /** * @throws JsonException * @throws \Bref\Event\InvalidLambdaEvent */ public function handle(mixed $event, Context $context): void { $kafkaEvent = new KafkaEvent($event); $this->logger->info("Processing records"); $records = $kafkaEvent->getRecords(); foreach ($records as $record) { try { $key = $record->getKey(); $this->logger->info("Key: $key"); $values = $record->getValue(); $this->logger->info(json_encode($values)); foreach ($values as $value) { $this->logger->info("Value: $value"); } } catch (Exception $e) { $this->logger->error($e->getMessage()); } } $totalRecords = count($records); $this->logger->info("Successfully processed $totalRecords records"); } } $logger = new StderrLogger(); return new Handler($logger);
Python
SDK untuk Python (Boto3)
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan Python.

import base64 def lambda_handler(event, context): # Iterate through keys for key in event['records']: print('Key:', key) # Iterate through records for record in event['records'][key]: print('Record:', record) # Decode base64 msg = base64.b64decode(record['value']).decode('utf-8') print('Message:', msg)
Ruby
SDK untuk Ruby
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan Ruby.

require 'base64' def lambda_handler(event:, context:) # Iterate through keys event['records'].each do |key, records| puts "Key: #{key}" # Iterate through records records.each do |record| puts "Record: #{record}" # Decode base64 msg = Base64.decode64(record['value']) puts "Message: #{msg}" end end end
Rust
SDK for Rust
catatan

Ada lebih banyak tentang GitHub. Temukan contoh lengkapnya dan pelajari cara mengatur dan menjalankannya di repositori contoh Nirserver.

Mengkonsumsi acara Amazon MSK dengan Lambda menggunakan Rust.

use aws_lambda_events::event::kafka::KafkaEvent; use lambda_runtime::{run, service_fn, tracing, Error, LambdaEvent}; use base64::prelude::*; use serde_json::{Value}; use tracing::{info}; /// Pre-Requisites: /// 1. Install Cargo Lambda - see https://www.cargo-lambda.info/guide/getting-started.html /// 2. Add packages tracing, tracing-subscriber, serde_json, base64 /// /// This is the main body for the function. /// Write your code inside it. /// There are some code example in the following URLs: /// - https://github.com/awslabs/aws-lambda-rust-runtime/tree/main/examples /// - https://github.com/aws-samples/serverless-rust-demo/ async fn function_handler(event: LambdaEvent<KafkaEvent>) -> Result<Value, Error> { let payload = event.payload.records; for (_name, records) in payload.iter() { for record in records { let record_text = record.value.as_ref().ok_or("Value is None")?; info!("Record: {}", &record_text); // perform Base64 decoding let record_bytes = BASE64_STANDARD.decode(record_text)?; let message = std::str::from_utf8(&record_bytes)?; info!("Message: {}", message); } } Ok(().into()) } #[tokio::main] async fn main() -> Result<(), Error> { // required to enable CloudWatch error logging by the runtime tracing::init_default_subscriber(); info!("Setup CW subscriber!"); run(service_fn(function_handler)).await }

Anda dapat memberikan kode fungsi ke Lambda Anda menggunakan konsol.

Untuk memperbarui kode fungsi menggunakan editor kode konsol
  1. Buka halaman Functions pada konsol Lambda dan pilih fungsi Anda.

  2. Pilih tab Kode.

  3. Di panel sumber kode, pilih file kode sumber Anda dan edit di editor kode terintegrasi.

  4. Di bagian DEPLOY, pilih De ploy untuk memperbarui kode fungsi Anda:

    Menyebarkan tombol di editor kode konsol Lambda

Uji fungsi Lambda Anda untuk memverifikasi bahwa ia terhubung ke topik Amazon MSK Anda

Anda sekarang dapat memverifikasi apakah Lambda Anda dipanggil oleh sumber peristiwa dengan memeriksa log peristiwa. CloudWatch

Untuk memverifikasi apakah fungsi Lambda Anda sedang dipanggil
  1. Gunakan host admin Kafka Anda untuk menghasilkan acara Kafka menggunakan CLI. kafka-console-producer Untuk informasi selengkapnya, lihat Menulis beberapa peristiwa ke dalam topik dalam dokumentasi Kafka. Kirim acara yang cukup untuk mengisi batch yang ditentukan oleh ukuran batch untuk pemetaan sumber peristiwa yang ditentukan pada langkah sebelumnya, atau Lambda menunggu informasi lebih lanjut untuk dipanggil.

  2. Jika fungsi Anda berjalan, Lambda menulis apa yang terjadi. CloudWatch Di konsol, navigasikan ke halaman detail fungsi Lambda Anda.

  3. Pilih tab Konfigurasi.

  4. Dari sidebar, pilih Alat Pem antauan dan operasi.

  5. Identifikasi grup CloudWatch log di bawah konfigurasi Logging. Grup log harus dimulai dengan/aws/lambda. Pilih tautan ke grup log.

  6. Di CloudWatch konsol, periksa peristiwa Log untuk peristiwa log yang telah dikirim Lambda ke aliran log. Identifikasi apakah ada peristiwa log yang berisi pesan dari acara Kafka Anda, seperti pada gambar berikut. Jika ada, Anda telah berhasil menghubungkan fungsi Lambda ke Amazon MSK dengan pemetaan sumber peristiwa Lambda.

    Peristiwa log dalam CloudWatch menampilkan informasi peristiwa yang diekstraksi oleh kode yang disediakan.