실시간 파이프라인 UI 가이드
GenD 는 Bytewax 워커가 Kafka 를 소비하는 실시간 추론 작업을 정의·운영할 수 있는 실시간 파이프라인 을 제공합니다. 사이드바 ML 허브 → 배포 & 서빙 의 탭(배치 작업 다음)으로 위치하며, 작업 정의 → 수동 트리거 → 실행 이력 → DLQ / consumer lag 까지 한 곳에서 확인할 수 있습니다.
본 가이드는 Epic #1084 M3 (UI partial) 의 산출물이며, 백엔드 M2 (PR #1241) 의 read-trio + mutation-quartet API 위에서 동작합니다.
진입점
- 사이드바 → ML 허브 → 배포 & 서빙 → 실시간 파이프라인 탭
- 직접 URL:
https://gend.genon.ai/ml-hub/serving?tab=streaming(구/streaming는 자동 리디렉션)
기능 요약
| 기능 | 엔드포인트 | 권한 | 비고 |
|---|---|---|---|
| 목록 조회 | GET /api/v1/streaming/jobs | viewer+ | 워크스페이스 펜스 자동 적용 |
| 단건 조회 | GET /api/v1/streaming/jobs/{id} | viewer+ | 작업 contract 표시용 |
| 실행 이력 | GET /api/v1/streaming/jobs/{id}/runs | viewer+ | 최신 순 정렬 |
| 작업 생성 | POST /api/v1/streaming/jobs | admin | UNIQUE (workspace_id, name) |
| 작업 수정 | PUT /api/v1/streaming/jobs/{id} | admin | description / model_alias / image / status 만 가능 |
| 수동 실행 | POST /api/v1/streaming/jobs/{id}/trigger | admin | M2 는 queued row 만 삽입, K8s scale + KEDA 는 M3 후속 |
| 작업 삭제 | DELETE /api/v1/streaming/jobs/{id} | admin | Soft delete (status='deleted'), 실행 이력은 CASCADE 보존 |
참고: 비관리자는 목록 조회까지 가능하지만, 모든 mutating 액션 버튼은 비활성화됩니다. 백엔드
require_admin이 fail-closed 로 동작하므로 403 round-trip 도 발생하지 않습니다.
신규 작업 등록
- 페이지 헤더 우측 신규 작업 버튼 클릭.
- 다이얼로그에서 다음 필드를 입력:
- 이름 — 워크스페이스 내 유일.
anomaly-fraud같은 stable identifier. - 엔진 —
Bytewax(M2 기본값).Flink옵션은 disabled — M3 후속에서 활성화 예정. - 소스 토픽 — Kafka 토픽 (예:
streaming.gend.feature_events). - 싱크 토픽 — 성공 prediction 출력 토픽.
- DLQ 토픽 — parse / inference 실패 → dead-letter. 클러스터에 사전 생성되어 있어야 합니다.
- MLflow 모델 — 등록 모델 이름. 자동완성 datalist 제공.
- MLflow Alias — default
Production. 워커 boot 시 concrete version 으로 resolve. - Consumer 그룹 — Kafka offset 커서 식별자. 작업당 고유.
- Bytewax 이미지 — default
genosprodacr.azurecr.io/gend-streaming:latest. 생성 시 read-only (수정 모드에서만 변경 가능).
- 이름 — 워크스페이스 내 유일.
- 생성 버튼 클릭.
불변 필드 (Immutable on PUT)
다음 필드는 작업 생성 후 변경할 수 없습니다. 변경이 필요하면 새 작업으로 정의 해야 합니다.
name/enginesource_topic/sink_topic/dlq_topicmodel_name/consumer_group
이유:
source_topic재바인딩은 같은 작업이 완전히 다른 이벤트 스트림을 scoring 하게 만들고,consumer_group재바인딩은 Kafka offset 커서를 silent 하게 rewind 시킵니다.model_name재바인딩은 live 데이터플로우의 scoring logic 을 silent 하게 교체합니다. 감사 추적 보존을 위해 PUT 에서 422 fail-closed 로 거부합니다.
수정 가능 필드 (Mutable on PUT)
description— 자유 텍스트.model_alias— 새 alias 로 재 mapping. 운영 중인 Pod 는 boot 시 frozen 된 alias 를 유지하며, 재시작 후에만 새 alias 가 반영됩니다.bytewax_image— 업그레이드된 워커 이미지. ACR allowlist (genosprodacr.azurecr.io/또는 dev 의gend/) 만 허용.status—active/paused(deleted는 soft-delete 전용).
수동 실행 (Run Now)
- 작업 목록 행에서 지금 실행 클릭, 또는 상세 시트 헤더의 동일 버튼.
- 새
streaming_job_runrow 가status='queued'+ 요청한 replica count 로 삽입됩니다. - M2 단계에서는 K8s
apps/v1Deployment scale + KEDA ScaledObject 생성은 M3 후속으로 미뤄져 있어, 운영자는 ops 대시보드에서 queued row 를 polling 하여 수동 reconcile 할 수 있습니다.
Replica 상한: 1..16. 공유 클러스터 리소스 보호를 위해 16 으로 cap 되어 있으며, M3 에서
Workspace.streaming_replica_quota컬럼으로 per-workspace cap 이 도입됩니다.
실행 이력 + DLQ + Consumer Lag
상세 시트 하단의 실행 이력 테이블 컬럼:
| 컬럼 | 비고 |
|---|---|
| 시작 | started_at (queued 상태에서는 created_at 표시) |
| 상태 | queued / running / success / failed / cancelled |
| Replicas | 요청한 replica 수 |
| 처리 건수 | processed_count — 누적 성공 prediction |
| DLQ | dlq_count — 누적 실패 (0 초과 시 빨간색 강조) |
| Lag | consumer_lag — 최신 Kafka consumer 지연 (events) |
| 오류 | error_message (tooltip 로 풀버전 확인) |
DLQ 카운트는
dlq_topic으로 라우팅된 이벤트의 누적 카운트입니다. 0 초과 시 운영자는 DLQ 토픽을 검사하여 schema mismatch / 모델 inference 오류를 진단해야 합니다.
삭제 (Soft delete)
- 삭제 버튼 클릭 → 확인 다이얼로그 → 실행 시
status='deleted'로 flip. - 실행 이력 (
streaming_job_run) 은 FK CASCADE 로 보존됩니다. - 삭제된 작업은 수정 / 트리거가 모두 차단되며, UI 도 액션 버튼을 비활성화합니다.
알려진 제약 (M3 partial)
- 백엔드 M2 는 manual trigger 시 queued row 만 삽입합니다. K8s
apps/v1Deployment scale + KEDA ScaledObject 생성 + Bytewax worker boot 은 M3 후속 PR 에서 추가됩니다. consumer_lag메트릭은 M2 에서는 ORM 컬럼 stub — 실제 값은 M3 의 KEDA Kafka trigger 와 worker Prometheus emit 으로 채워집니다.engine은bytewax만 허용됩니다.flink는 M3 / M4 에서 ORM CheckConstraint 와 함께 확장됩니다 (다이얼로그의 disabled 옵션이 roadmap 노출 용도).- 작업 생성 시
bytewax_image는 read-only 로 default ACR FQDN 이 stamp 됩니다. 운영자 override 가 필요하면 수정 모드에서 변경하세요 (registry allowlist 위반 시 422).
관련 자료
- 백엔드 라우터:
apps/api/src/gend_api/routers/streaming_jobs.py - Pydantic 모델:
apps/api/src/gend_api/models/streaming_job.py - ORM 모델:
apps/api/src/gend_api/db/models/streaming.py - 배치 작업 UI 가이드:
ml-batch-ui.md