기술 소개

데이터브릭스 머신러닝을 활용한 실시간에 가까운 이상 탐지 (Anomaly Detection)

스누피18 2023. 11. 20. 18:18

(원문: 링크)

 

이상 탐지가 중요한 이유는 무엇인가요?

리테일, 금융, 사이버 보안, 기타 모든 산업에서 비정상적인 행동이 발생하는 즉시 이를 발견하는 것은 매우 중요한 과제입니다. 이를 탐지할 수 있는 역량이 부족하면 매출 손실, 규제 기관으로부터의 벌금, 보안 침해로 인한 고객 개인정보 및 신뢰 상실로 이어질 수 있습니다. 따라서 비정상적인 신용카드 거래를 발견하거나, 의심스러운 행동을 하는 사용자를 발견하거나, 웹 서비스에 대한 요청에서 이상한 패턴을 식별하는 것은 천국과 지옥의 차이만큼이나 큰 차이를 만들어낼 수 있습니다. 

 

이상 탐지의 도전 과제

이상 탐지에는 몇 가지 도전 과제가 있습니다. 첫 번째는 '이상 징후'가 무엇인지에 대한 데이터 사이언스적 질문입니다. 다행히도 머신러닝에는 데이터에서 정상 패턴과 비정상 패턴을 구별하는 방법을 학습할 수 있는 강력한 도구가 있습니다. 이상 탐지의 경우, 모든 이상 징후가 어떻게 생겼는지 알 수 없기 때문에 머신러닝 모델 학습을 위한 데이터 세트에 라벨을 붙이는 것이 불가능합니다. 따라서 라벨이 지정되지 않은 데이터에서 패턴을 학습하는 비지도 학습을 사용하여 이상 징후를 탐지해야 합니다.

이상 탐지를 위한 완벽한 비지도 (unsupervised) 머신러닝 모델을 찾았다고 해도 실제 문제는 이제 시작에 불과합니다. 원본 시스템에서 데이터가 도착하자마자 각 관측이 수집되고, 변환되고, 최종적으로 모델로 점수를 매길 수 있도록 이 모델을 프로덕션에 적용해야 합니다. 그것도 거의 실시간, 또는 5분~10분 정도의 짧은 간격으로. 이를 위해서는 정교한 extract, load, transform (ELT) 파이프라인을 구축하고 이를 비정상적인 레코드를 정확하게 식별할 수 있는 비지도 머신러닝 모델과 통합해야 합니다. 또한 이 end-to-end 파이프라인은 수집부터 모델 추론에 이르기까지 데이터 품질을 보장하면서 항상 실행되는 프로덕션급이어야 하며, 기본 인프라를 유지 관리해야 합니다.

데이터브릭스 레이크하우스 플랫폼으로 과제 해결하기

데이터브릭스를 사용하면 이 과정이 복잡하지 않습니다. 실시간에 가까운 이상 탐지 파이프라인을 100% SQL로 구축할 수 있으며, 머신러닝 모델을 학습하는 데 Python만 사용하는 것도 가능합니다. 데이터 수집, 변환, 모델 추론은 모두 SQL로 수행할 수 있습니다.

이 블로그에서는 특히 비정상 레코드를 탐지하는 데 적합한 격리 포리스트 알고리즘 (isolation forest algorithm) 을 훈련하고, 훈련된 모델을 델타 라이브 테이블 (Delta Live Table, DLT)을 사용하여 생성된 스트리밍 데이터 파이프라인에 통합하는 방법에 대해 설명합니다. DLT는 데이터 엔지니어링 프로세스를 자동화하는 ETL 프레임워크입니다. DLT는 간단한 선언적 접근 방식을 사용하여 신뢰할 수 있는 데이터 파이프라인을 생성하고 배치 및 스트리밍 데이터를 위한 대규모 기본 인프라를 완벽하게 관리합니다. 그 결과 거의 실시간에 가까운 이상 탐지 시스템을 구축할 수 있습니다. 특히, 이 블로그에 사용된 데이터는 Kaggle에서 신용카드 거래를 시뮬레이션하기 위해 생성된 합성 데이터의 샘플이며, 이렇게 탐지된 이상 거래는 사기 거래로 보시면 됩니다. 

블로그에 소개된 ML 및 델타 라이브 테이블 기반 이상 탐지 솔루션의 아키텍처 개요


Scikit-learn 격리 포리스트 알고리즘 구현은 기본적으로 데이터브릭스 머신러닝 런타임에서 사용할 수 있으며, MLflow 프레임워크를 사용하여 학습되는 동안 이상 탐지 모델을 추적하고 로깅합니다. ETL 파이프라인은 전적으로 델타 라이브 테이블을 사용하여 SQL로 개발됩니다.

레이블이 지정되지 않은 데이터에 대한 이상 탐지를 위한 격리 포리스트

격리 포리스트는 랜덤 포리스트와 유사한 트리 기반 앙상블 알고리즘의 한 유형입니다. 이 알고리즘은 주어진 관측값 집합의 아웃라이어(outlier)보다 인라이어(Inlier)를 분리하기가 더 어렵다고 가정하도록 설계되었습니다. 높은 수준에서 보면, 비정상적이지 않은 지점(예: 일반적인 신용카드 거래)은 분리하기 어렵기 때문에 의사 결정 트리에서 더 깊숙이 위치할 것이고, 비정상적인 지점은 그 반대의 경우입니다. 이 알고리즘은 레이블이 없는 관찰 집합에 대해 학습한 다음 이전에 볼 수 없었던 데이터에서 비정상적인 기록을 예측하는 데 사용할 수 있습니다.

Outlier를 격리하는 것이 Inlier를 격리하는 것보다 쉽습니다.


데이터브릭스가 모델 훈련 및 추적에 어떤 도움을 주나요?

데이터브릭스에서 머신러닝과 관련된 작업을 할 때는 머신러닝 런타임과 함께 클러스터를 사용하는 것이 필수입니다. 데이터 사이언스 및 머신러닝 관련 작업에 일반적으로 사용되는 많은 오픈 소스 라이브러리는 ML 런타임에서 기본적으로 사용할 수 있습니다. Scikit-learn은 이러한 라이브러리 중 하나이며, 격리 포리스트 알고리즘을 훌륭하게 구현한 라이브러리입니다.

모델을 정의하는 방법은 다음과 같습니다.

from sklearn.ensemble import IsolationForest
isolation_forest = IsolationForest(n_jobs=-1, warm_start=True, random_state=42)

 

이 런타임은 무엇보다도 머신러닝 실험 추적, 모델 스테이징 및 배포를 위해 노트북 환경과 MLflow를 긴밀하게 통합할 수 있게 해줍니다.

 

ML 클러스터에 연결된 노트북 환경에서 수행되는 모든 모델 훈련 또는 하이퍼파라미터 최적화는 기본적으로 활성화된 기능인 MLflow 자동 로깅을 통해 자동으로 기록됩니다.

모델이 로깅되면 다양한 방법으로 MLflow 내에서 모델을 등록하고 배포할 수 있습니다. 특히, 이 모델을 Apache Spark™를 통한 분산 인스트림 또는 배치 추론을 위해 벡터화된 사용자 정의 함수(User Defined Function, UDF)로 배포하기 위해  MLflow는 사용자 인터페이스(UI) 자체 내에서 UDF를 생성하고 등록하기 위한 코드를 생성합니다 (아래 이미지 참고).

MLflow는 모델 추론을 위해 Apache Spark UDF를 생성하고 등록하기 위한 코드를 생성합니다.


이 외에도 MLflow REST API를 사용하면 다음과 같이 몇 줄의 코드를 함수에 깔끔하게 담아 운영 중인 기존 모델을 아카이브하고 새로 학습된 모델을 운영 환경에 투입할 수 있습니다.

def train_model(mlFlowClient, loaded_model, model_name, run_name)->str:
  """
  Trains, logs, registers and promotes the model to production. Returns the URI of the model in prod
  """
  with mlflow.start_run(run_name=run_name) as run:

    # 0. Fit the model 
    loaded_model.fit(X_train)

    # 1. Get predictions 
    y_train_predict = loaded_model.predict(X_train)

    # 2. Create model signature 
    signature = infer_signature(X_train, y_train_predict)
    runID = run.info.run_id

    # 3. Log the model alongside the model signature 
    mlflow.sklearn.log_model(loaded_model, model_name, signature=signature, registered_model_name= model_name)

    # 4. Get the latest version of the model 
    model_version = mlFlowClient.get_latest_versions(model_name,stages=['None'])[0].version

    # 5. Transition the latest version of the model to production and archive the existing versions
    client.transition_model_version_stage(name= model_name, version = model_version, stage='Production', archive_existing_versions= True)


    return mlFlowClient.get_latest_versions(model_name, stages=["Production"])[0].source

 

프로덕션 시나리오에서는 단일 레코드가 모델에 의해 한 번만 스코어되기를 원할 것입니다. 데이터브릭스에서는 Auto Loader를 사용하여 이 "정확히 한 번" 동작을 보장할 수 있습니다. Auto Loader는 Python 또는 SQL을 사용하는 델타 라이브 테이블, 구조화된 스트리밍 애플리케이션과 함께 작동합니다.

고려해야 할 또 다른 중요한 요소는 환경적이든 행동적이든 비정상적인 발생의 특성이 시간에 따라 변한다는 것입니다. 따라서 새로운 데이터가 도착하면 모델을 재학습시켜야 합니다.

모델 학습 로직이 포함된 노트북을 데이터브릭스 워크플로우에서 예약된 작업으로 생성하면, 작업이 실행될 때마다 최신 모델을 효과적으로 재학습하고 생산에 투입할 수 있습니다.

델타 라이브 테이블로 실시간에 가까운 이상 탐지 달성

머신러닝 측면은 이 과제의 일부에 불과합니다. 데이터 수집, 변환, 모델 추론을 결합한 실시간에 가까운 프로덕션급 데이터 파이프라인을 구축하는 것이 사실 더 어려운 과제입니다. 복잡하고 시간이 많이 소요되며 오류가 발생하기 쉽기 때문이죠. 

이를 상시 가동할 수 있는 인프라를 구축하고 유지 관리하며 오류를 처리하려면 데이터 엔지니어링보다 더 많은 소프트웨어 엔지니어링 노하우가 필요합니다. 또한 전체 파이프라인을 통해 데이터 품질을 보장해야 합니다. 특정 애플리케이션에 따라 복잡성이 추가될 수도 있습니다. 

 

바로 이 부분에서 델타 라이브 테이블(DLT)이 등장합니다.

DLT 용어로 노트북 라이브러리는 기본적으로 DLT 파이프라인에 대한 코드의 일부 또는 전부가 포함된 노트북입니다. DLT 파이프라인에는 하나 이상의 노트북이 연결될 수 있으며, 각 노트북은 SQL 또는 Python 구문을 사용할 수 있습니다. 첫 번째 노트북 라이브러리에는 파이썬으로 구현된 로직이 포함되어 MLflow 모델 레지스트리에서 모델을 가져오고, 수집된 레코드가 파이프라인의 다운스트림에서 기능화되면 모델 추론 함수를 사용할 수 있도록 UDF를 등록합니다. 🍯꿀팁: DLT Python 노트북에서 새 패키지는 첫 번째 셀에 %pip magic 명령으로 설치해야 합니다.

두 번째 DLT 라이브러리 노트북은 Python 또는 SQL 구문으로 구성할 수 있습니다. DLT의 다재다능함을 증명하기 위해 SQL을 사용하여 데이터 수집, 변환 및 모델 추론을 수행했습니다. 이 노트북에는 파이프라인을 구성하는 실제 데이터 변환 로직이 포함되어 있습니다.

수집은 객체 스토리지로 스트리밍된 데이터를 점진적으로 로드할 수 있는 Auto Loader로 수행됩니다. 이것은 메달리온 아키텍처의 브론즈(raw data) 테이블로 읽혀집니다. 또한 아래 구문에서 스트리밍 라이브 테이블은 오브젝트 스토리지에서 데이터가 지속적으로 수집되는 곳입니다. Auto Loader는 데이터가 수집될 때 스키마를 감지하도록 구성됩니다. Auto Loader는 진화하는 스키마도 처리할 수 있으며, 이는 많은 실제 이상 탐지 시나리오에 적용됩니다.

CREATE OR REFRESH STREAMING LIVE TABLE transaction_readings_raw
COMMENT "The raw transaction readings, ingested from landing directory"
TBLPROPERTIES ("quality" = "bronze")
AS SELECT * FROM cloud_files("/FileStore/tables/transaction_landing_dir", "json", map("cloudFiles.inferColumnTypes", "true"))


또한 DLT를 사용하면 데이터 품질 제약 조건을 정의할 수 있으며 개발자나 애널리스트가 오류를 수정할 수 있는 기능을 제공합니다. 주어진 레코드가 주어진 제약 조건을 충족하지 않는 경우, DLT는 해당 레코드를 유지하거나 삭제하거나 파이프라인을 완전히 중단할 수 있습니다. 아래 예시에서는 트랜잭션 시간이나 금액이 지정되지 않은 경우 레코드를 삭제하는 변환 단계 중 하나에 제약 조건이 정의되어 있습니다.

 

CREATE OR REFRESH STREAMING LIVE TABLE transaction_readings_cleaned(
  CONSTRAINT valid_transaction_reading EXPECT (AMOUNT IS NOT NULL AND TIME IS NOT NULL) ON VIOLATION DROP ROW
)
TBLPROPERTIES ("quality" = "silver")

COMMENT "Drop all rows with nulls for Time and store these records in a silver delta table"
AS SELECT * FROM STREAM(live.transaction_readings_raw)



델타 라이브 테이블은 사용자 정의 함수(User Defined Functions, UDF)도 지원합니다. UDF는 SQL을 사용하여 스트리밍 DLT 파이프라인에서 모델 추론을 활성화하는 데 사용할 수 있습니다. 아래 예시에서는 훈련된 격리 포레스트 모델을 캡슐화하는 이전에 등록된 Apache Spark™ 벡터화된 UDF를 사용하고 있습니다.

CREATE OR REFRESH STREAMING LIVE TABLE predictions
COMMENT "Use the isolation forest vectorized udf registered in the previous step to predict anomalous transaction readings"
TBLPROPERTIES ("quality" = "gold")
AS SELECT cust_id, detect_anomaly(<enter by="" column="" commas="" names="" separated="">) as 
anomalous from STREAM(live.transaction_readings_cleaned)

 

데이터 사이언티스트가 Python으로 훈련한 머신 러닝 모델(예: scikit-learn, xgboost 또는 기타 머신 러닝 라이브러리)을 사용하여 전체 SQL 데이터 파이프라인에서 추론할 수 있으므로 SQL을 선호하는 SQL 분석가 및 데이터 엔지니어에게는 매우 재밌는 기능이죠. 

이 노트북은 DLT 파이프라인을 생성하는 데 사용됩니다 (아래의 구성 세부 정보 섹션에 자세히 설명되어 있습니다). 리소스, 테이블을 설정하고 종속성을 파악하는 짧은 시간 (그리고 최종 사용자로부터 DLT가 추상화하는 다른 모든 복잡한 작업)이 지나면, 학습된 머신러닝 모델을 통해 데이터가 지속적으로 처리되고 거의 실시간으로 비정상적인 레코드가 감지되는 DLT 파이프라인이 UI에 렌더링됩니다.

DLT 사용자 인터페이스에서 본 엔드 투 엔드 델타 라이브 테이블 파이프라인


이 파이프라인이 실행되는 동안, 데이터브릭스 SQL을 사용하여 식별된 비정상 레코드를 시각화할 수 있으며, 데이터브릭스 SQL 대시보드 새로고침 기능을 통해 지속적으로 업데이트할 수 있습니다. '예측' 테이블에 대해 실행된 쿼리를 기반으로 시각화된 대시보드는 아래에서 확인할 수 있습니다.

예측된 비정상 레코드를 대화형으로 표시하도록 구축된 Databricks SQL 대시보드

정리하자면 이 블로그에서는 이상 탐지을 위한 격리 포리스트 알고리즘을 훈련하는 데 사용되는 데이터브릭스 머신러닝 및 워크플로우에서 사용할 수 있는 기능과 거의 실시간으로 이 기능을 수행할 수 있는 Delta 라이브 테이블 파이프라인을 정의하는 프로세스에 대해 자세히 설명합니다. 델타 라이브 테이블은 최종 사용자로부터 프로세스의 복잡성을 추상화하여 자동화합니다.

이 블로그에서는 델타 라이브 테이블의 전체 기능 중 일부만 소개했습니다. 데이터브릭스의 주요 기능에 대한 이해하기 쉬운 설명서는 https://docs.databricks.com/data-engineering/delta-live-tables/index.html 에서 확인할 수 있습니다.

 

모범 사례

데이터브릭스 워크플로우 사용자 인터페이스를 사용하여 델타 라이브 테이블 파이프라인을 생성할 수 있습니다.


실시간에 가까운 방식으로 이상 탐지를 수행하려면 DLT 파이프라인을 연속 모드에서 실행해야 합니다. 공식 빠른 시작에 설명된 프로세스에 따라 이 블로그의 리포지토리에서 사용할 수 있는 앞서 설명한 Python 및 SQL 노트북을 사용하여 생성할 수 있습니다. 다른 구성은 원하는 대로 입력할 수 있습니다.

간헐적인 파이프라인 실행이 허용되는 사용 사례(예: 소스 시스템에서 일괄 수집한 레코드에 대한 이상 탐지)의 경우, 파이프라인을 트리거 모드에서 10분 정도의 짧은 간격으로 실행할 수 있습니다. 그런 다음 이 트리거된 파이프라인이 실행되도록 일정을 지정할 수 있으며, 각 실행에서 데이터는 증분 방식으로 파이프라인을 통해 처리됩니다.

그 후, 클러스터 자동 확장이 활성화된 파이프라인 구성을 저장하고 파이프라인을 시작할 수 있습니다(처리 병목 현상 없이 파이프라인을 통해 전달되는 다양한 레코드 부하를 처리하기 위해). 또는 이러한 모든 구성을 JSON 형식으로 깔끔하게 설명하여 동일한 입력 양식에 입력할 수 있습니다.

델타 라이브 테이블은 클러스터 구성, 기본 테이블 최적화 및 최종 사용자를 위한 기타 여러 가지 중요한 세부 정보를 파악합니다. 파이프라인을 실행하기 위해 반복 개발에 도움이 되는 개발 모드 또는 프로덕션에 맞춰진 프로덕션 모드를 선택할 수 있습니다. 후자의 경우, DLT는 자동으로 재시도 및 클러스터 재시작을 수행합니다.

위에서 설명한 모든 작업은 델타 라이브 테이블 REST API를 통해 수행할 수 있다는 점을 강조하는 것이 중요합니다. 이는 이 블로그의 앞부분에서 언급한 것처럼 예약된 작업을 통해 격리 포리스트가 재훈련될 때마다 다운타임 없이 연속 모드에서 실행되는 DLT 파이프라인을 즉시 편집할 수 있는 프로덕션 시나리오에 특히 유용합니다.

이 예제에서 델타 라이브 테이블 파이프라인에 대한 구성. 생성된 델타 테이블을 저장할 대상 데이터베이스 이름을 입력합니다.


데이터브릭스로 직접 구축하기

이 솔루션을 다시 만들기 위한 노트북과 단계별 지침은 모두 다음 리포지토리 (https://github.com/sathishgang-db/anomaly_detection_using_databricks)에 포함되어 있습니다.

모델 학습 작업에는 반드시 데이터브릭스 머신러닝 런타임이 포함된 클러스터를 사용해야 합니다. 여기에 제공된 예는 다소 단순하지만, 더 복잡한 변환에도 동일한 원칙이 적용되며, 델타 라이브 테이블은 이러한 파이프라인 구축에 내재된 복잡성을 줄이기 위해 만들어졌습니다. 이 블로그의 아이디어를 여러분의 사용 사례에 맞게 조정해 보시기 바랍니다.

이 외에도, DLT 기능에 대한 훌륭한 데모와 워크스루는 다음 영상에서 확인할 수 있습니다: https://www.youtube.com/watch?v=BIxwoO65ylY&t=1s

데이터브릭스에 대한 포괄적인 엔드투엔드 머신러닝 워크플로는 여기에서 확인할 수 있습니다:
https://www.youtube.com/watch?v=5CpaimNhMzs