Lambda로 Amazon Kinesis Data Streams의 레코드 처리하기

Lambda로 Amazon Kinesis Data Streams의 레코드 처리하기

Lambda 함수를 사용해 Amazon Kinesis 데이터 스트림의 레코드를 처리할 수 있어요. Lambda 함수를 Kinesis Data Streams 공유 처리량 컨슈머(표준 반복자)에 매핑하거나, enhanced fan-out을 사용하는 전용 처리량 컨슈머에 매핑할 수 있습니다. 표준 반복자의 경우 Lambda는 HTTP 프로토콜을 사용해 Kinesis 스트림의 각 샤드에서 레코드를 폴링합니다. 이벤트 소스 매핑은 샤드의 다른 컨슈머와 읽기 처리량을 공유합니다.

출처: AWS Lambda 개발자 안내서

본문

Kinesis 데이터 스트림에 대한 자세한 내용은 Amazon Kinesis Data Streams에서 데이터 읽기를 참고하세요.

시작하려면 Lambda로 Amazon Kinesis Data Streams 레코드 처리하기를 참고하세요. 이벤트 소스 매핑 파라미터를 구성하고, 집계용으로 텀블링 윈도우를 사용하며, 예시 함수 코드를 확인할 수 있습니다.

참고

Kinesis는 샤드마다 비용을 청구하고, enhanced fan-out의 경우 스트림에서 읽은 데이터에 대해서도 비용을 청구합니다. 요금 세부 정보는 Amazon Kinesis 요금을 참고하세요.

스트림 폴링 및 배칭

Lambda는 데이터 스트림에서 레코드를 읽고, 스트림 레코드를 담은 이벤트와 함께 함수를 동기적으로 호출합니다. Lambda는 레코드를 배치로 읽어 함수를 호출해 배치의 레코드를 처리하게 합니다. 각 배치에는 단일 샤드/데이터 스트림의 레코드가 들어 있습니다.

여러분의 Lambda 함수는 데이터 스트림의 컨슈머 애플리케이션이에요. 함수는 각 샤드에서 배치 하나를 한 번에 처리합니다. Lambda 함수를 공유 처리량 컨슈머(표준 반복자)에 매핑하거나 enhanced fan-out이 있는 전용 처리량 컨슈머에 매핑할 수 있어요.

  • 표준 반복자: Lambda는 Kinesis 스트림의 각 샤드를 초당 1회의 기본 속도로 폴링해 레코드를 가져옵니다. 더 많은 레코드가 사용 가능하면 함수가 스트림을 따라잡을 때까지 Lambda는 배치 처리를 계속합니다. 이벤트 소스 매핑은 샤드의 다른 컨슈머와 읽기 처리량을 공유합니다.
  • Enhanced fan-out: 지연 시간을 최소화하고 읽기 처리량을 최대화하려면 enhanced fan-out으로 데이터 스트림 컨슈머를 만드세요. Enhanced fan-out 컨슈머는 각 샤드에 대한 전용 연결을 받으며, 스트림에서 읽는 다른 애플리케이션에 영향을 주지 않습니다. 스트림 컨슈머는 HTTP/2를 사용해 오래 지속되는 연결을 통해 레코드를 Lambda로 푸시하고 요청 헤더를 압축함으로써 지연 시간을 줄입니다. Kinesis RegisterStreamConsumer API로 스트림 컨슈머를 만들 수 있어요.
aws kinesis register-stream-consumer \
--consumer-name con1 \
--stream-arn arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream

다음과 같은 출력이 보일 거예요:

{
    "Consumer": {
        "ConsumerName": "con1",
        "ConsumerARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream/consumer/con1:1540591608",
        "ConsumerStatus": "CREATING",
        "ConsumerCreationTimestamp": 1540591608.0
    }
}

함수가 레코드를 처리하는 속도를 높이려면 데이터 스트림에 샤드를 추가하세요. Lambda는 각 샤드의 레코드를 순서대로 처리합니다. 함수가 오류를 반환하면 샤드에서 추가 레코드 처리를 중단합니다. 샤드가 많을수록 한 번에 처리되는 배치가 많아져 오류가 동시성에 미치는 영향이 줄어듭니다.

함수가 동시 배치 총량을 처리하도록 확장할 수 없다면 쿼터 증가를 요청하거나 함수에 동시성을 예약하세요.

기본적으로 Lambda는 레코드를 사용할 수 있게 되는 즉시 함수를 호출합니다. 이벤트 소스에서 Lambda가 읽는 배치에 레코드가 하나만 있으면 Lambda는 함수에 레코드를 하나만 보냅니다. 적은 수의 레코드로 함수가 호출되는 것을 피하려면 배칭 윈도우를 구성해 이벤트 소스가 최대 5분 동안 레코드를 버퍼링하게 할 수 있어요. 함수를 호출하기 전에 Lambda는 전체 배치를 모으거나, 배칭 윈도우가 만료되거나, 배치가 6MB 페이로드 한도에 도달할 때까지 이벤트 소스에서 레코드를 계속 읽습니다. 자세한 내용은 배칭 동작을 참고하세요.

경고

Lambda 이벤트 소스 매핑은 각 이벤트를 최소 한 번 처리하며, 레코드의 중복 처리가 발생할 수 있어요. 중복 이벤트로 인한 잠재적 문제를 피하려면 함수 코드를 멱등(idempotent)하게 만들 것을 강력히 권장합니다. 자세한 내용은 AWS Knowledge Center의 Lambda 함수를 멱등하게 만드는 방법을 참고하세요.

Lambda는 다음 배치를 처리하기 위해 보내기 전에 구성된 확장(extension)이 완료되기를 기다리지 않습니다. 즉, Lambda가 다음 레코드 배치를 처리하는 동안 확장이 계속 실행될 수 있어요. 이로 인해 계정의 동시성 설정이나 한도를 위반하면 스로틀링 문제가 발생할 수 있습니다. 이것이 잠재적 문제인지 감지하려면 함수를 모니터링하고 이벤트 소스 매핑에 대해 예상보다 높은 동시성 지표가 보이는지 확인하세요. 호출 사이의 시간이 짧기 때문에 Lambda는 샤드 수보다 높은 동시성 사용량을 잠시 보고할 수 있어요. 이는 확장이 없는 Lambda 함수에서도 마찬가지일 수 있습니다.

ParallelizationFactor 설정을 구성해 하나의 Kinesis 데이터 스트림 샤드를 여러 Lambda 호출로 동시에 처리할 수 있어요. 1(기본값)부터 10까지의 병렬화 팩터를 지정해 Lambda가 샤드에서 폴링하는 동시 배치 수를 정할 수 있습니다. 예를 들어 ParallelizationFactor를 2로 설정하면 100개의 Kinesis 데이터 샤드를 처리하기 위해 최대 200개의 동시 Lambda 호출을 가질 수 있습니다(실제로는 ConcurrentExecutions 지표에서 다른 값을 볼 수도 있어요). 이는 데이터 볼륨이 변동적이고 IteratorAge가 높을 때 처리 처리량을 확장하는 데 도움이 됩니다. 샤드당 동시 배치 수를 늘려도 Lambda는 파티션 키 수준에서 순서 처리를 보장합니다.

ParallelizationFactor는 Kinesis 집계와 함께 사용할 수도 있어요. 이벤트 소스 매핑의 동작은 enhanced fan-out 사용 여부에 따라 달라집니다:

  • Enhanced fan-out 없음: 집계된 이벤트 안의 모든 이벤트는 동일한 파티션 키를 가져야 합니다. 파티션 키는 집계된 이벤트의 파티션 키와도 일치해야 해요. 집계된 이벤트 안의 이벤트들이 서로 다른 파티션 키를 가지면 Lambda는 파티션 키별로 이벤트의 순서 처리를 보장할 수 없습니다.
  • Enhanced fan-out 사용: 먼저 Lambda는 집계된 이벤트를 개별 이벤트로 디코딩합니다. 집계된 이벤트는 포함된 이벤트와 다른 파티션 키를 가질 수 있어요. 그러나 파티션 키와 일치하지 않는 이벤트는 삭제되어 유실됩니다. Lambda는 이런 이벤트를 처리하지 않으며 구성된 실패 대상으로도 보내지 않습니다.

예시 이벤트

{
    "Records": [
        {
            "kinesis": {
                "kinesisSchemaVersion": "1.0",
                "partitionKey": "1",
                "sequenceNumber": "49590338271490256608559692538361571095921575989136588898",
                "data": "SGVsbG8sIHRoaXMgaXMgYSB0ZXN0Lg==",
                "approximateArrivalTimestamp": 1545084650.987
            },
            "eventSource": "aws:kinesis",
            "eventVersion": "1.0",
            "eventID": "shardId-000000000006:***",
            "eventName": "aws:kinesis:record",
            "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
            "awsRegion": "us-east-2",
            "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
        },
        {
            "kinesis": {
                "kinesisSchemaVersion": "1.0",
                "partitionKey": "1",
                "sequenceNumber": "49590338271490256608559692540925702759324208523137515618",
                "data": "VGhpcyBpcyBvbmx5IGEgdGVzdC4=",
                "approximateArrivalTimestamp": 1545084711.166
            },
            "eventSource": "aws:kinesis",
            "eventVersion": "1.0",
            "eventID": "shardId-000000000006:***",
            "eventName": "aws:kinesis:record",
            "invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-role",
            "awsRegion": "us-east-2",
            "eventSourceARN": "arn:aws:kinesis:us-east-2:123456789012:stream/lambda-stream"
        }
    ]
}

더 알아보기 (Learn more)