Amazon MSK 및 자체 관리형 Apache Kafka 이벤트 소스의 버려진 배치 캡처

Amazon MSK 및 자체 관리형 Apache Kafka 이벤트 소스의 버려진 배치 캡처

실패한 이벤트 소스 매핑 호출의 기록을 보존하려면 함수의 이벤트 소스 매핑에 대상(destination)을 추가하세요. 대상으로 보내지는 각 레코드는 실패한 호출에 대한 메타데이터를 담은 JSON 문서입니다. Amazon S3 대상의 경우 Lambda는 메타데이터와 함께 전체 호출 레코드도 보냅니다. Amazon SNS 주제, Amazon SQS 큐, Amazon S3 버킷, 또는 Kafka 중 어떤 것이든 대상으로 구성할 수 있어요.

출처: AWS Lambda 개발자 안내서

본문

Amazon S3 대상을 사용하면 Amazon S3 이벤트 알림 기능을 사용해 대상 S3 버킷에 객체가 업로드될 때 알림을 받을 수 있어요. 또한 S3 이벤트 알림을 구성해 실패한 배치에 대해 자동 처리를 수행하는 다른 Lambda 함수를 호출할 수도 있습니다.

실행 역할에는 대상에 대한 권한이 있어야 해요:

Kafka 이벤트 소스 매핑의 on-failure 대상으로 Kafka 토픽을 구성할 수 있어요. Lambda가 재시도 시도를 소진한 후에도 레코드를 처리할 수 없거나 레코드가 최대 수명을 초과하면, Lambda는 실패한 레코드를 지정된 Kafka 토픽으로 보내 나중에 처리하게 합니다. on-failure 대상으로 Kafka 토픽 사용하기를 참고하세요.

Kafka 클러스터 VPC 내부에 on-failure 대상 서비스용 VPC 엔드포인트를 배포해야 합니다.

또한 대상에 KMS 키를 구성했다면 대상 유형에 따라 Lambda에 다음 권한이 필요합니다:

  • S3 대상에 자체 KMS 키로 암호화를 활성화했다면 kms:GenerateDataKey가 필요합니다. KMS 키와 S3 버킷 대상이 Lambda 함수·실행 역할과 다른 계정에 있다면, kms:GenerateDataKey를 허용하도록 실행 역할을 신뢰하게 KMS 키를 구성하세요.
  • SQS 대상에 자체 KMS 키로 암호화를 활성화했다면 kms:Decrypt와 kms:GenerateDataKey가 필요합니다. KMS 키와 SQS 큐 대상이 다른 계정에 있다면 kms:Decrypt, kms:GenerateDataKey, kms:DescribeKey, kms:ReEncrypt를 허용하도록 실행 역할을 신뢰하게 KMS 키를 구성하세요.
  • SNS 대상에 자체 KMS 키로 암호화를 활성화했다면 kms:Decrypt와 kms:GenerateDataKey가 필요합니다. KMS 키와 SNS 주제 대상이 다른 계정에 있다면 kms:Decrypt, kms:GenerateDataKey, kms:DescribeKey, kms:ReEncrypt를 허용하도록 실행 역할을 신뢰하게 KMS 키를 구성하세요.

Kafka 이벤트 소스 매핑의 on-failure 대상 구성하기

콘솔에서 on-failure 대상을 구성하려면 다음 단계를 따르세요:

  1. Lambda 콘솔의 Functions 페이지를 엽니다.
  2. 함수를 선택합니다.
  3. Function overview(함수 개요) 아래에서 Add destination(대상 추가)을 선택합니다.
  4. Source(소스)에서 Event source mapping invocation(이벤트 소스 매핑 호출)을 선택합니다.
  5. Event source mapping(이벤트 소스 매핑)에서 이 함수에 구성된 이벤트 소스를 선택합니다.
  6. Condition(조건)에서 On failure(실패 시)를 선택합니다. 이벤트 소스 매핑 호출의 경우 이것이 허용되는 유일한 조건이에요.
  7. Destination type(대상 유형)에서 Lambda가 호출 레코드를 보낼 대상 유형을 선택합니다.
  8. Destination(대상)에서 리소스를 선택합니다.
  9. Save(저장)을 선택합니다.

AWS CLI로도 on-failure 대상을 구성할 수 있어요. 예를 들어 다음 create-event-source-mapping 명령은 MyFunction에 SQS on-failure 대상을 가진 이벤트 소스 매핑을 추가합니다:

aws lambda create-event-source-mapping \
--function-name "MyFunction" \
--event-source-arn arn:aws:kafka:us-east-1:123456789012:cluster/vpc-2priv-2pub/751d2973-a626-431c-9d4e-d7975eb44dd7-2 \
--destination-config '{"OnFailure": {"Destination": "arn:aws:sqs:us-east-1:123456789012:dest-queue"}}'

다음 update-event-source-mapping 명령은 입력 uuid와 연결된 이벤트 소스에 S3 on-failure 대상을 추가합니다:

aws lambda update-event-source-mapping \
--uuid f89f8514-cdd9-4602-9e1f-01a5b77d449b \
--destination-config '{"OnFailure": {"Destination": "arn:aws:s3:::dest-bucket"}}'

대상을 제거하려면 destination-config 파라미터의 인자로 빈 문자열을 제공하세요:

aws lambda update-event-source-mapping \
--uuid f89f8514-cdd9-4602-9e1f-01a5b77d449b \
--destination-config '{"OnFailure": {"Destination": ""}}'

Amazon S3 대상의 보안 모범 사례

함수 구성에서 대상을 제거하지 않고 대상으로 구성된 S3 버킷을 삭제하면 보안 위험이 발생할 수 있어요. 다른 사용자가 여러분의 대상 버킷 이름을 알게 되면 그들은 자신의 AWS 계정에서 해당 버킷을 다시 만들 수 있습니다. 실패한 호출의 기록이 그들의 버킷으로 전송되어 함수의 데이터가 노출될 수 있어요.

경고

함수의 호출 기록이 다른 AWS 계정의 S3 버킷으로 전송되지 않도록 하려면 함수의 실행 역할에 s3:PutObject 권한을 여러분의 계정의 버킷으로 제한하는 조건을 추가하세요.

다음 예시는 함수의 s3:PutObject 권한을 계정의 버킷으로 제한하는 IAM 정책입니다. 이 정책은 Lambda가 S3 버킷을 대상으로 사용하는 데 필요한 s3:ListBucket 권한도 부여합니다.

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "S3BucketResourceAccountWrite",
            "Effect": "Allow",
            "Action": [
                "s3:PutObject",
                "s3:ListBucket"
            ],
            "Resource": [
                "arn:aws:s3:::*/*",
                "arn:aws:s3:::*"
            ],
            "Condition": {
                "StringEquals": {
                    "s3:ResourceAccount": "111122223333"
                }
            }
        }
    ]
}

AWS Management Console이나 AWS CLI로 함수의 실행 역할에 권한 정책을 추가하는 방법은 다음 절차의 지침을 참고하세요.

SNS 및 SQS 예시 호출 레코드

다음 예시는 실패한 Kafka 이벤트 소스 호출에 대해 Lambda가 SNS 주제 또는 SQS 큐 대상으로 보내는 내용을 보여줘요. recordsInfo 아래의 각 키에는 하이픈으로 구분된 Kafka 토픽과 파티션이 모두 들어 있습니다. 예를 들어 키 "Topic-0"에서 Topic은 Kafka 토픽이고 0은 파티션입니다. 각 토픽과 파티션에 대해 offsets와 timestamp 데이터를 사용해 원래 호출 레코드를 찾을 수 있어요.

{
    "requestContext": {
        "requestId": "316aa6d0-8154-xmpl-9af7-85d5f4a6bc81",
        "functionArn": "arn:aws:lambda:us-east-1:123456789012:function:myfunction",
        "condition": "RetryAttemptsExhausted" | "MaximumPayloadSizeExceeded",
        "approximateInvokeCount": 1
    },
    "responseContext": {
        // null if record is MaximumPayloadSizeExceeded
        "statusCode": 200,
        "executedVersion": "$LATEST",
        "functionError": "Unhandled"
    },
    "version": "1.0",
    "timestamp": "2019-11-14T00:38:06.021Z",
    "KafkaBatchInfo": {
        "batchSize": 500,
        "eventSourceArn": "arn:aws:kafka:us-east-1:123456789012:cluster/vpc-2priv-2pub/751d2973-a626-431c-9d4e-d7975eb44dd7-2",
        "bootstrapServers": "...",
        "payloadSize": 2039086, // In bytes
        "recordsInfo": {
            "Topic-0": {
                "firstRecordOffset": "49601189658422359378836298521827638475320189012309704722",
                "lastRecordOffset": "49601189658422359378836298522902373528957594348623495186",
                "firstRecordTimestamp": "2019-11-14T00:38:04.835Z",
                "lastRecordTimestamp": "2019-11-14T00:38:05.580Z"
            },
            "Topic-1": {
                "firstRecordOffset": "49601189658422359378836298521827638475320189012309704722",
                "lastRecordOffset": "49601189658422359378836298522902373528957594348623495186",
                "firstRecordTimestamp": "2019-11-14T00:38:04.835Z",
                "lastRecordTimestamp": "2019-11-14T00:38:05.580Z"
            }
        }
    }
}

S3 대상 예시 호출 레코드

S3 대상의 경우 Lambda는 메타데이터와 함께 전체 호출 레코드를 대상으로 보냅니다. 다음 예시는 실패한 Kafka 이벤트 소스 호출에 대해 Lambda가 S3 버킷 대상으로 보내는 내용을 보여줘요. SQS 및 SNS 대상에 대한 이전 예시의 모든 필드에 더해, payload 필드는 이스케이프된 JSON 문자열로 된 원래 호출 레코드를 포함합니다.

{
    "requestContext": {
        "requestId": "316aa6d0-8154-xmpl-9af7-85d5f4a6bc81",
        "functionArn": "arn:aws:lambda:us-east-1:123456789012:function:myfunction",
        "condition": "RetryAttemptsExhausted" | "MaximumPayloadSizeExceeded",
        "approximateInvokeCount": 1
    },
    "responseContext": {
        // null if record is MaximumPayloadSizeExceeded
        "statusCode": 200,
        "executedVersion": "$LATEST",
        "functionError": "Unhandled"
    },
    "version": "1.0",
    "timestamp": "2019-11-14T00:38:06.021Z",
    "KafkaBatchInfo": {
        "batchSize": 500,
        "eventSourceArn": "arn:aws:kafka:us-east-1:123456789012:cluster/vpc-2priv-2pub/751d2973-a626-431c-9d4e-d7975eb44dd7-2",
        "bootstrapServers": "...",
        "payloadSize": 2039086, // In bytes
        "recordsInfo": {
            "Topic-0": {
                "firstRecordOffset": "49601189658422359378836298521827638475320189012309704722",
                "lastRecordOffset": "49601189658422359378836298522902373528957594348623495186",
                "firstRecordTimestamp": "2019-11-14T00:38:04.835Z",
                "lastRecordTimestamp": "2019-11-14T00:38:05.580Z"
            },
            "Topic-1": {
                "firstRecordOffset": "49601189658422359378836298521827638475320189012309704722",
                "lastRecordOffset": "49601189658422359378836298522902373528957594348623495186",
                "firstRecordTimestamp": "2019-11-14T00:38:04.835Z",
                "lastRecordTimestamp": "2019-11-14T00:38:05.580Z"
            }
        }
    },
    "payload": "<Whole Event>" // Only available in S3
}

팁

대상 버킷에 S3 버전 관리(versioning)를 활성화할 것을 권장합니다.

더 알아보기 (Learn more)