AWS용 샘플 비동기 원격 서비스

AWS용 샘플 비동기 원격 서비스 (Sample Asynchronous Remote Service for AWS)

이 문서는 AWS용 비동기(asynchronous) AWS Lambda Function(원격 서비스) 샘플을 제공해 드려요. 1단계: Management Console에서 원격 서비스(AWS Lambda function) 만들기 에 설명된 것과 같은 단계로 이 샘플 함수를 만들 수 있어요. 비동기 원격 서비스가 직면하는 제약과, POST 핸들러·데이터 처리·GET 핸들러 세 프로세스가 공유 저장소로 협력하는 구조를 설명하고 샘플 코드를 보여줘요.

출처: Snowflake SQL Reference

본문

이 문서는 샘플 비동기 AWS Lambda Function(원격 서비스)을 포함해요. 1단계: Management Console에서 원격 서비스(AWS Lambda function) 만들기 에 설명된 것과 같은 단계로 이 샘플 함수를 만들 수 있어요.

코드 개요

이 문서 섹션은 AWS에서 비동기 외부 함수를 만드는 방법에 대한 정보를 제공해요. (첫 번째 비동기 외부 함수를 구현하기 전에, 비동기 외부 함수의 개념적 개요 를 읽는 것이 좋아요.)

AWS에서 비동기 원격 서비스는 다음 제약을 극복해야 해요:

  • HTTP POST와 GET은 별도의 요청이므로, 원격 서비스는 나중에 GET 요청이 상태를 조회할 수 있도록 POST 요청이 시작한 워크플로에 대한 정보를 유지해야 해요.

보통 각 HTTP POST와 HTTP GET은 별도의 프로세스나 스레드에서 핸들러 함수의 별도 인스턴스를 호출해요. 별도의 인스턴스는 메모리를 공유하지 않아요. GET 핸들러가 상태나 처리된 데이터를 읽으려면, GET 핸들러는 AWS에서 사용할 수 있는 공유 저장소 리소스에 접근해야 해요.

  • POST 핸들러가 초기 HTTP 202 응답 코드를 보내는 유일한 방법은 return 문(또는 동등한 것)을 통해서예요. 이는 핸들러의 실행을 종료해요. 따라서 HTTP 202를 반환하기 전에, POST 핸들러는 원격 서비스의 실제 데이터 처리 작업을 수행할 독립 프로세스(또는 스레드)를 시작해야 해요. 이 독립 프로세스는 보통 GET 핸들러가 볼 수 있는 저장소에 접근해야 해요.

비동기 원격 서비스가 이러한 제약을 극복하는 한 가지 방법은 3개의 프로세스(또는 스레드)와 공유 저장소를 사용하는 거예요:

이 모델에서 프로세스는 다음과 같은 책임을 가져요:

  • HTTP POST 핸들러:

    • 입력 데이터를 읽어요. Lambda Function에서 이는 핸들러 함수의 event 입력 파라미터의 body에서 읽어요.
    • 배치 ID를 읽어요. Lambda Function에서 이는 event 입력 파라미터의 헤더에서 읽어요.
    • 데이터 처리 프로세스를 시작하고, 그 프로세스에 데이터와 배치 ID를 전달해요. 데이터는 보통 호출 중에 전달되지만, 외부 저장소에 써서 전달할 수도 있어요.
    • 데이터 처리 프로세스와 HTTP GET 핸들러 프로세스 모두가 접근할 수 있는 공유 저장소에 배치 ID를 기록해요.
    • 필요하면 이 배치의 처리가 아직 끝나지 않았음을 기록해요.
    • 오류가 감지되지 않으면 HTTP 202를 반환해요.
  • 데이터 처리 코드:

    • 입력 데이터를 읽어요.
    • 데이터를 처리해요.
    • 결과를 GET 핸들러가 사용할 수 있게 해요 (결과 데이터를 공유 저장소에 쓰거나, 결과를 조회할 API를 제공하는 방식으로).
    • 보통 이 배치의 상태를 업데이트해요 (예: IN_PROGRESS 에서 SUCCESS 로) — 결과를 읽을 준비가 되었다는 뜻이에요.
    • 종료해요. 선택적으로 이 프로세스는 오류 표시를 반환할 수 있어요. Snowflake는 이를 직접 보지는 못해요 (Snowflake는 POST 핸들러와 GET 핸들러의 HTTP 반환 코드만 봐요). 하지만 데이터 처리 프로세스에서 오류 표시를 반환하면 디버깅에 도움이 될 수 있어요.
  • GET 핸들러:

    • 배치 ID를 읽어요. Lambda Function에서 이는 event 입력 파라미터의 헤더에서 읽어요.

    • 저장소를 읽어 이 배치의 현재 상태를 얻어요 (예: IN_PROGRESS 또는 SUCCESS).

    • 처리가 아직 진행 중이면 202를 반환해요.

    • 처리가 성공적으로 끝났다면:

      • 결과를 읽어요.
      • 저장소를 정리해요.
      • 결과를 HTTP 코드 200과 함께 반환해요.
    • 저장된 상태가 오류를 나타내면:

      • 저장소를 정리해요.
      • 오류 코드를 반환해요.

처리가 충분히 오래 걸려 여러 HTTP GET 요청이 전송된다면, GET 핸들러는 배치에 대해 여러 번 호출될 수 있다는 점에 주의해요.

이 모델에는 많은 변형이 가능해요. 예를 들어:

  • 배치 ID와 상태는 POST 프로세스의 끝이 아니라 데이터 처리 프로세스의 시작에 기록될 수도 있어요.
  • 데이터 처리는 별도의 함수(예: 별도의 Lambda function)에서, 또는 완전히 별개의 서비스로 수행될 수도 있어요.
  • 데이터 처리 코드가 반드시 공유 저장소에 쓸 필요는 없어요. 대신 처리된 데이터를 다른 방식으로 사용할 수 있게 할 수도 있어요. 예를 들어 API가 배치 ID를 파라미터로 받아 데이터를 반환할 수 있어요.

구현 코드는 처리가 너무 오래 걸리거나 실패할 가능성을 고려해야 해요. 따라서 저장 공간 낭비를 피하려면 부분 결과를 정리해야 해요.

저장 메커니즘은 여러 프로세스(또는 스레드)가 공유할 수 있어야 해요. 가능한 저장 메커니즘은 다음과 같아요:

  • AWS가 제공하는 저장 메커니즘. 예:

  • AWS 바깥에 있지만 AWS에서 접근할 수 있는 저장소.

위의 3개 프로세스 각각의 코드는 3개의 별도 Lambda Functions(POST 핸들러용, 데이터 처리 함수용, GET 핸들러용)로 작성하거나, 하나의 함수를 다른 방식으로 호출하도록 작성할 수 있어요.

아래 샘플 Python 코드는 POST, 데이터 처리, GET 프로세스에 대해 별도로 호출할 수 있는 단일 Lambda Function이에요.

샘플 코드

이 코드는 샘플 쿼리와 출력을 보여줘요. 이 예제의 초점은 세 프로세스와 그 상호 작용 방식이지, 공유 저장 메커니즘(DynamoDB)이나 데이터 변환(감성 분석)이 아니에요. 코드는 예제 저장 메커니즘과 데이터 변환을 다른 것으로 쉽게 바꿀 수 있도록 구조화되어 있어요.

단순함을 위해 이 예제는:

  • 일부 중요한 값(예: AWS 리전)을 하드코딩해요.
  • 일부 리소스(예: Dynamo의 Jobs 테이블)가 존재한다고 가정해요.
import json
import time
import boto3

HTTP_METHOD_STRING = "httpMethod"
HEADERS_STRING = "headers"
BATCH_ID_STRING = "sf-external-function-query-batch-id"
DATA_STRING = "data"
REGION_NAME = "us-east-2"

TABLE_NAME = "Jobs"
IN_PROGRESS_STATUS = "IN_PROGRESS"
SUCCESS_STATUS = "SUCCESS"

def lambda_handler(event, context):
    # this is called from either the GET or POST
    if (HTTP_METHOD_STRING in event):
        method = event[HTTP_METHOD_STRING]
        if method == "POST":
            return initiate(event, context)
        elif method == "GET":
            return poll(event, context)
        else:
            return create_response(400, "Function called from invalid method")

    # if not called from GET or POST, then this lambda was called to
    # process data
    else:
        return process_data(event, context)


# Reads batch_ID and data from the request, marks the batch_ID as being processed, and
# starts the processing service.
def initiate(event, context):
    batch_id = event[HEADERS_STRING][BATCH_ID_STRING]
    data = json.loads(event["body"])[DATA_STRING]

    lambda_name = context.function_name

    write_to_storage(batch_id, IN_PROGRESS_STATUS, "NULL")
    lambda_response = invoke_process_lambda(batch_id, data, lambda_name)

    # lambda response returns 202, because we are invoking it with
    # InvocationType = 'Event'
    if lambda_response["StatusCode"] != 202:
        response = create_response(400, "Error in initiate: processing lambda not started")
    else:
        response = {
            'statusCode': lambda_response["StatusCode"]
        }

    return response


# Processes the data passed to it from the POST handler. In this example,
# the processing is to perform sentiment analysis on text.
def process_data(event, context):
    data = event[DATA_STRING]
    batch_id = event[BATCH_ID_STRING]

    def process_data_impl(data):
        comprehend = boto3.client(service_name='comprehend', region_name=REGION_NAME)
        # create return rows
        ret = []
        for i in range(len(data)):
            text = data[i][1]
            sentiment_response = comprehend.detect_sentiment(Text=text, LanguageCode='en')
            sentiment_score = json.dumps(sentiment_response['SentimentScore'])
            ret.append([i, sentiment_score])
        return ret

    processed_data = process_data_impl(data)
    write_to_storage(batch_id, SUCCESS_STATUS, processed_data)

    return create_response(200, "No errors in process")


# Repeatedly checks on the status of the batch_ID, and returns the result after the
# processing has been completed.
def poll(event, context):
    batch_id = event[HEADERS_STRING][BATCH_ID_STRING]
    processed_data = read_data_from_storage(batch_id)

    def parse_processed_data(response):
        # in this case, the response is the response from DynamoDB
        response_metadata = response['ResponseMetadata']
        status_code = response_metadata['HTTPStatusCode']

        # Take action depending on item status
        item = response['Item']
        job_status = item['status']
        if job_status == SUCCESS_STATUS:
            # the row number is stored at index 0 as a Decimal object,
            # we need to convert it into a normal int to be serialized to JSON
            data = [[int(row[0]), row[1]] for row in item['data']]
            return {
                'statusCode': 200,
                'body': json.dumps({
                    'data': data
                })
            }
        elif job_status == IN_PROGRESS_STATUS:
            return {
                'statusCode': 202,
                "body": "{}"
            }
        else:
            return create_response(500, "Error in poll: Unknown item status.")

    return parse_processed_data(processed_data)


def create_response(code, msg):
    return {
        'statusCode': code,
        'body': msg
    }


def invoke_process_lambda(batch_id, data, lambda_name):
    # Create payload to be sent to processing lambda
    invoke_payload = json.dumps({
        BATCH_ID_STRING: batch_id,
        DATA_STRING: data
    })

    # Invoke processing lambda asynchronously by using InvocationType='Event'.
    # This allows the processing to continue while the POST handler returns HTTP 202.
    lambda_client = boto3.client('lambda', region_name=REGION_NAME,)
    lambda_response = lambda_client.invoke(
        FunctionName=lambda_name,
        InvocationType='Event',
        Payload=invoke_payload
    )
    # returns 202 on success if InvocationType = 'Event'
    return lambda_response


def write_to_storage(batch_id, status, data):
    # we assume that the table has already been created
    client = boto3.resource('dynamodb')
    table = client.Table(TABLE_NAME)

    # Put in progress item in table
    item_to_store = {
        'batch_id': batch_id,
        'status': status,
        'data': data,
        'timestamp': "{}".format(time.time())
    }
    db_response = table.put_item(
        Item=item_to_store
    )


def read_data_from_storage(batch_id):
    # we assume that the table has already been created
    client = boto3.resource('dynamodb')
    table = client.Table(TABLE_NAME)

    response = table.get_item(Key={'batch_id': batch_id},
                          ConsistentRead=True)
    return response

샘플 호출과 출력

다음은 비동기 외부 함수에 대한 샘플 호출과 감성 분석 결과를 포함한 샘플 출력이에요:

create table test_tb(a string);
insert into test_tb values
    ('hello world'),
    ('I am happy');
select ext_func_async(a) from test_tb;

Row | EXT_FUNC_ASYNC(A)
0   | {"Positive": 0.47589144110679626, "Negative": 0.07314028590917587, "Neutral": 0.4493273198604584, "Mixed": 0.0016409909585490823}
1   | {"Positive": 0.9954453706741333, "Negative": 0.00039307220140472054, "Neutral": 0.002452891319990158, "Mixed": 0.0017087293090298772}

샘플 코드에 대한 참고 사항

  • 데이터 처리 함수는 다음 호출로 호출돼요:
lambda_response = lambda_client.invoke(
 ...
 InvocationType='Event',
 ...
)

위와 같이 InvocationType은 'Event'여야 해요. 두 번째 프로세스(또는 스레드)가 비동기여야 하고 Eventinvoke() 메서드를 통해 사용할 수 있는 유일한 비차단 호출 유형이기 때문이에요.

  • 데이터 처리 함수는 HTTP 200 코드를 반환해요. 하지만 이 HTTP 200 코드는 Snowflake에 직접 반환되지 않아요. GET이 상태를 폴링해 데이터 처리 함수가 이 배치를 성공적으로 처리했다는 것을 볼 때까지 Snowflake는 어떤 HTTP 200도 보지 못해요.

더 알아보기 (Learn more)