Kafka 토픽을 on-failure 대상으로 사용
Kafka 토픽을 on-failure 대상으로 사용
Kafka 이벤트 소스 매핑의 on-failure 대상으로 Kafka 토픽을 구성할 수 있습니다. Lambda가 재시도 시도를 소진한 뒤에도 레코드를 처리할 수 없거나, 레코드가 최대 보관 시간을 초과하면 Lambda는 실패한 레코드를 지정된 Kafka 토픽으로 보내 나중에 처리합니다. 무한 재시도와 on-failure 대상을 모두 구성하면 Lambda는 자동으로 최대 10회 재시도를 적용합니다.
본문
Kafka on-failure 대상 작동 방식
Kafka 토픽을 on-failure 대상으로 구성하면 Lambda가 Kafka 프로듀서로 작동해 실패한 레코드를 대상 토픽에 씁니다. 이는 Kafka 인프라 안에 데드 레터 토픽(DLT, dead letter topic) 패턴을 만듭니다.
- 같은 클러스터 요건 – 대상 토픽은 소스 토픽과 같은 Kafka 클러스터에 존재해야 합니다.
- 실제 레코드 내용 – Kafka 대상은 실패 메타데이터와 함께 실제 실패한 레코드를 받습니다.
- 재귀 방지 – Lambda는 소스 토픽과 대상 토픽이 같은 구성을 차단해 무한 루프를 방지합니다.
Kafka on-failure 대상 구성
Kafka 이벤트 소스 매핑을 만들거나 갱신할 때 Kafka 토픽을 on-failure 대상으로 구성할 수 있습니다.
Kafka 대상 구성(콘솔)
- Lambda 콘솔의 Functions 페이지를 엽니다.
- 함수 이름을 선택합니다.
- 다음 중 하나를 수행합니다.
- 새 Kafka 트리거를 추가하려면 Function overview 아래에서 Add trigger를 선택합니다.
- 기존 Kafka 트리거를 수정하려면 트리거를 선택한 뒤 Edit을 선택합니다.
- Additional settings 아래의 On-failure destination에서 Kafka topic을 선택합니다.
- Topic name에 실패한 레코드를 보낼 Kafka 토픽의 이름을 입력합니다.
- 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": "*"
}
]
}
- 재귀 없음 – 대상 토픽 이름은 어떤 소스 토픽 이름과도 같을 수 없습니다.