Skip to content

component data processor

JJong-03 edited this page Jun 17, 2026 · 4 revisions

컴포넌트 — Lambda Data Processor

IoT Core가 전달한 canonical JSON을 정규화하고, Safety Score와 pipeline freshness를 계산해 DynamoDB/S3 read model을 갱신하는 Lambda다.


역할

Data Processor는 현재 데이터 플레인의 핵심 처리 단위다. 초기 설계에서 독립 컴포넌트로 나뉘었던 Risk Normalizer, Risk Score Engine, Pipeline Status Aggregator는 현재 모두 이 Lambda 내부 Python 모듈로 구현되어 있다.

Edge publisher / Dummy generator
  -> AWS IoT Core
  -> IoT Rule
      -> S3 raw
      -> Lambda data-processor
          -> envelope validation
          -> source_type별 normalize
          -> pipeline_status 계산
          -> Safety Score 계산
          -> DynamoDB LATEST 부분 갱신
          -> DynamoDB HISTORY#STATE snapshot 저장
          -> S3 processed 저장

구현 위치

apps/data-processor/
├── lambda_function.py
└── processor/
    ├── envelope.py
    ├── normalizer.py
    ├── risk.py
    ├── pipeline_status.py
    ├── dynamo.py
    └── s3_writer.py
모듈 역할
lambda_function.py Lambda handler, source_type 분기, refresh action 처리
envelope.py canonical JSON 필수 필드 검증
normalizer.py factory_state, infra_state, image_snapshot 정규화
pipeline_status.py last_infra_state_at 기준 freshness 판정
risk.py Safety Score 계산
dynamo.py LATEST, HISTORY#STATE 쓰기
s3_writer.py S3 processed/ object 쓰기

메시지 처리 흐름

factory_state

parse envelope
  -> normalize_factory_state(payload)
  -> DynamoDB LATEST 조회
  -> last_infra_state_at 기준 pipeline_status 계산
  -> 최신 infra_state + 현재 factory_state + pipeline_status로 Safety Score 계산
  -> LATEST.factory_state / LATEST.risk / LATEST.pipeline_status 갱신
  -> HISTORY#STATE#{updated_at} 저장
  -> S3 processed factory_state / risk_score / state_snapshot 저장

factory_state는 온도, 습도, 기압, AI score, AI sample count를 제공한다. 최신 infra_state가 이미 있으면 인프라 상태도 함께 score에 반영된다.

infra_state

parse envelope
  -> normalize_infra_state(payload)
  -> source_timestamp 기준 pipeline_status 계산
  -> DynamoDB LATEST 조회
  -> 최신 factory_state가 있으면 Safety Score 재계산
  -> LATEST.infra_state / LATEST.pipeline_status / 선택적 LATEST.risk 갱신
  -> HISTORY#STATE#{updated_at} 저장
  -> S3 processed infra_state / state_snapshot 저장

infra_state가 방금 도착한 경우 freshness age는 거의 0으로 계산된다.

image_snapshot

parse envelope
  -> normalize_image_snapshot(payload)
  -> LATEST.latest_image_snapshot 갱신
  -> HISTORY#STATE#{updated_at} 저장
  -> S3 processed image_snapshot / state_snapshot 저장

이미지 binary는 DynamoDB에 넣지 않고 S3 bucket/key/content metadata만 최신 상태에 연결한다.


Freshness Refresh

Data Processor는 일반 IoT 메시지 외에 아래 action도 처리한다.

{"action": "refresh_pipeline_status", "factories": ["factory-a", "factory-b", "factory-c"]}

Terraform은 EventBridge Scheduler rate(1 minute)로 이 action을 호출한다. 새 메시지가 없어도 현재 시각 기준으로 last_infra_state_at의 age가 증가하기 때문에, refresh가 없으면 Dashboard에 과거의 safe 점수가 계속 남는다.

Refresh는 각 factory의 LATEST를 읽고 다음만 갱신한다.

pipeline_status
risk
updated_at
HISTORY#STATE snapshot
S3 processed/{factory_id}/state_snapshot/

마지막 factory_state/infra_state 원본 값을 새로 만든다는 뜻은 아니다. 기존 최신 데이터를 유지하되, 데이터 단절 자체를 data_freshness 위험으로 재평가한다.


DynamoDB 저장 계약

항목 pk sk 보존
최신 상태 FACTORY#{factory_id} LATEST overwrite
상태 이력 FACTORY#{factory_id} HISTORY#STATE#{updated_at} HISTORY_TTL_HOURS 기준 TTL

LATEST는 공장별 현재 상태 1건이며, Dashboard 현재 카드의 1차 read model이다. HISTORY#STATELATEST를 복사한 short-term snapshot이고 GraphAggregator5m의 입력이다.

HISTORY#STATE의 목표 TTL은 ADR 0025 기준 2시간이지만, 마지막으로 문서화된 운영 환경값은 48시간이다. 2시간 정책을 실제 적용하려면 data-pipeline에 HISTORY_TTL_HOURS=2를 주입해 재배포하고 기존 item이 자연 만료될 시간을 둬야 한다. GRAPH#5M TTL은 48시간 기준을 유지한다.


S3 Processed 출력

processed/{factory_id}/factory_state/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json
processed/{factory_id}/risk_score/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json
processed/{factory_id}/infra_state/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json
processed/{factory_id}/image_snapshot/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json
processed/{factory_id}/state_snapshot/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{updated_at}.json

S3 state_snapshot은 DynamoDB HISTORY#STATE와 같은 상태 snapshot이지만 DynamoDB TTL 정책 필드인 ttl은 제외한다.


GraphAggregator5m과의 구분

Data Processor는 개별 메시지와 freshness refresh를 처리한다. 5분 그래프 집계는 이 Lambda가 아니라 별도 GraphAggregator5m Lambda가 담당한다.

구분 Data Processor GraphAggregator5m
실행 주기 IoT Rule trigger + 1분 refresh 5분 scheduler
입력 IoT 메시지 또는 LATEST HISTORY#STATE window query
출력 LATEST, HISTORY#STATE, S3 processed/ GRAPH#5M, S3 processed_agg/
목적 현재 상태와 이벤트별 처리 결과 Dashboard 그래프용 5분 집계

관련 문서

Aegis-Pi Wiki

· 대표 문서 목록은 홈의 문서 탐색 표 참조

시작하기

요구사항

핵심 개념

아키텍처

컴포넌트 (Edge → Cloud → Dashboard)

Dashboard & 운영

시나리오 · 사례 · 참조

Clone this wiki locally