Lambda에서 상태 유지형 Kinesis Data Streams 처리 구현
Lambda에서 상태 유지형 Kinesis Data Streams 처리 구현
Lambda 함수는 연속 스트림 처리 애플리케이션을 실행할 수 있어요. 스트림은 애플리케이션을 통해 계속 흐르는 무한한 데이터를 나타내죠. 이렇게 지속적으로 갱신되는 입력의 정보를 분석하려면 시간 기준으로 정의한 윈도우(window)로 포함된 레코드를 묶을 수 있어요.
텀블링 윈도우(tumbling window)는 정해진 간격으로 열리고 닫히는 뚜렷한 시간 윈도우예요. 기본적으로 Lambda 호출은 상태를 유지하지 않아서(stateless), 외부 데이터베이스 없이는 여러 연속 호출에 걸쳐 데이터를 처리하는 데 사용할 수 없어요. 하지만 텀블링 윈도우를 사용하면 호출 간에 상태를 유지할 수 있죠.
이 상태에는 현재 윈도우에 대해 이전에 처리된 메시지의 집계 결과가 담겨요. 상태는 샤드별로 최대 1MB일 수 있어요. 이 크기를 초과하면 Lambda가 윈도우를 일찍 종료해요.
스트림의 각 레코드는 특정 윈도우에 속해요. Lambda는 각 레코드를 최소 한 번 처리하지만, 각 레코드가 정확히 한 번만 처리되도록 보장하진 않아요. 오류 처리 같은 드문 경우 일부 레코드가 두 번 이상 처리될 수도 있어요. 레코드는 항상 첫 번째로 순서대로 처리돼요. 레코드가 두 번 이상 처리되면 순서가 어긋나게 처리될 수 있어요.
본문
집계 및 처리
사용자 관리 함수는 집계와 해당 집계의 최종 결과 처리 모두를 위해 호출돼요. Lambda는 윈도우에서 받은 모든 레코드를 집계해요. 이러한 레코드는 여러 배치로 받을 수 있으며, 각 배치는 별도의 호출이에요. 각 호출은 상태를 받아요.
따라서 텀블링 윈도우를 사용할 때 Lambda 함수 응답에는 state 속성이 포함되어야 해요. 응답에 state 속성이 없으면 Lambda는 이를 실패한 호출로 간주해요. 이 조건을 충족하려면 함수가 다음 JSON 형태의 TimeWindowEventResponse 객체를 반환할 수 있어요:
예제 TimeWindowEventResponse 값
{
"state": {
"1": 282,
"2": 715
},
"batchItemFailures": []
}
참고
Java 함수의 경우 상태를 나타내는 데 Map<String, String>을 사용하는 것을 권장해요.
윈도우 끝에서 isFinalInvokeForWindow 플래그가 true로 설정되어 최종 상태임을 나타내고 처리할 준비가 되었음을 알려줘요. 처리 후 윈도우가 완료되고 최종 호출이 완료되며, 상태는 버려져요.
윈도우 끝에서 Lambda는 집계 결과에 대한 작업을 위해 최종 처리를 사용해요. 최종 처리는 동기식으로 호출돼요. 성공적으로 호출된 후 함수는 시퀀스 번호를 체크포인트하고 스트림 처리가 계속돼요. 호출이 실패하면 Lambda 함수는 성공적인 호출이 있을 때까지 추가 처리를 중단해요.
예제 KinesisTimeWindowEvent
{
"Records": [
{
"kinesis": {
"kinesisSchemaVersion": "1.0",
"partitionKey": "1",
"sequenceNumber": "49590338271490256608559692538361571095921575989136588898",
"data": "SGVsbG8sIHRoaXMgaXMgYSB0ZXN0Lg==",
"approximateArrivalTimestamp": 1607497475.000
},
"eventSource": "aws:kinesis",
"eventVersion": "1.0",
"eventID": "shardId-000000000006:***",
"eventName": "aws:kinesis:record",
"invokeIdentityArn": "arn:aws:iam::123456789012:role/lambda-kinesis-role",
"awsRegion": "us-east-1",
"eventSourceARN": "arn:aws:kinesis:us-east-1:123456789012:stream/lambda-stream"
}
],
"window": {
"start": "2020-12-09T07:04:00Z",
"end": "2020-12-09T07:06:00Z"
},
"state": {
"1": 282,
"2": 715
},
"shardId": "shardId-000000000006",
"eventSourceARN": "arn:aws:kinesis:us-east-1:123456789012:stream/lambda-stream",
"isFinalInvokeForWindow": false,
"isWindowTerminatedEarly": false
}
구성
이벤트 소스 매핑을 만들거나 업데이트할 때 텀블링 윈도우를 구성할 수 있어요. 텀블링 윈도우를 구성하려면 윈도우를 초 단위로 지정하세요(TumblingWindowInSeconds). 다음 예제 AWS Command Line Interface(AWS CLI) 명령은 120초의 텀블링 윈도우가 있는 스트리밍 이벤트 소스 매핑을 만들어요. 집계와 처리를 위해 정의된 Lambda 함수 이름은 tumbling-window-example-function이에요.
aws lambda create-event-source-mapping \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/lambda-stream \
--function-name tumbling-window-example-function \
--starting-position TRIM_HORIZON \
--tumbling-window-in-seconds {{120}}
Lambda는 레코드가 스트림에 삽입된 시간을 기준으로 텀블링 윈도우 경계를 결정해요. 모든 레코드에는 Lambda가 경계 결정에 사용할 수 있는 대략적인 타임스탬프가 있어요.
텀블링 윈도우 집계는 리샤딩(resharding)을 지원하지 않아요. 샤드가 끝나면 Lambda는 현재 윈도우가 닫힌 것으로 간주하고, 모든 하위 샤드는 새로운 상태에서 자체 윈도우를 시작해요. 현재 윈도우에 새 레코드가 추가되지 않으면 Lambda는 윈도우가 끝났다고 가정하기 전에 최대 2분을 기다려요. 이렇게 하면 레코드가 간헐적으로 추가되어도 함수가 현재 윈도우의 모든 레코드를 읽도록 보장해요.
텀블링 윈도우는 기존 재시도 정책 maxRetryAttempts와 maxRecordAge를 완전히 지원해요.
예제 Handler.py – 집계 및 처리
다음 Python 함수는 집계 후 최종 상태를 처리하는 방법을 보여줘요:
def lambda_handler(event, context):
print('Incoming event: ', event)
print('Incoming state: ', event['state'])
#Check if this is the end of the window to either aggregate or process.
if event['isFinalInvokeForWindow']:
# logic to handle final state of the window
print('Destination invoke')
else:
print('Aggregate invoke')
#Check for early terminations
if event['isWindowTerminatedEarly']:
print('Window terminated early')
#Aggregation logic
state = event['state']
for record in event['Records']:
state[record['kinesis']['partitionKey']] = state.get(record['kinesis']['partitionKey'], 0) + 1
print('Returning state: ', state)
return {'state': state}
더 알아보기 (Learn more)
- 텀블링 윈도우로 Kinesis 스트림의 상태 유지형 집계·최종 처리를 구현하고,
state응답과isFinalInvokeForWindow플래그 활용을 익혀 보세요.