Kafka 토픽을 on-failure 대상으로 사용

Kafka 토픽을 on-failure 대상으로 사용

Kafka 이벤트 소스 매핑의 on-failure 대상으로 Kafka 토픽을 구성할 수 있습니다. Lambda가 재시도 시도를 소진한 뒤에도 레코드를 처리할 수 없거나, 레코드가 최대 보관 시간을 초과하면 Lambda는 실패한 레코드를 지정된 Kafka 토픽으로 보내 나중에 처리합니다. 무한 재시도와 on-failure 대상을 모두 구성하면 Lambda는 자동으로 최대 10회 재시도를 적용합니다.

출처: AWS Lambda 개발자 안내서

본문

Kafka on-failure 대상 작동 방식

Kafka 토픽을 on-failure 대상으로 구성하면 Lambda가 Kafka 프로듀서로 작동해 실패한 레코드를 대상 토픽에 씁니다. 이는 Kafka 인프라 안에 데드 레터 토픽(DLT, dead letter topic) 패턴을 만듭니다.

  • 같은 클러스터 요건 – 대상 토픽은 소스 토픽과 같은 Kafka 클러스터에 존재해야 합니다.
  • 실제 레코드 내용 – Kafka 대상은 실패 메타데이터와 함께 실제 실패한 레코드를 받습니다.
  • 재귀 방지 – Lambda는 소스 토픽과 대상 토픽이 같은 구성을 차단해 무한 루프를 방지합니다.

Kafka on-failure 대상 구성

Kafka 이벤트 소스 매핑을 만들거나 갱신할 때 Kafka 토픽을 on-failure 대상으로 구성할 수 있습니다.

Kafka 대상 구성(콘솔)

  1. Lambda 콘솔의 Functions 페이지를 엽니다.
  2. 함수 이름을 선택합니다.
  3. 다음 중 하나를 수행합니다.
    • 새 Kafka 트리거를 추가하려면 Function overview 아래에서 Add trigger를 선택합니다.
    • 기존 Kafka 트리거를 수정하려면 트리거를 선택한 뒤 Edit을 선택합니다.
  4. Additional settings 아래의 On-failure destination에서 Kafka topic을 선택합니다.
  5. Topic name에 실패한 레코드를 보낼 Kafka 토픽의 이름을 입력합니다.
  6. Add 또는 Save를 선택합니다.

Kafka 대상 구성(AWS CLI)

kafka:// 접두사를 사용해 Kafka 토픽을 on-failure 대상으로 지정합니다.

Kafka 대상으로 이벤트 소스 매핑 생성

다음 예제는 Kafka 토픽을 on-failure 대상으로 하는 Amazon MSK 이벤트 소스 매핑을 만듭니다.

aws lambda create-event-source-mapping \
  --function-name my-kafka-function \
  --topics AWSKafkaTopic \
  --event-source-arn arn:aws:kafka:us-east-1:123456789012:cluster/my-cluster/abc123 \
  --starting-position LATEST \
  --provisioned-poller-config MinimumPollers=1,MaximumPollers=3 \
  --destination-config '{"OnFailure":{"Destination":"kafka://failed-records-topic"}}'

자체 관리형 Kafka의 경우 같은 구문을 사용합니다.

aws lambda create-event-source-mapping \
  --function-name my-kafka-function \
  --topics AWSKafkaTopic \
  --self-managed-event-source '{"Endpoints":{"KAFKA_BOOTSTRAP_SERVERS":["abc.xyz.com:9092"]}}' \
  --starting-position LATEST \
  --provisioned-poller-config MinimumPollers=1,MaximumPollers=3 \
  --destination-config '{"OnFailure":{"Destination":"kafka://failed-records-topic"}}'

Kafka 대상 갱신

update-event-source-mapping 명령으로 Kafka 대상을 추가하거나 수정합니다.

aws lambda update-event-source-mapping \
  --uuid 12345678-1234-1234-1234-123456789012 \
  --destination-config '{"OnFailure":{"Destination":"kafka://failed-records-topic"}}'

Kafka 대상의 레코드 형식

Lambda가 실패한 레코드를 Kafka 토픽으로 보낼 때 각 메시지에는 실패에 대한 메타데이터와 실제 레코드 내용이 모두 포함됩니다.

실패 메타데이터

메타데이터에는 레코드가 실패한 이유와 원래 배치에 대한 세부 정보가 포함됩니다.

{
  "requestContext": {
    "requestId": "e4b46cbf-b738-xmpl-8880-a18cdf61200e",
    "functionArn": "arn:aws:lambda:us-east-1:123456789012:function:my-function:$LATEST",
    "condition": "RetriesExhausted",
    "approximateInvokeCount": 3
  },
  "responseContext": {
    "statusCode": 200,
    "executedVersion": "$LATEST",
    "functionError": "Unhandled"
  },
  "version": "1.0",
  "timestamp": "2019-11-14T18:16:05.568Z",
  "KafkaBatchInfo": {
    "batchSize": 1,
    "eventSourceArn": "arn:aws:kafka:us-east-1:123456789012:cluster/my-cluster/abc123",
    "bootstrapServers": "b-1.mycluster.abc123.kafka.us-east-1.amazonaws.com:9098",
    "payloadSize": 1162,
    "recordInfo": {
      "offset": "49601189658422359378836298521827638475320189012309704722",
      "timestamp": "2019-11-14T18:16:04.835Z"
    }
  },
  "payload": {
    "bootstrapServers": "b-1.mycluster.abc123.kafka.us-east-1.amazonaws.com:9098",
    "eventSource": "aws:kafka",
    "eventSourceArn": "arn:aws:kafka:us-east-1:123456789012:cluster/my-cluster/abc123",
    "records": {
      "my-topic-0": [
        {
          "headers": [],
          "key": "dGVzdC1rZXk=",
          "offset": 100,
          "partition": 0,
          "timestamp": 1749116692330,
          "timestampType": "CREATE_TIME",
          "topic": "my-topic",
          "value": "dGVzdC12YWx1ZQ=="
        }
      ]
    }
  }
}

파티션 키 동작

Lambda는 대상 토픽에 생산할 때 원본 레코드의 같은 파티션 키를 사용합니다. 원본 레코드에 키가 없었다면 Lambda는 대상 토픽의 사용 가능한 모든 파티션에 Kafka의 기본 라운드 로빈 파티셔닝을 사용합니다.

요구 사항과 제한

  • 프로비저닝 모드 필요 – Kafka on-failure 대상은 프로비저닝 모드(provisioned mode)가 활성화된 이벤트 소스 매핑에서만 사용할 수 있습니다.
  • 같은 클러스터만 – 대상 토픽은 소스 토픽과 같은 Kafka 클러스터에 존재해야 합니다.
  • 토픽 권한 – 이벤트 소스 매핑에 대상 토픽에 대한 쓰기 권한이 있어야 합니다. 예시:
{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "ClusterPermissions",
            "Effect": "Allow",
            "Action": [
                "kafka-cluster:Connect",
                "kafka-cluster:DescribeCluster",
                "kafka-cluster:DescribeTopic",
                "kafka-cluster:WriteData",
                "kafka-cluster:ReadData"
            ],
            "Resource": [
                "arn:aws:kafka:*:*:cluster/*"
            ]
        },
        {
            "Sid": "TopicPermissions",
            "Effect": "Allow",
            "Action": [
                "kafka-cluster:DescribeTopic",
                "kafka-cluster:WriteData",
                "kafka-cluster:ReadData"
            ],
            "Resource": [
                "arn:aws:kafka:*:*:topic/*/*"
            ]
        },
        {
            "Effect": "Allow",
            "Action": [
                "kafka:DescribeCluster",
                "kafka:GetBootstrapBrokers",
                "kafka:Produce"
            ],
            "Resource": "arn:aws:kafka:*:*:cluster/*"
        },
        {
            "Effect": "Allow",
            "Action": [
                "ec2:CreateNetworkInterface",
                "ec2:DescribeNetworkInterfaces",
                "ec2:DeleteNetworkInterface",
                "ec2:DescribeSubnets",
                "ec2:DescribeSecurityGroups"
            ],
            "Resource": "*"
        }
    ]
}
  • 재귀 없음 – 대상 토픽 이름은 어떤 소스 토픽 이름과도 같을 수 없습니다.

더 알아보기 (Learn more)