-
Notifications
You must be signed in to change notification settings - Fork 0
component 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 쓰기 |
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에 반영된다.
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으로 계산된다.
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만 최신 상태에 연결한다.
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 위험으로 재평가한다.
| 항목 | 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#STATE는 LATEST를 복사한 short-term snapshot이고 GraphAggregator5m의 입력이다.
HISTORY#STATE의 목표 TTL은 ADR 0025 기준 2시간이지만, 마지막으로 문서화된 운영 환경값은
48시간이다. 2시간 정책을 실제 적용하려면 data-pipeline에 HISTORY_TTL_HOURS=2를 주입해
재배포하고 기존 item이 자연 만료될 시간을 둬야 한다. GRAPH#5M TTL은 48시간 기준을 유지한다.
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은 제외한다.
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분 집계 |
관련 문서
- 시스템 아키텍처
- 제어 & 데이터 플레인
- Dashboard VPC 설계
- 하드웨어 배치
- Hub EKS 네임스페이스
- Tailscale Mesh VPN
- 데이터 생명주기
- 데이터 조회 모델
- 실시간 갱신 구조
- IoT 데이터 계약
- Reporting Pipeline
- 로컬 스토리지
- 클라우드 스토리지
- Edge Agent
- Edge AI 탐지
- Factory-A Log Adapter
- Dummy Sensor
- Edge IoT Publisher
- Lambda Data Processor
- Risk Normalizer
- Risk Score Engine
- Pipeline Status Aggregator
- Graph Aggregator 5m
- Cloud Infra Collector
- Daily Report Generator
- Risk Alert Dispatcher
- Image Snapshot Pipeline
- Dashboard Backend
- Dashboard Web
- AI 채팅 어시스턴트