모델 추론

모델 추론 (Model Inference)

Flink SQL은 SQL 쿼리에서 모델 추론(model inference)을 수행하기 위한 ML_PREDICT 테이블 값 함수(TVF)를 제공합니다. 이 함수를 사용하면 데이터 스트림에 머신러닝 모델을 SQL에서 직접 적용할 수 있습니다. 모델 생성 방법은 Model Creation을 참고하세요.

출처: 문서

본문

ML_PREDICT 함수

ML_PREDICT 함수는 테이블 입력을 받아 모델을 적용하고, 모델의 예측 결과가 포함된 새 테이블을 반환합니다. 이 함수는 기본 모델이 둘 다 허용하는 경우 동기/비동기 추론 모드를 지원합니다.

구문 (Syntax)

SELECT * FROM
ML_PREDICT(
  TABLE input_table,
  MODEL model_name,
  DESCRIPTOR(feature_columns),
  [CONFIG => MAP['key', 'value']]
)

파라미터 (Parameters)

  • input_table: 처리할 데이터를 포함하는 입력 테이블
  • model_name: 추론에 사용할 모델의 이름
  • feature_columns: 입력 테이블의 어떤 컬럼을 모델의 피처로 사용할지 지정하는 디스크립터
  • config: (선택) 모델 추론을 위한 구성 옵션 맵

구성 옵션 (Configuration Options)

config 맵에 다음 구성 옵션을 지정할 수 있습니다.

키 (Key) 기본값 (Default) 타입 (Type) 설명 (Description)
async (none) Boolean 값은 'true' 또는 'false'일 수 있으며, 플래너가 해당 predict 함수를 선택하도록 힌트를 줍니다. 백엔드 predict 함수 제공자가 제안된 모드를 지원하지 않으면 사용자에게 알리기 위해 예외를 던집니다.
max-concurrent-operations (none) Integer async ml predict가 트리거할 수 있는 async i/o 연산의 최대 개수.
output-mode (none) Enum 비동기 연산의 출력 모드로, {@see AsyncDataStream.OutputMode}로 변환됩니다. 기본값은 ORDERED입니다. ALLOW_UNORDERED로 설정하면 결과의 정확성에 영향을 주지 않을 때 {@see AsyncDataStream.OutputMode.UNORDERED}를 사용하려 시도하고, 그렇지 않으면 ORDERED를 사용합니다. 가능한 값: "ORDERED", "ALLOW_UNORDERED"
timeout (none) Duration 첫 호출부터 비동기 연산의 최종 완료까지의 타임아웃. 여러 번의 재시도를 포함할 수 있으며, 장애 조치(failover) 시 리셋됩니다.

예시 (Example)

-- Basic usage
SELECT * FROM ML_PREDICT(
  TABLE input_table,
  MODEL my_model,
  DESCRIPTOR(feature1, feature2)
);

-- With configuration options
SELECT * FROM ML_PREDICT(
  TABLE input_table,
  MODEL my_model,
  DESCRIPTOR(feature1, feature2),
  MAP['async', 'true', 'timeout', '100s']
);

-- Using named parameters
SELECT * FROM ML_PREDICT(
  INPUT => TABLE input_table,
  MODEL => MODEL my_model,
  ARGS => DESCRIPTOR(feature1, feature2),
  CONFIG => MAP['async', 'true']
);

출력 (Output)

출력 테이블은 입력 테이블의 모든 컬럼에 모델의 예측 컬럼을 더한 것을 포함합니다. 예측 컬럼은 모델의 출력 스키마를 기반으로 추가됩니다.

참고 사항 (Notes)

  • 모델은 ML_PREDICT와 함께 사용하기 전에 카탈로그에 등록되어 있어야 합니다.
  • 디스크립터에 지정된 피처 컬럼의 개수는 모델의 입력 스키마와 일치해야 합니다.
  • 출력의 컬럼 이름이 입력 테이블의 기존 컬럼 이름과 충돌하면 충돌을 피하기 위해 출력 컬럼 이름에 인덱스가 추가됩니다. 예를 들어 출력 컬럼 이름이 prediction이라면 입력 테이블에 같은 이름의 컬럼이 이미 있으면 prediction0으로 이름이 바뀝니다.
  • 비동기 추론의 경우 모델 제공자가 AsyncPredictRuntimeProvider 인터페이스를 지원해야 합니다.
  • ML_PREDICT는 append-only 테이블만 지원합니다. ML_PREDICT 결과는 비결정적이므로 CDC(Change Data Capture) 테이블은 지원되지 않습니다.

모델 제공자 (Model Provider)

ML_PREDICT 함수는 ModelProvider를 사용해 실제 모델 추론을 수행합니다. 제공자는 모델을 등록할 때 지정한 provider 식별자를 기반으로 조회됩니다. 두 가지 유형의 모델 제공자가 있습니다.

  • PredictRuntimeProvider: 동기 모델 추론용
    • 동기 예측 함수를 만드는 createPredictFunction 메서드를 구현합니다.
    • config에서 asyncfalse로 설정된 경우 사용됩니다.
  • AsyncPredictRuntimeProvider: 비동기 모델 추론용
    • 비동기 예측 함수를 만드는 createAsyncPredictFunction 메서드를 구현합니다.
    • config에서 asynctrue로 설정된 경우 사용됩니다.
    • timeout과 buffer capacity에 대한 추가 구성이 필요합니다.

config에 async가 설정되어 있지 않으면 시스템이 동기 또는 비동기 모델 제공자를 선택하며, 둘 다 있으면 비동기 모델 제공자를 선호합니다.

오류 처리 (Error Handling)

이 함수는 다음 경우에 예외를 던집니다.

  • 모델이 카탈로그에 존재하지 않는 경우
  • 피처 컬럼의 개수가 모델의 입력 스키마와 일치하지 않는 경우
  • 모델 파라미터가 누락된 경우
  • 너무 적거나 너무 많은 인자가 제공된 경우

성능 고려 사항 (Performance Considerations)

  • 높은 처리량 시나리오에서는 비동기 추론 모드 사용을 고려하세요.
  • 비동기 추론에는 적절한 timeout과 buffer capacity 값을 구성하세요.
  • 함수의 성능은 기본 모델 제공자 구현에 따라 달라집니다.

관련 문 (Related Statements)

지원되는 모델 제공자 (Supported Model Providers)

Flink는 현재 다음 모델 제공자를 지원합니다.

더 알아보기 (Learn more)