튜토리얼: Amazon MSK 이벤트 소스 매핑을 사용해 Lambda 함수 호출

튜토리얼: Amazon MSK 이벤트 소스 매핑을 사용해 Lambda 함수 호출

이 튜토리얼에서는 다음을 수행합니다.

출처: AWS Lambda 개발자 안내서

본문

  1. 기존 Amazon MSK 클러스터와 같은 AWS 계정에 Lambda 함수를 만듭니다.
  2. Lambda가 Amazon MSK와 통신하도록 네트워킹과 인증을 구성합니다.
  3. 토픽에 이벤트가 나타나면 Lambda 함수를 실행하는 Lambda Amazon MSK 이벤트 소스 매핑을 설정합니다.

이 단계들을 마치면 Amazon MSK에 이벤트가 전송될 때 자체 커스텀 Lambda 코드로 이벤트를 자동으로 처리하도록 Lambda 함수를 설정할 수 있습니다.

이 기능으로 무엇을 할 수 있나요?

예제 솔루션: MSK 이벤트 소스 매핑으로 고객에게 실시간 점수 전달

다음 시나리오를 고려해보세요. 회사가 고객이 스포츠 경기 같은 라이브 이벤트에 대한 정보를 볼 수 있는 웹 애플리케이션을 운영합니다. 경기의 정보 업데이트는 Amazon MSK의 Kafka 토픽을 통해 팀에 제공됩니다. MSK 토픽의 업데이트를 소비해 개발한 애플리케이션 안의 고객에게 라이브 이벤트의 갱신된 화면을 제공하는 솔루션을 설계하려고 합니다. 다음과 같은 설계 방식을 정했습니다. 클라이언트 애플리케이션은 AWS에 호스팅된 서버리스 백엔드와 통신합니다. 클라이언트는 Amazon API Gateway WebSocket API를 사용해 WebSocket 세션으로 연결합니다.

이 솔루션에서는 MSK 이벤트를 읽고, 그 이벤트를 애플리케이션 레이어용으로 준비하는 커스텀 로직을 수행한 뒤, 그 정보를 API Gateway API로 전달하는 구성 요소가 필요합니다. Lambda 함수에 커스텀 로직을 제공한 뒤 AWS Lambda Amazon MSK 이벤트 소스 매핑으로 호출해 이 구성 요소를 AWS Lambda로 구현할 수 있습니다.

Amazon API Gateway WebSocket API를 사용해 솔루션을 구현하는 방법에 대한 자세한 내용은 API Gateway 문서의 'WebSocket API tutorials'를 참고하세요.

사전 요구 사항

다음과 같이 미리 구성된 리소스가 있는 AWS 계정:

이 사전 요구 사항을 충족하려면 Amazon MSK 문서의 'Getting started using Amazon MSK'를 따를 것을 권장합니다.

  • Amazon MSK 클러스터. 'Getting started using Amazon MSK'의 'Create an Amazon MSK cluster'를 참고하세요. 다음 구성이 필요합니다.
    • 클러스터 보안 설정에서 IAM 역할 기반 인증(role-based authentication)이 Enabled인지 확인하세요. 이는 Lambda 함수가 필요한 Amazon MSK 리소스에만 접근하도록 제한해 보안을 개선합니다. 새 Amazon MSK 클러스터에서는 기본적으로 활성화됩니다.
    • 클러스터 네트워킹 설정에서 Public access가 off인지 확인하세요. Amazon MSK 클러스터의 인터넷 접근을 제한하면 데이터를 처리하는 중개자의 수를 제한해 보안이 개선됩니다. 새 Amazon MSK 클러스터에서는 기본적으로 활성화됩니다.
  • 이 솔루션에 사용할 Amazon MSK 클러스터의 Kafka 토픽. 'Getting started using Amazon MSK'의 'Create a topic'을 참고하세요.
  • Kafka 클러스터에서 정보를 검색하고 Kafka 이벤트를 토픽에 보내 테스트하는 Kafka admin 호스트(예: Kafka admin CLI와 Amazon MSK IAM 라이브러리를 설치한 Amazon EC2 인스턴스). 'Getting started using Amazon MSK'의 'Create a client machine'을 참고하세요.

이 리소스들을 설정한 뒤 계속 준비가 되었는지 확인하려면 AWS 계정에서 다음 정보를 수집하세요.

  • Amazon MSK 클러스터 이름. Amazon MSK 콘솔에서 찾을 수 있습니다.
  • 클러스터 UUID. Amazon MSK 클러스터 ARN의 일부로, Amazon MSK 콘솔에서 찾을 수 있습니다. Amazon MSK 문서의 'Listing clusters' 절차를 따라 확인하세요.
  • Amazon MSK 클러스터와 연결된 보안 그룹. Amazon MSK 콘솔에서 찾을 수 있습니다. 다음 단계에서 이것들을 clusterSecurityGroups라고 부릅니다.
  • Amazon MSK 클러스터를 포함하는 Amazon VPC의 ID. Amazon MSK 콘솔에서 Amazon MSK 클러스터와 연결된 서브넷을 식별한 뒤, Amazon VPC Console에서 그 서브넷과 연결된 Amazon VPC를 식별해 찾을 수 있습니다.
  • 솔루션에 사용되는 Kafka 토픽 이름. Kafka admin 호스트에서 Kafka topics CLI로 Amazon MSK 클러스터를 호출해 찾을 수 있습니다. topics CLI에 대한 자세한 내용은 Kafka 문서의 'Adding and removing topics'을 참고하세요.
  • Kafka 토픽의 소비자 그룹 이름. Lambda 함수에 사용하기 적합한 이름입니다. 이 그룹은 Lambda가 자동으로 만들 수 있으므로 Kafka CLI로 만들 필요는 없습니다. 소비자 그룹을 관리해야 한다면 consumer-groups CLI에 대해 더 알아보려면 Kafka 문서의 'Managing Consumer Groups'을 참고하세요.

AWS 계정의 다음 권한:

  • Lambda 함수를 만들고 관리할 권한.
  • IAM 정책을 만들고 Lambda 함수에 연결할 권한.
  • Amazon VPC 엔드포인트를 만들고 Amazon MSK 클러스터를 호스팅하는 Amazon VPC의 네트워킹 구성을 변경할 권한.

아직 AWS Command Line Interface를 설치하지 않았다면 'Installing or updating the latest version of the AWS CLI'의 단계를 따라 설치하세요.

튜토리얼은 명령을 실행할 커맨드라인 터미널 또는 셸이 필요합니다. Linux와 macOS에서는 원하는 셸과 패키지 관리자를 사용하세요.

Windows에서는 Lambda에서 흔히 쓰는 일부 Bash CLI 명령(예: zip)이 운영 체제의 내장 터미널에서 지원되지 않습니다. Windows 통합 버전의 Ubuntu와 Bash를 얻으려면 Windows Subsystem for Linux를 설치하세요.

Lambda가 Amazon MSK와 통신하도록 네트워크 연결 구성

AWS PrivateLink를 사용해 Lambda와 Amazon MSK를 연결합니다. Amazon VPC 콘솔에서 인터페이스 Amazon VPC 엔드포인트를 만들어 이렇게 할 수 있습니다. 네트워킹 구성에 대한 자세한 내용은 'Lambda용 Amazon MSK 클러스터 및 Amazon VPC 네트워크 구성'을 참고하세요.

Amazon MSK 이벤트 소스 매핑이 Lambda 함수를 대신해 실행될 때 Lambda 함수의 실행 역할을 가정합니다. 이 IAM 역할은 IAM으로 보호되는 리소스(예: Amazon MSK 클러스터)에 접근할 권한을 매핑에 부여합니다. 구성 요소들이 실행 역할을 공유하지만, Amazon MSK 매핑과 Lambda 함수는 다음 다이어그램처럼 각자의 작업에 대해 별도의 연결 요구 사항이 있습니다.

이벤트 소스 매핑은 Amazon MSK 클러스터 보안 그룹에 속합니다. 이 네트워킹 단계에서 Amazon MSK 클러스터 VPC에서 Lambda와 STS 서비스로 이벤트 소스 매핑을 연결하는 Amazon VPC 엔드포인트를 만드세요. Amazon MSK 클러스터 보안 그룹의 트래픽을 받도록 이 엔드포인트를 보호하세요. 그런 다음 이벤트 소스 매핑이 Amazon MSK 클러스터와 통신하도록 Amazon MSK 클러스터 보안 그룹을 조정하세요.

다음 단계는 AWS Management Console로 구성할 수 있습니다.

인터페이스 Amazon VPC 엔드포인트용 보안 그룹 만들기 – clusterSecurityGroups에서 포트 443의 인바운드 TCP 트래픽을 허용하는 endpointSecurityGroup 보안 그룹을 만드세요. Amazon EC2 문서의 'Create a security group' 절차를 따라 보안 그룹을 만든 뒤, 'Add rules to a security group' 절차를 따라 적절한 규칙을 추가하세요. 다음 정보로 보안 그룹을 만드세요.

  • 인바운드 규칙을 추가할 때 clusterSecurityGroups의 각 보안 그룹에 대해 규칙을 만드세요. 각 규칙에 대해:
    • Type에서 HTTPS를 선택합니다.
    • Source에서 clusterSecurityGroups 중 하나를 선택합니다.

Lambda 서비스를 Amazon VPC에 연결하는 엔드포인트 만들기 – Amazon MSK 클러스터를 포함하는 Amazon VPC에 Lambda 서비스를 연결하는 엔드포인트를 만드세요. 'Create an interface endpoint' 절차를 따라 다음 정보로 인터페이스 엔드포인트를 만드세요.

  • Service name에서 com.amazonaws.regionName.lambda를 선택합니다. 여기서 regionName은 Lambda 함수를 호스팅하는 리전입니다.
  • VPC에서 Amazon MSK 클러스터를 포함하는 Amazon VPC를 선택합니다.
  • Security groups에서 앞서 만든 endpointSecurityGroup을 선택합니다.
  • Subnets에서 Amazon MSK 클러스터를 호스팅하는 서브넷을 선택합니다.
  • Policy에서 lambda:InvokeFunction 액션에 대해 Lambda 서비스 보안 주체(principal)가 엔드포인트를 사용하도록 보호하는 다음 정책 문서를 제공합니다.
{
    "Statement": [
        {
            "Action": "lambda:InvokeFunction",
            "Effect": "Allow",
            "Principal": {
                "Service": [
                    "lambda.amazonaws.com"
                ]
            },
            "Resource": "*"
        }
    ]
}
  • Enable DNS name이 설정된 상태로 유지되는지 확인합니다.

AWS STS 서비스를 Amazon VPC에 연결하는 엔드포인트 만들기 – Amazon MSK 클러스터를 포함하는 Amazon VPC에 AWS STS 서비스를 연결하는 엔드포인트를 만드세요. 'Create an interface endpoint' 절차를 따라 다음 정보로 인터페이스 엔드포인트를 만드세요.

  • Service name에서 AWS STS를 선택합니다.
  • VPC에서 Amazon MSK 클러스터를 포함하는 Amazon VPC를 선택합니다.
  • Security groups에서 endpointSecurityGroup을 선택합니다.
  • Subnets에서 Amazon MSK 클러스터를 호스팅하는 서브넷을 선택합니다.
  • Policy에서 sts:AssumeRole 액션에 대해 Lambda 서비스 보안 주체가 엔드포인트를 사용하도록 보호하는 다음 정책 문서를 제공합니다.
{
    "Statement": [
        {
            "Action": "sts:AssumeRole",
            "Effect": "Allow",
            "Principal": {
                "Service": [
                    "lambda.amazonaws.com"
                ]
            },
            "Resource": "*"
        }
    ]
}
  • Enable DNS name이 설정된 상태로 유지되는지 확인합니다.

Amazon MSK 클러스터 보안 그룹 조정 – Amazon MSK 클러스터와 연결된 각 보안 그룹(clusterSecurityGroups에 있는)에 대해 다음을 허용하세요.

  • clusterSecurityGroups 전체(자기 자신 포함)에 대해 포트 9098의 모든 인바운드·아웃바운드 TCP 트래픽을 허용합니다.
  • 포트 443의 모든 아웃바운드 TCP 트래픽을 허용합니다.

이 트래픽 중 일부는 기본 보안 그룹 규칙으로 허용되므로, 클러스터가 단일 보안 그룹에 연결되어 있고 그 그룹에 기본 규칙이 있다면 추가 규칙은 필요하지 않습니다. 보안 그룹 규칙을 조정하려면 Amazon EC2 문서의 'Add rules to a security group' 절차를 따르세요. 다음 정보로 보안 그룹에 규칙을 추가하세요.

  • 포트 9098의 각 인바운드·아웃바운드 규칙에 대해:
    • Type에서 Custom TCP를 선택합니다.
    • Port range에 9098을 제공합니다.
    • Source에 clusterSecurityGroups 중 하나를 제공합니다.
  • 포트 443의 각 인바운드 규칙에 대해 Type에서 HTTPS를 선택합니다.

Lambda가 Amazon MSK 토픽에서 읽도록 IAM 역할 만들기

Lambda가 Amazon MSK 토픽에서 읽기 위한 인증 요구 사항을 식별한 뒤 정책에 정의합니다. Lambda에 해당 권한을 사용하도록 승인하는 역할 lambdaAuthRole을 만드세요. kafka-cluster IAM 액션으로 Amazon MSK 클러스터의 작업을 승인합니다. 그런 다음 Lambda가 Amazon MSK 클러스터를 발견·연결하는 데 필요한 Amazon MSK kafka와 Amazon EC2 액션, Lambda가 한 일을 로그할 수 있도록 CloudWatch 액션을 승인합니다.

클러스터 인증 정책 작성 – clusterAuthPolicy라는 IAM 정책 문서(JSON 문서)를 작성해 Lambda가 Kafka 소비자 그룹을 사용해 Amazon MSK 클러스터의 Kafka 토픽에서 읽을 수 있게 합니다. Lambda는 읽을 때 Kafka 소비자 그룹이 설정되어 있어야 합니다. 다음 템플릿을 사전 요구 사항에 맞게 수정하세요.

{
    "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"
            ]
        }
    ]
}

자세한 내용은 'Amazon MSK 이벤트 소스 매핑용 Lambda 권한 구성'을 참고하세요. 정책을 작성할 때:

  • us-east-1과 111122223333을 Amazon MSK 클러스터의 AWS 리전과 AWS 계정으로 바꾸세요.
  • mskClusterName에 Amazon MSK 클러스터의 이름을 제공하세요.
  • cluster-uuid에 Amazon MSK 클러스터 ARN의 UUID를 제공하세요.
  • mskTopicName에 Kafka 토픽의 이름을 제공하세요.
  • mskGroupName에 Kafka 소비자 그룹의 이름을 제공하세요.

MSK 실행 역할 정책 식별 – Lambda가 Amazon MSK 클러스터를 발견·연결하고 이벤트를 로그하는 데 필요한 Amazon MSK, Amazon EC2, CloudWatch 권한을 식별하세요. AWSLambdaMSKExecutionRole 관리형 정책이 필요한 권한을 관대하게 정의합니다. 다음 단계에서 이를 사용하세요. 프로덕션 환경에서는 최소 권한 원칙에 따라 실행 역할 정책을 제한하도록 AWSLambdaMSKExecutionRole을 평가한 뒤, 이 관리형 정책을 대체하는 역할 정책을 작성하세요. IAM 정책 언어에 대한 자세한 내용은 IAM 문서를 참고하세요.

IAM 정책 만들기 – 정책 문서를 작성했으니 역할에 연결할 수 있는 IAM 정책을 만드세요. 콘솔로 다음 절차를 사용할 수 있습니다.

  1. AWS Management Console에 로그인하고 https://console.aws.amazon.com/iam/에서 IAM 콘솔을 엽니다.
  2. 왼쪽 탐색 창에서 Policies를 선택합니다.
  3. Create policy를 선택합니다.
  4. Policy editor 섹션에서 JSON 옵션을 선택합니다.
  5. clusterAuthPolicy를 붙여넣습니다.
  6. 정책에 권한 추가를 마쳤으면 Next를 선택합니다.
  7. Review and create 페이지에서 만들 정책의 Policy Name과 Description(선택)을 입력합니다. 이 정책이 부여하는 권한을 보려면 Permissions defined in this policy를 검토하세요.
  8. Create policy를 선택해 새 정책을 저장합니다.

자세한 내용은 IAM 문서의 'Creating IAM policies'를 참고하세요.

IAM 역할 만들기 – 적절한 IAM 정책을 갖추었으니 역할을 만들고 정책을 연결하세요. 콘솔로 다음 절차를 사용할 수 있습니다.

  1. IAM 콘솔의 Roles 페이지를 엽니다.
  2. Create role을 선택합니다.
  3. Trusted entity type 아래에서 AWS service를 선택합니다.
  4. Use case 아래에서 Lambda를 선택합니다.
  5. Next를 선택합니다.
  6. 다음 정책을 선택합니다.
    • clusterAuthPolicy
    • AWSLambdaMSKExecutionRole
  7. Next를 선택합니다.
  8. Role name에 lambdaAuthRole을 입력한 뒤 Create role을 선택합니다.

자세한 내용은 '실행 역할로 Lambda 함수 권한 정의'를 참고하세요.

Amazon MSK 토픽에서 읽도록 Lambda 함수 만들기

IAM 역할을 사용하도록 구성된 Lambda 함수를 만드세요. 콘솔로 Lambda 함수를 만들 수 있습니다.

  1. Lambda 콘솔을 열고 헤더에서 Create function을 선택합니다.
  2. Author from scratch를 선택합니다.
  3. Function name에 원하는 적절한 이름을 제공합니다.
  4. Runtime에서 이 튜토리얼의 코드를 사용할 가장 최신 지원 버전의 Node.js를 선택합니다.
  5. Change default execution role을 선택합니다.
  6. Use an existing role을 선택합니다.
  7. Existing role에서 lambdaAuthRole을 선택합니다.

프로덕션 환경에서는 보통 Lambda 함수가 Amazon MSK 이벤트를 의미 있게 처리하도록 실행 역할에 추가 정책을 더해야 합니다. 역할에 정책을 추가하는 방법에 대한 자세한 내용은 IAM 문서의 'Add or remove identity permissions'을 참고하세요.

Lambda 함수에 이벤트 소스 매핑 만들기

Amazon MSK 이벤트 소스 매핑은 적절한 Amazon MSK 이벤트가 발생할 때 Lambda를 호출하는 데 필요한 정보를 Lambda 서비스에 제공합니다. 콘솔로 Amazon MSK 매핑을 만들 수 있습니다. Lambda 트리거를 만들면 이벤트 소스 매핑이 자동으로 설정됩니다.

  1. Lambda 함수의 개요 페이지로 이동합니다.
  2. 함수 개요 섹션에서 왼쪽 아래의 Add trigger를 선택합니다.
  3. Select a source 드롭다운에서 Amazon MSK를 선택합니다.
  4. 인증은 설정하지 마세요.
  5. MSK cluster에서 클러스터 이름을 선택합니다.
  6. Batch size에 1을 입력합니다. 이 단계는 이 기능을 테스트하기 더 쉽게 하기 위한 것으로, 프로덕션에는 이상적인 값이 아닙니다.
  7. Topic name에 Kafka 토픽의 이름을 제공합니다.
  8. Consumer group ID에 Kafka 소비자 그룹의 id를 제공합니다.

스트리밍 데이터를 읽도록 Lambda 함수 갱신

Lambda는 event 메서드 파라미터를 통해 Kafka 이벤트에 대한 정보를 제공합니다. Amazon MSK 이벤트 예제 구조는 '예제 이벤트'를 참고하세요. Lambda가 전달한 Amazon MSK 이벤트를 해석하는 방법을 이해한 뒤에는 제공하는 정보를 사용하도록 Lambda 함수 코드를 변경할 수 있습니다.

테스트 목적으로 Lambda Amazon MSK 이벤트의 내용을 로그하기 위해 Lambda 함수에 다음 코드를 제공하세요.

GitHub에 더 많은 내용이 있습니다. 완전한 예제와 설정·실행 방법은 Serverless examples 저장소에서 확인하세요.

.NET으로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

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)
      }) 
    }
}

TypeScript로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

<?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으로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

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로 Lambda에서 Amazon MSK 이벤트 소비 :

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
}

콘솔로 Lambda에 함수 코드를 제공할 수 있습니다.

  1. Lambda 콘솔의 Functions 페이지를 열고 함수를 선택합니다.
  2. Code 탭을 선택합니다.
  3. Code source 창에서 소스 코드 파일을 선택하고 통합 코드 편집기에서 편집합니다.
  4. DEPLOY 섹션에서 Deploy를 선택해 함수의 코드를 갱신합니다.

Lambda 함수를 테스트해 Amazon MSK 토픽에 연결되었는지 확인

이제 CloudWatch 이벤트 로그를 검사해 Lambda가 이벤트 소스에 의해 호출되고 있는지 여부를 확인할 수 있습니다.

Kafka admin 호스트에서 kafka-console-producer CLI로 Kafka 이벤트를 생성하세요. 자세한 내용은 Kafka 문서의 'Write some events into the topic'을 참고하세요. 이전 단계에서 정의한 이벤트 소스 매핑의 배치 크기로 정의된 배치를 채울 충분한 이벤트를 보내세요. 그렇지 않으면 Lambda는 호출할 더 많은 정보를 기다립니다.

함수가 실행되면 Lambda는 일어난 일을 CloudWatch에 기록합니다. 콘솔에서 Lambda 함수의 상세 페이지로 이동하세요.

  1. Configuration 탭을 선택합니다.
  2. 사이드바에서 Monitoring and operations tools를 선택합니다.
  3. Logging configuration 아래의 CloudWatch 로그 그룹을 식별합니다. 로그 그룹은 /aws/lambda로 시작해야 합니다. 로그 그룹의 링크를 선택합니다.
  4. CloudWatch 콘솔에서 Lambda가 로그 스트림에 보낸 로그 이벤트에 대한 Log events를 검사합니다. 다음 이미지처럼 Kafka 이벤트의 메시지를 포함한 로그 이벤트가 있는지 식별하세요. 있다면 Lambda 이벤트 소스 매핑으로 Lambda 함수를 Amazon MSK에 성공적으로 연결한 것입니다.

더 알아보기 (Learn more)