본문으로 건너뛰기

실시간 파이프라인 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/jobsviewer+워크스페이스 펜스 자동 적용
단건 조회GET /api/v1/streaming/jobs/{id}viewer+작업 contract 표시용
실행 이력GET /api/v1/streaming/jobs/{id}/runsviewer+최신 순 정렬
작업 생성POST /api/v1/streaming/jobsadminUNIQUE (workspace_id, name)
작업 수정PUT /api/v1/streaming/jobs/{id}admindescription / model_alias / image / status 만 가능
수동 실행POST /api/v1/streaming/jobs/{id}/triggeradminM2 는 queued row 만 삽입, K8s scale + KEDA 는 M3 후속
작업 삭제DELETE /api/v1/streaming/jobs/{id}adminSoft delete (status='deleted'), 실행 이력은 CASCADE 보존

참고: 비관리자는 목록 조회까지 가능하지만, 모든 mutating 액션 버튼은 비활성화됩니다. 백엔드 require_admin 이 fail-closed 로 동작하므로 403 round-trip 도 발생하지 않습니다.

신규 작업 등록

  1. 페이지 헤더 우측 신규 작업 버튼 클릭.
  2. 다이얼로그에서 다음 필드를 입력:
    • 이름 — 워크스페이스 내 유일. 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 (수정 모드에서만 변경 가능).
  3. 생성 버튼 클릭.

불변 필드 (Immutable on PUT)

다음 필드는 작업 생성 후 변경할 수 없습니다. 변경이 필요하면 새 작업으로 정의 해야 합니다.

  • name / engine
  • source_topic / sink_topic / dlq_topic
  • model_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/) 만 허용.
  • statusactive / paused (deleted 는 soft-delete 전용).

수동 실행 (Run Now)

  1. 작업 목록 행에서 지금 실행 클릭, 또는 상세 시트 헤더의 동일 버튼.
  2. streaming_job_run row 가 status='queued' + 요청한 replica count 로 삽입됩니다.
  3. M2 단계에서는 K8s apps/v1 Deployment 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
DLQdlq_count — 누적 실패 (0 초과 시 빨간색 강조)
Lagconsumer_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/v1 Deployment scale + KEDA ScaledObject 생성 + Bytewax worker boot 은 M3 후속 PR 에서 추가됩니다.
  • consumer_lag 메트릭은 M2 에서는 ORM 컬럼 stub — 실제 값은 M3 의 KEDA Kafka trigger 와 worker Prometheus emit 으로 채워집니다.
  • enginebytewax 만 허용됩니다. 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