Lambda에서 상태 저장 DynamoDB 스트림 처리 구현하기

Lambda에서 상태 저장 DynamoDB 스트림 처리 구현하기

Lambda 함수는 지속적인 스트림 처리 애플리케이션을 실행할 수 있어요. 스트림은 애플리케이션을 통해 끊임없이 흐르는 무한한 데이터를 나타냅니다. 계속 갱신되는 입력의 정보를 분석하려면 시간으로 정의된 윈도우를 사용해 포함된 레코드를 경계 지을 수 있습니다.

출처: AWS Lambda 개발자 안내서

본문

텀블링 윈도우(tumbling window)는 일정한 간격으로 열리고 닫히는 뚜렷한 시간 윈도우예요. 기본적으로 Lambda 호출은 상태가 없습니다(stateless). 그래서 외부 데이터베이스 없이는 여러 연속 호출에 걸쳐 데이터를 처리하는 데 사용할 수 없습니다. 그러나 텀블링 윈도우를 사용하면 호출 전반에 걸쳐 상태를 유지할 수 있어요.

이 상태에는 현재 윈도우에 대해 이전에 처리된 메시지의 집계 결과가 들어 있습니다. 상태는 샤드당 최대 1MB일 수 있어요. 이 크기를 초과하면 Lambda는 윈도우를 조기 종료합니다.

스트림의 각 레코드는 특정 윈도우에 속합니다. Lambda는 각 레코드를 최소 한 번 처리하지만, 각 레코드가 한 번만 처리된다고 보장하지는 않아요. 오류 처리 같은 드문 경우에는 일부 레코드가 두 번 이상 처리될 수 있습니다. 레코드는 항상 처음에는 순서대로 처리됩니다. 레코드가 두 번 이상 처리되면 순서가 뒤섞여 처리될 수 있어요.

집계 및 처리

사용자 관리 함수는 집계와 그 집계의 최종 결과 처리 양쪽 모두에 호출됩니다. Lambda는 윈도우에서 받은 모든 레코드를 집계합니다. 이 레코드를 여러 배치로 받을 수 있고, 각각은 별도의 호출이에요. 각 호출은 상태(state)를 받습니다.

따라서 텀블링 윈도우를 사용할 때 Lambda 함수 응답에는 state 속성이 반드시 포함되어야 합니다. 응답에 state 속성이 없으면 Lambda는 이를 실패한 호출로 간주합니다. 이 조건을 충족하려면 함수가 다음 JSON 형태를 가진 TimeWindowEventResponse 객체를 반환하면 됩니다:

예시 TimeWindowEventResponse 값

{
    "state": {
        "1": 282,
        "2": 715
    },
    "batchItemFailures": []
}

참고

Java 함수의 경우 Map<String, String>으로 상태를 표현할 것을 권장합니다.

윈도우가 끝나면 isFinalInvokeForWindow 플래그가 true로 설정되어 이것이 최종 상태이며 처리가 준비되었음을 나타냅니다. 처리가 끝나면 윈도우가 완료되고, 최종 호출이 완료된 뒤 상태는 폐기됩니다.

윈도우가 끝나면 Lambda는 집계 결과에 대한 작업에 최종 처리를 사용합니다. 최종 처리는 동기적으로 호출됩니다. 호출이 성공하면 함수가 시퀀스 번호를 체크포인트하고 스트림 처리가 계속됩니다. 호출이 실패하면 Lambda 함수는 성공적인 호출이 있을 때까지 추가 처리를 중단합니다.

예시 DynamodbTimeWindowEvent

{
   "Records":[
      {
         "eventID":"1",
         "eventName":"INSERT",
         "eventVersion":"1.0",
         "eventSource":"aws:dynamodb",
         "awsRegion":"us-east-1",
         "dynamodb":{
            "Keys":{
               "Id":{
                  "N":"101"
               }
            },
            "NewImage":{
               "Message":{
                  "S":"New item!"
               },
               "Id":{
                  "N":"101"
               }
            },
            "SequenceNumber":"111",
            "SizeBytes":26,
            "StreamViewType":"NEW_AND_OLD_IMAGES"
         },
         "eventSourceARN":"stream-ARN"
      },
      {
         "eventID":"2",
         "eventName":"MODIFY",
         "eventVersion":"1.0",
         "eventSource":"aws:dynamodb",
         "awsRegion":"us-east-1",
         "dynamodb":{
            "Keys":{
               "Id":{
                  "N":"101"
               }
            },
            "NewImage":{
               "Message":{
                  "S":"This item has changed"
               },
               "Id":{
                  "N":"101"
               }
            },
            "OldImage":{
               "Message":{
                  "S":"New item!"
               },
               "Id":{
                  "N":"101"
               }
            },
            "SequenceNumber":"222",
            "SizeBytes":59,
            "StreamViewType":"NEW_AND_OLD_IMAGES"
         },
         "eventSourceARN":"stream-ARN"
      },
      {
         "eventID":"3",
         "eventName":"REMOVE",
         "eventVersion":"1.0",
         "eventSource":"aws:dynamodb",
         "awsRegion":"us-east-1",
         "dynamodb":{
            "Keys":{
               "Id":{
                  "N":"101"
               }
            },
            "OldImage":{
               "Message":{
                  "S":"This item has changed"
               },
               "Id":{
                  "N":"101"
               }
            },
            "SequenceNumber":"333",
            "SizeBytes":38,
            "StreamViewType":"NEW_AND_OLD_IMAGES"
         },
         "eventSourceARN":"stream-ARN"
      }
   ],
   "window":{
        "start": "2020-07-30T17:00:00Z",
        "end": "2020-07-30T17:05:00Z"
    },
   "state":{
        "1": "state1"
    },
   "shardId": "shard123456789",
   "eventSourceARN": "stream-ARN",
   "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:dynamodb:us-east-2:123456789012:table/my-table/stream/2024-06-10T19:26:16.525 \
--function-name tumbling-window-example-function \
--starting-position TRIM_HORIZON \
--tumbling-window-in-seconds 120

Lambda는 레코드가 스트림에 삽입된 시간을 기준으로 텀블링 윈도우 경계를 결정합니다. 모든 레코드에는 Lambda가 경계 결정에 사용하는 근사 타임스탬프가 있습니다.

텀블링 윈도우 집계는 리샤딩(resharding)을 지원하지 않아요. 샤드가 끝나면 Lambda는 윈도우가 닫힌 것으로 간주하고, 자식 샤드는 빈 상태에서 자신의 윈도우를 시작합니다.

텀블링 윈도우는 기존 재시도 정책 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['dynamodb']['NewImage']['Id']] = state.get(record['dynamodb']['NewImage']['Id'], 0) + 1

    print('Returning state: ', state)
    return {'state': state}

더 알아보기 (Learn more)