-
Notifications
You must be signed in to change notification settings - Fork 0
architecture data lifecycle
Edge 로컬 저장소에서 AWS IoT Core, S3, DynamoDB, 보고서, 이미지 snapshot까지 이어지는 현재 구현 기준의 데이터 흐름이다.
기준 코드:
docs/specs/data_storage_pipeline.mddocs/ops/23_data_pipeline.mddocs/ops/24_daily_factory_report.mddocs/ops/26_dynamodb_key_model.mddocs/ops/32_image_snapshot_pipeline.mdapps/*/s3_writer.pyinfra/foundation/{s3,dynamodb}.tfinfra/data-pipeline/variables.tf
Edge workloads
- factory-a-log-adapter
- dummy-data-generator
- snapshot-uploader
|
v
Longhorn / outbox
/var/lib/aegis/outbox
|
v
edge-iot-publisher
|
v
AWS IoT Core topic: aegis/{factory_id}/{source_type}
|
+--> IoT Rule
| S3 raw/{factory_id}/{source_type}/yyyy=YYYY/mm=MM/dd=DD/{message_id}.json
|
+--> AEGIS-Lambda-DataProcessor
DynamoDB FACTORY#{factory_id}/LATEST
DynamoDB FACTORY#{factory_id}/HISTORY#STATE#{updated_at}
S3 processed/{factory_id}/{dataset}/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{message_id}.json
S3 processed/{factory_id}/state_snapshot/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{updated_at}.json
EventBridge rate(1 minute)
-> DataProcessor refresh_pipeline_status
-> LATEST refresh, HISTORY#STATE, state_snapshot
EventBridge rate(5 minutes)
-> GraphAggregator5m
-> DynamoDB FACTORY#{factory_id}/GRAPH#5M#{bucket_start}
-> S3 processed_agg/{factory_id}/metrics_5m/yyyy=YYYY/mm=MM/dd=DD/hh=HH/mm=MM.json
EventBridge rate(1 minute)
-> CloudInfraFastCollector
-> DynamoDB CLOUD#infra/LATEST.fast
-> DynamoDB CLOUD#infra/HISTORY#FAST#{updated_at}
-> S3 processed/cloud_infra/fast/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{updated_at}.json
EventBridge rate(5 minutes)
-> CloudInfraSlowCollector
-> DynamoDB CLOUD#infra/LATEST.slow
-> DynamoDB CLOUD#infra/HISTORY#SLOW#{updated_at}
-> S3 processed/cloud_infra/slow/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{updated_at}.json
S3 ObjectCreated on selected processed snapshots
-> RiskAlertDispatcher
-> DynamoDB ALERT#{scope}/{severity}#{reason}#{status}
-> Slack webhook
Daily report pipeline
-> reads S3 processed and processed/cloud_infra
-> writes S3 reports/daily/yyyy=YYYY/mm=MM/dd=DD/{factory_id_or_cloud-infra}/...
Dashboard API/Web
-> reads DynamoDB LATEST, GRAPH#5M, short HISTORY#STATE, CLOUD#infra/LATEST
-> reads S3 reports and selected processed objects for detail/audit
| 저장 계층 | 저장 위치 / 키 | 목적 | 보존 기간 | 생성 주체 | 조회 주체 |
|---|---|---|---|---|---|
| Edge InfluxDB | safe_edge_db |
factory-a 센서, AI score, 오디오 감지 결과의 로컬 단기 원천 | InfluxDB retention 1d
|
sensor/AI/audio workloads |
factory-a-log-adapter, 운영자 |
| Longhorn/outbox |
/var/lib/aegis/outbox PVC |
canonical JSON publish 버퍼, atomic write, 재시도 대기 | publish 성공 시 파일 삭제. 실패/스키마 오류는 outbox/quarantine에 남음 |
factory-a-log-adapter, dummy-data-generator, snapshot-uploader
|
edge-iot-publisher, 운영자 |
| snapshot hostPath | /var/lib/safe-edge/snapshots |
AI event 이미지의 노드 로컬 임시 저장 | 24시간 초과 이미지 정리, 매일 03:00 KST purge | safe-edge-integrated-ai |
snapshot-uploader, 운영자 |
| S3 image_snapshot binary | image_snapshot/factory_id={factory_id}/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{filename} |
이미지 원본 바이너리 저장. IoT Core에는 bytes를 보내지 않음 | S3 bucket lifecycle에 별도 만료 rule 없음 |
snapshot-uploader가 presigned PUT 사용 |
Dashboard/API, 운영자, 보고서 확장 |
| S3 raw | raw/{factory_id}/{source_type}/yyyy=YYYY/mm=MM/dd=DD/{message_id}.json |
IoT payload 원본 보존, 감사, 재처리 근거 | 90일 후 Glacier Instant Retrieval 전환. 삭제 TTL 아님 | IoT Rule | 운영자, 재처리, cloud infra freshness |
| S3 processed | processed/{factory_id}/{factory_state,risk_score,infra_state,image_snapshot}/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{message_id}.json |
DataProcessor 정규화 결과, risk 결과, image metadata | 365일 후 Standard-IA 전환. 삭제 TTL 아님 | DataProcessor | Daily report, Dashboard detail/audit, 운영자 |
| S3 state snapshot | processed/{factory_id}/state_snapshot/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{updated_at}.json |
DynamoDB HISTORY#STATE와 같은 상태 snapshot의 장기 사본. ttl 필드는 제외 |
365일 후 Standard-IA 전환. 삭제 TTL 아님 | DataProcessor, DataProcessorRefresh1m | RiskAlertDispatcher, Daily report, 운영자 |
| DynamoDB LATEST |
pk=FACTORY#{factory_id}, sk=LATEST
|
factory별 현재 상태 read model | TTL 없음 | DataProcessor, DataProcessorRefresh1m | Dashboard API/Web, GraphAggregator, CloudInfraFastCollector |
| DynamoDB HISTORY#STATE |
pk=FACTORY#{factory_id}, sk=HISTORY#STATE#{updated_at}
|
단기 raw-resolution 상태 이력, 1h Dashboard chart, GraphAggregator 입력 |
dynamodb_history_ttl_hours=2 시간 |
DataProcessor, DataProcessorRefresh1m | Dashboard 1h history, GraphAggregator |
| DynamoDB GRAPH#5M |
pk=FACTORY#{factory_id}, sk=GRAPH#5M#{bucket_start}
|
5분 집계 read model. 6h/12h/24h Dashboard chart |
graph_bucket_ttl_hours=48 시간 |
GraphAggregator5m | Dashboard API/Web |
| S3 processed_agg | processed_agg/{factory_id}/metrics_5m/yyyy=YYYY/mm=MM/dd=DD/hh=HH/mm=MM.json |
GRAPH#5M의 장기 보조 산출물, audit/replay | 명시 lifecycle rule 없음. processed/와 같은 장기 분석 산출물로 취급 |
GraphAggregator5m | Cloud infra freshness, 운영자 |
| DynamoDB CLOUD#infra LATEST |
pk=CLOUD#infra, sk=LATEST
|
cloud infra 현재 상태 read model. fast와 slow 필드 부분 갱신 |
TTL 없음 | CloudInfraFastCollector, CloudInfraSlowCollector | Cloud infra Dashboard/API |
| DynamoDB CLOUD#infra history |
HISTORY#FAST#{updated_at}, HISTORY#SLOW#{updated_at}
|
cloud infra 최근 fast/slow snapshot | fast 6시간, slow 24시간 | CloudInfraFastCollector, CloudInfraSlowCollector | Cloud infra Dashboard/API, 운영자 |
| S3 cloud infra snapshots | processed/cloud_infra/{fast,slow}/yyyy=YYYY/mm=MM/dd=DD/hh=HH/{updated_at}.json |
cloud infra snapshot 장기 사본. ttl 필드는 제외 |
365일 후 Standard-IA 전환. 삭제 TTL 아님 | CloudInfra collectors | RiskAlertDispatcher, cloud-infra daily report |
| DynamoDB ALERT# state |
pk=ALERT#{scope}, sk={severity}#{reason}#{status} 또는 OBSERVATION#{severity}#{reason}#{status}
|
Slack alert cooldown, dedupe, 연속 관측 확인 상태 |
risk_alert_state_ttl_seconds=604800초, 즉 7일 |
RiskAlertDispatcher | RiskAlertDispatcher |
| S3 reports | reports/daily/yyyy=YYYY/mm=MM/dd=DD/{factory_id_or_cloud-infra}/... |
Bedrock 기반 일일 운영 보고서와 중간 집계 | 명시 lifecycle rule 없음 | Daily report Step Functions/Lambda | Dashboard Reports API, 운영자 |
| source_type | Edge 생성 | raw S3 | processed S3 | DynamoDB 반영 |
|---|---|---|---|---|
factory_state |
factory-a-log-adapter 또는 dummy generator |
raw/{factory_id}/factory_state/yyyy=YYYY/mm=MM/dd=DD/{message_id}.json |
processed/{factory_id}/factory_state/.../hh=HH/{message_id}.json, risk_score, state_snapshot
|
LATEST.factory_state, risk, pipeline_status, HISTORY#STATE
|
infra_state |
factory-a-log-adapter 또는 dummy generator |
raw/{factory_id}/infra_state/yyyy=YYYY/mm=MM/dd=DD/{message_id}.json |
processed/{factory_id}/infra_state/.../hh=HH/{message_id}.json, state_snapshot
|
LATEST.infra_state, pipeline_status, HISTORY#STATE
|
image_snapshot |
snapshot-uploader metadata |
raw/{factory_id}/image_snapshot/yyyy=YYYY/mm=MM/dd=DD/{message_id}.json |
processed/{factory_id}/image_snapshot/.../hh=HH/{message_id}.json |
LATEST.latest_image_snapshot, last_image_snapshot_at
|
raw/는 source type과 무관하게 hh partition이 없다. processed/, processed_agg/, processed/cloud_infra/, image_snapshot/ binary는 hh partition을 가진다.
DynamoDB TTL과 S3 lifecycle은 다른 정책이다.
| 대상 | 정책 |
|---|---|
DynamoDB LATEST
|
TTL 없음. 최신 상태 1건을 계속 overwrite |
DynamoDB HISTORY#STATE
|
item attribute ttl, 목표 2시간(마지막 문서화된 운영값 48시간) |
DynamoDB GRAPH#5M
|
item attribute ttl, 48시간 |
DynamoDB HISTORY#FAST
|
item attribute ttl, 6시간 |
DynamoDB HISTORY#SLOW
|
item attribute ttl, 24시간 |
DynamoDB ALERT#
|
item attribute ttl, 7일 |
S3 raw/
|
90일 후 Glacier IR 전환. 자동 삭제 아님 |
S3 processed/
|
365일 후 Standard-IA 전환. 자동 삭제 아님 |
S3 reports/, processed_agg/, image_snapshot/
|
현재 foundation lifecycle에 별도 prefix rule 없음 |
초기 구조는 state history를 장시간 저장하고 Dashboard가 window=24h 조회에서도 같은 prefix를 읽었다. factory_state 3초 주기와 3개 factory가 겹치면 24시간 조회가 DynamoDB Query 페이지 수십 개로 커지고, Dashboard Backend의 동시 요청 제한이 포화되어 504 cascade가 발생했다.
현재 구조는 read model을 나눈다.
| 조회 window | 조회 모델 | 이유 |
|---|---|---|
| 현재 상태 | LATEST |
한 factory당 1건 |
| 1h | HISTORY#STATE |
raw-resolution이 필요한 짧은 구간 |
| 6h/12h/24h | GRAPH#5M |
5분 bucket 집계로 query item 수를 고정 |
| 장기 보고서/audit | S3 processed/, processed_agg/, reports/
|
DynamoDB TTL 이후에도 보존 |
자세한 모델 설명은 데이터 read model을 따른다.
aegis-bucket-data S3 bucket과 AEGIS-DynamoDB-FactoryStatus DynamoDB table은 infra/foundation이 생성하고 소유한다. data-pipeline, reporting, dashboard 계층은 이 foundation 리소스를 참조해서 쓰거나 읽는다. 따라서 Dashboard/EKS/reporting stack을 destroy해도 데이터 bucket/table의 소유권은 foundation에 남는다.
- 시스템 아키텍처
- 제어 & 데이터 플레인
- 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 채팅 어시스턴트