-
Notifications
You must be signed in to change notification settings - Fork 0
component data processor
IoT Core로 수신된 메시지를 정규화·계산해 DynamoDB와 S3에 저장하는 클라우드 처리 Lambda다.
Lambda Data Processor는 Aegis-Pi 데이터 플레인의 클라우드 처리 단계다.
Edge Agent / Dummy Sensor
→ AWS IoT Core (MQTT)
→ IoT Rule (S3 raw 적재 + Lambda 트리거)
→ Lambda data-processor ← 이 컴포넌트
├── factory_state 처리
│ → normalize_factory_state()
│ → Risk Score 계산
│ → pipeline_status 계산
│ → DynamoDB LATEST 갱신 + HISTORY 스냅샷
│ → S3 processed (factory_state / risk_score / state_snapshot)
└── infra_state 처리
→ normalize_infra_state()
→ pipeline_status 계산
→ DynamoDB LATEST 갱신 + HISTORY 스냅샷
→ S3 processed (infra_state / state_snapshot)
초기 설계에서 별도 K8s 서비스로 계획했던 Risk Normalizer, Risk Score Engine, Pipeline Status Aggregator 세 컴포넌트를 Lambda 단일 실행 단위 내부 모듈로 통합했다.
apps/data-processor/
├── lambda_function.py # 진입점, 메시지 분기
└── processor/
├── envelope.py # 메시지 유효성 검증
├── normalizer.py # factory_state / infra_state 정규화
├── risk.py # Risk Score 계산
├── pipeline_status.py # 파이프라인 건강 상태 계산
├── dynamo.py # DynamoDB LATEST / HISTORY 쓰기
└── s3_writer.py # S3 processed 쓰기
source_type 값에 따라 처리 경로가 분기된다.
| source_type | 처리 경로 |
|---|---|
factory_state |
정규화 → Risk Score → pipeline_status → DynamoDB + S3 |
infra_state |
정규화 → pipeline_status → DynamoDB + S3 |
envelope.py는 메시지 필수 필드(schema_version, message_id, factory_id, source_type 등)를 검증한다. 필드가 없거나 허용값이 아니면 EnvelopeError로 처리를 건너뛴다.
테이블명: aegis-factory-status
| 항목 | pk | sk | 내용 |
|---|---|---|---|
| 최신 상태 | FACTORY#{factory_id} |
LATEST |
가장 최근 factory_state + infra_state + risk + pipeline_status |
| 이력 | FACTORY#{factory_id} |
HISTORY#STATE#{updated_at} |
LATEST를 복사한 스냅샷. TTL 48시간 |
LATEST 항목은 factory_state와 infra_state 메시지가 각각 도착할 때마다 해당 필드만 부분 갱신한다. 두 source_type이 독립적으로 쓰여도 단일 항목에 최신 상태가 유지된다.
processor/risk.py가 factory_state payload에서 Risk Score를 계산한다.
MVP에서 활성화된 필드:
| 필드 | 가중치 | 정상→위험 임계값 |
|---|---|---|
temperature |
15% | warning ≥ 32°C, critical ≥ 38°C |
humidity |
10% | warning ≥ 70%, critical ≥ 85% |
ai_event_rate |
10% | fire/fall/bend score 합산 |
Score 범위: 0~100 (높을수록 안전)
| 등급 | Score |
|---|---|
| safe | ≥ 85 |
| warning | 50~84 |
| danger | < 50 |
인프라 기반 지표(node_status, pod_health, data_freshness 등)는 M6에서 추가한다.
processor/pipeline_status.py가 last_infra_state_at 기준 경과 시간으로 파이프라인 건강 상태를 판정한다.
| 상태 | 조건 |
|---|---|
normal |
경과 < 40초 |
warning |
40초 ≤ 경과 < 60초 |
critical |
경과 ≥ 60초 또는 last_infra_state_at 없음 |
infra_state는 20초 주기로 수신되므로 1분 내 이상을 감지할 수 있다.
Lambda가 처리 결과를 저장하는 S3 prefix:
processed/{factory_id}/{dataset}/yyyy={YYYY}/mm={MM}/dd={DD}/hh={HH}/{message_id}.json
| dataset | 생산 시점 | 내용 |
|---|---|---|
factory_state |
factory_state 수신 | 정규화된 센서·AI 값 |
risk_score |
factory_state 수신 | Risk Score 계산 결과 |
infra_state |
infra_state 수신 | 정규화된 노드·워크로드 상태 |
state_snapshot |
모든 메시지 수신 | DynamoDB LATEST 전체 복사본 |
infra/data-pipeline/
Terraform infra/data-pipeline으로 IoT Rule 3개(factory-a/b/c)와 Lambda를 함께 배포한다. Lambda 실행 역할은 DynamoDB 쓰기, S3 쓰기, CloudWatch 로그 권한을 갖는다.
DynamoDB 테이블(aegis-factory-status)은 infra/foundation에서 영구 리소스로 관리한다. Hub EKS를 내려도 테이블과 데이터는 유지된다.
관련 문서
- 시스템 아키텍처
- 제어 & 데이터 플레인
- 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 채팅 어시스턴트