Apache Flink 플러그인
Apache Flink 플러그인은 Flink 클러스터를 Konduo 리소스로 등록하고 JobManager REST API와 연결된 Prometheus 메트릭을 이용해 작업 실행, 체크포인트, 처리량, 백프레셔 및 JVM 상태를 보여줍니다.
주요 기능
- JobManager REST API로 클러스터와 실행 중인 작업의 현재 상태를 확인합니다.
- Prometheus 매핑팩으로 19개 Flink 논리 메트릭을 제공합니다.
- 16개 기본 패널에서 클러스터 용량, 처리량, 체크포인트 및 JVM 자원을 모니터링합니다.
- 현재 상태와 제한된 기간의 메트릭 근거를 분리해 진단합니다.
- 백프레셔, Task Slot, 체크포인트, 작업 재시작 및 JVM 압박에 대한 관리 경보 규칙을 제공합니다.
등록 전 준비
- Konduo 백엔드에서 접근 가능한 JobManager REST URL을 준비합니다. 기본값은
http://localhost:8081이지만 컨테이너 환경에서localhost는 Konduo 백엔드 컨테이너 자체를 의미할 수 있습니다. - REST API가 Bearer 인증을 요구하면 API 토큰을 준비합니다. 토큰은 민감 설정으로 저장되며 진단 요청의
Authorization헤더에만 사용됩니다. - JobManager와 모든 TaskManager에서 Prometheus reporter를 활성화하고 Prometheus가 해당 엔드포인트를 수집하는지 확인합니다.
- Flink 리소스에 대상 Prometheus 리소스를 메트릭 소스로 연결합니다.
metrics.scope.*또는 reporter 접두어를 변경했다면 기본 매핑팩 대신 환경에 맞는 매핑 규칙을 준비합니다.
연결 설정
| 항목 | 필수 여부 | 기본값 | 설명 |
|---|---|---|---|
endpoint | 필수 | http://localhost:8081 | JobManager REST API의 기준 URL |
api_token | 선택 | 없음 | REST API 요청에 사용하는 Bearer 토큰 |
| 메트릭 소스 연결 | 메트릭 사용 시 필수 | 없음 | 시계열 조회를 수행할 prometheus-plugin 리소스 |
브라우저에서 열리는 URL이라도 Konduo 백엔드 네트워크에서 접근할 수 없으면 진단은 실패합니다. 컨테이너를 사용하는 경우 서비스 DNS 이름, 포트 공개 범위 및 방화벽을 함께 확인하십시오.
등록 및 검증 절차
endpoint와 필요한 경우api_token을 입력하고 연결 테스트를 실행합니다.- 리소스를 저장한 뒤 클러스터 진단에서 Flink 버전, 등록 TaskManager와 전체·사용 가능 Task Slot을 확인합니다.
- Prometheus에서 JobManager와 모든 TaskManager target이
UP인지 확인합니다. - 해당 Prometheus 리소스를 Flink의 메트릭 소스로 연결합니다.
- 실행 중인 작업이 없어도 클러스터와 JVM 패널에 최신 샘플이 표시되는지 확인합니다.
- 시험 작업이 있다면 입력·출력 처리율, busy time, 백프레셔와 체크포인트 패널의 작업·태스크 라벨이 올바른지 확인합니다.
- 관리 경보 9개가 등록되었는지 확인하고, 기본 비활성 상태인 체크포인트 경과 시간·소요 시간 규칙은 워크로드 기준을 정한 뒤 활성화합니다.
REST 연결 성공과 Prometheus 수집 성공은 별개의 검증입니다. REST 진단이 정상이더라도 reporter 또는 메트릭 소스 연결이 잘못되면 대시보드와 기간 기반 진단은 비어 있을 수 있습니다.
대시보드와 메트릭
기본 대시보드는 다음 네 영역으로 구성됩니다.
- 클러스터 개요: 등록 TaskManager, 실행 중인 작업, 전체·사용 가능 Task Slot과 사용률
- 워크로드와 백프레셔: 초당 입력·출력 레코드, 태스크 busy time, 백프레셔
- 체크포인트와 복구: 체크포인트 소요 시간·경과 시간·실패 및 작업 재시작
- JVM 자원: JobManager와 TaskManager CPU 및 힙 사용률
작업, 태스크 및 체크포인트 패널은 실행 중인 작업이 없을 때 정상적으로 데이터가 없을 수 있습니다. 이 경우 클러스터 개요와 JVM 패널에 데이터가 있다면 전체 메트릭 수집 실패로 판단하지 않습니다.
기본 매핑은 Flink Prometheus reporter의 flink_ 접두어와 기본 논리 scope를 가정합니다. 패널 전체 또는 특정 역할의 패널이 비어 있으면 Prometheus target 상태, 실제 메트릭 이름, 라벨과 적용된 매핑팩을 차례로 확인하십시오.
진단 화면
진단 요약은 현재 REST 근거와 메트릭 근거 준비 상태를 함께 보여주며, 상세 화면은 다음 영역으로 나뉩니다.
- 클러스터: JobManager REST 접근성, 등록 TaskManager 및 Task Slot 용량
- 작업: 실행 중인 작업 상태와 실패·재시작 징후
- 체크포인트: 활성 작업의 최근 체크포인트 결과
- 성능 근거: 백프레셔, 처리량, 체크포인트 및 JVM 논리 메트릭
현재 진단은 /overview, /jobs/overview, /jobs/{jobid}/checkpoints를 제한 시간과 응답 크기 한도 안에서 조회합니다. 현재 REST 진단의 전체 제한 시간은 8초이고 응답별 최대 크기는 2 MiB이며, 체크포인트 조회는 활성 작업 최대 10개로 제한됩니다.
FAILED 작업은 과거 상태 맥락으로 표시되지만 그 사실만으로 현재 경고를 만들지 않습니다. FAILING, RESTARTING, SUSPENDED처럼 현재 실행에 영향을 주는 상태가 실시간 안정성 경고 대상입니다. 배치 또는 stateless 작업은 체크포인트를 사용하지 않을 수 있으므로 체크포인트 활동이 없다는 이유만으로 실패로 판정하지 않습니다.
기간 기반 진단은 연결된 Prometheus에 논리 메트릭 조회를 위임합니다. 메트릭 소스가 없거나 조회할 수 없으면 정상으로 추정하지 않고 사용 불가 근거와 후속 조회 경로를 표시합니다.
진단 실행 버튼은 현재 근거를 다시 수집할 뿐 결과 이력을 저장하지 않습니다. 플러그인 재시작 여부와 관계없이 이전 진단 결과를 조회하는 영구 저장소는 제공하지 않습니다.
경보 규칙
flink-alert-rules-v1은 다음 9개 규칙을 제공합니다.
| 규칙 | 심각도 | 기본 조건 | 지속 조건 | 기본 상태 |
|---|---|---|---|---|
| 태스크 백프레셔 높음 | 경고 | 5분 구간에서 50% 초과 | 3분 | 활성 |
| Task Slot 사용률 높음 | 경고 | 10분 구간에서 90% 이상 | 5분 | 활성 |
| 체크포인트 실패 감지 | 심각 | 5분 구간에서 1회 이상 증가 | 1분 | 활성 |
| 작업 재시작 감지 | 경고 | 10분 구간에서 1회 이상 증가 | 2분 | 활성 |
| 최근 체크포인트 경과 시간 높음 | 경고 | 10분 구간에서 900초 이상 | 5분 | 비활성 |
| 체크포인트 소요 시간 높음 | 경고 | 10분 구간에서 60초 이상 | 5분 | 비활성 |
| TaskManager JVM 힙 사용률 높음 | 경고 | 5분 구간에서 85% 이상 | 3분 | 활성 |
| JobManager JVM 힙 사용률 높음 | 경고 | 5분 구간에서 85% 이상 | 3분 | 활성 |
| JobManager JVM CPU 사용률 높음 | 경고 | 10분 구간에서 90% 이상 | 5분 | 활성 |
체크포인트 경과 시간과 소요 시간 규칙은 작업별 주기와 상태 크기에 따라 정상 범위가 크게 달라 기본적으로 비활성화됩니다. 활성화하기 전에 대상 작업의 체크포인트 주기, 상태 백엔드와 정상 소요 시간을 기준으로 임계값을 조정하십시오.
경보는 단독 결론이 아니라 조사 시작점입니다. 예를 들어 백프레셔 경보는 다운스트림 처리량, busy time, 체크포인트 소요 시간과 TaskManager 용량을 함께 확인해야 합니다.
대표 점검 절차
작업이 반복해서 재시작될 때
- 작업 상세에서 현재 상태와 최근 예외를 확인합니다.
- 최근 체크포인트 실패와 마지막 성공 체크포인트 시각을 확인합니다.
- 체크포인트 스토리지 지연과 접근성을 확인합니다.
- TaskManager CPU·힙과 백프레셔가 같은 시각에 증가했는지 비교합니다.
- 근거를 확인한 후 재시작 전략 또는 병렬도를 조정합니다.
처리량이 감소할 때
- 입력과 출력 레코드 비율을 비교합니다.
- 백프레셔가 높은 연산자와 busy time이 높은 태스크를 찾습니다.
- Task Slot 여유와 TaskManager별 자원 편차를 확인합니다.
- sink 또는 외부 의존 서비스의 지연을 확인한 뒤 확장 여부를 결정합니다.
체크포인트가 느리거나 실패할 때
- 영향을 받는 작업과 최근 성공·실패 체크포인트를 확인합니다.
- 체크포인트 소요 시간, 경과 시간 및 실패 증가량을 함께 봅니다.
- 상태 크기, 정렬 지연, 백프레셔와 스토리지 상태를 비교합니다.
- 실패 원인을 확인하기 전에는 반복 재시작으로 문제를 숨기지 않습니다.
관리 경계
- CE 플러그인은 관찰과 진단을 제공하며 작업 배포, 취소, savepoint, rescale 같은 Flink 작업 변경을 자동 실행하지 않습니다.
- 메트릭이 없다는 사실을 Flink 장애로 단정하지 않습니다. REST 상태와 메트릭 수집 상태를 분리해 판단합니다.
- API 토큰을 진단 결과, 로그 또는 화면 설명에 노출하지 않습니다.
- 기본 경보 임계값은 출발점입니다. 워크로드 특성에 맞게 검토한 뒤 변경하십시오.
문제 해결
| 증상 | 확인 순서 |
|---|---|
| REST 진단만 실패 | endpoint, JobManager 상태, 백엔드 네트워크 경로와 방화벽을 확인 |
REST가 401 또는 403 반환 | Bearer 토큰 값, 만료, 앞단 프록시의 인증 전달과 대상 REST 경로 권한을 확인 |
| REST 진단 제한 시간 초과 | 백엔드에서 JobManager까지의 연결 지연, 프록시, JobManager 부하를 확인; 진단 제한 시간은 8초 |
| 응답 형식 오류 | 프록시 오류 페이지가 JSON 대신 반환되는지, REST API 버전과 응답 크기가 2 MiB 한도를 넘는지 확인 |
| 모든 메트릭 패널이 비어 있음 | Prometheus target, JobManager·TaskManager reporter와 메트릭 소스 연결을 확인 |
| 작업·체크포인트 패널만 비어 있음 | 실행 중인 작업과 체크포인트 사용 여부를 확인; 유휴 클러스터에서는 정상일 수 있음 |
| 일부 역할 또는 작업만 비어 있음 | 실제 metrics.scope.*, flink_ 접두어, 라벨과 매핑팩 규칙을 확인 |
| 기간 진단이 사용 불가 | Prometheus 연결과 해당 논리 메트릭의 범위 조회 가능 여부를 확인 |
Apache Flink Enterprise 확장
Apache Flink Enterprise 확장은 Community Flink 리소스 플러그인에 다중 신호 이상징후 규칙과 읽기 전용 MCP 설명자를 추가합니다. 연결, JobManager REST 진단, 대시보드, 논리 메트릭, 매핑 팩과 관리 경보는 Community 플러그인이 계속 소유합니다.
등록 전 확인
- Community 매뉴얼에 따라 JobManager REST 연결과 선택적인 Bearer 인증을 구성합니다.
- JobManager와 모든 TaskManager의 메트릭을 수집하는 Prometheus 리소스를 Flink 리소스에 연결합니다.
- Enterprise 라이선스에서
mcp.gateway와anomaly.engine기능이 활성화되어 있는지 확인합니다. - 이상징후 규칙을 조회하는 사용자에게
viewer이상의 역할과flink-plugin.anomaly.read권한을 부여합니다. - MCP 호출자에게 허용된 MCP 읽기 범위와 대상 Flink 리소스 인스턴스 접근 권한을 함께 부여합니다.
MCP와 이상징후 분석은 Flink REST 연결을 대체하지 않습니다. REST 상태, Prometheus 수집과 논리 메트릭 매핑을 각각 검증해야 합니다.
이상징후 규칙
flink-anomaly-rules-v1은 모든 조건을 함께 만족할 때 규칙을 일치시키는 5개 규칙을 제공합니다.
| 규칙 키 | 심각도·점수 | 조건 | 우선 확인 |
|---|---|---|---|
flink.worker_availability_risk | 심각·0.98 | 실행 작업 1개 이상, 등록 TaskManager 1개 미만 | TaskManager 등록·heartbeat·네트워크, JobManager 리더십과 Task Slot |
flink.checkpoint_restart_cascade | 심각·0.94 | 체크포인트 실패 증가량 0 초과, 작업 재시작 증가량 0 초과 | 작업 예외, 최근 성공·실패 체크포인트, 체크포인트 스토리지 |
flink.backpressure_capacity_saturation | 경고·0.86 | 백프레셔 50% 이상, Task Slot 사용률 90% 이상 | 병목 연산자, 업스트림·다운스트림 처리량, Slot과 병렬도 |
flink.taskmanager_resource_pressure | 경고·0.82 | TaskManager CPU 90% 이상, 힙 사용률 85% 이상 | TaskManager별 편차, GC, busy time과 백프레셔 |
flink.throughput_stall_pattern | 경고·0.80 | 입력 1건/초 이상, 출력 0.001건/초 이하, busy time 85% 이상 | sink, 비동기 I/O, 연산자 예외와 태스크별 쏠림 |
이 규칙은 Community 관리 경보를 대체하지 않습니다. 관리 경보는 개별 임계값과 지속 조건을 평가하고, Enterprise 규칙은 같은 평가 맥락에서 여러 논리 메트릭을 결합해 장애 분류 근거를 제공합니다.
실행 작업이 없는 유휴 클러스터에서는 작업·체크포인트 메트릭이 없을 수 있습니다. 누락된 시계열이나 해석할 수 없는 매핑을 정상 또는 0으로 간주하지 말고 메트릭 소스와 매핑 상태를 먼저 확인합니다.
MCP 리소스와 도구
MCP 카탈로그는 8개 읽기 리소스와 10개 읽기 도구를 제공합니다.
| 종류 | 제공 항목 |
|---|---|
| 리소스 | 모니터링 개요, 현재 진단, 이력 진단, 메트릭 카탈로그, 매핑 팩 카탈로그, 관리 경보 규칙, 이상징후 규칙, 소프트웨어 인벤토리 |
| 진단 도구 | 모니터링 개요, 진단 요약, 클러스터, 작업, 체크포인트 및 성능 진단 읽기 |
| 메트릭 도구 | 논리 메트릭 쿼리 해석, 매핑 팩 해석 |
| 규칙 도구 | 관리 경보 규칙과 Enterprise 이상징후 규칙 읽기 |
metrics_query_resolve 도구는 다음 입력을 사용합니다.
| 입력 | 필수 여부 | 값 |
|---|---|---|
logical_metric_key | 필수 | Community Flink 메트릭 카탈로그에 있는 논리 키 |
query_mode | 선택 | instant 또는 range |
MCP 설명자는 조회 경로만 제공합니다. diagnostics/summary/run 진단 새로고침은 쓰기 권한이 필요한 별도 작업이므로 읽기 전용 MCP 카탈로그에 포함되지 않습니다. 작업 배포, 취소, savepoint와 rescale도 MCP에서 제공하지 않습니다.
운영 절차
워커 가용성 위험
- 현재 실행 작업과 등록 TaskManager 수를 확인합니다.
- JobManager REST 진단에서 TaskManager 등록과 Task Slot을 확인합니다.
- TaskManager 프로세스, heartbeat, 네트워크와 리소스 상태를 점검합니다.
- 워커 복구 후 사용 가능한 Slot과 작업 복구 상태를 확인한 뒤 재시작 여부를 결정합니다.
체크포인트·재시작 연쇄
- 영향받은 작업과 최근 예외를 확인합니다.
- 최근 성공·실패 체크포인트, 소요 시간과 스토리지 접근성을 비교합니다.
- 같은 구간의 백프레셔, TaskManager CPU·힙과 재시작 증가량을 확인합니다.
- 원인을 확인하기 전에 재시작 전략이나 병렬도를 반복 변경하지 않습니다.
처리량 정체
- 작업·태스크별 입력, 출력과 busy time을 같은 시간 구간에서 비교합니다.
- 출력이 멈춘 연산자 체인의 sink, 비동기 I/O와 외부 의존 서비스를 확인합니다.
- 백프레셔와 Slot 여유를 확인한 뒤 태스크 쏠림인지 전체 용량 부족인지 구분합니다.
문제 해결
| 증상 | 확인 순서 |
|---|---|
| EE 확장 또는 MCP 항목이 보이지 않음 | Enterprise 라이선스, mcp.gateway, anomaly.engine, 플러그인 버전과 contribution 로딩 상태를 확인 |
| 이상징후 규칙 조회가 거부됨 | 대상 리소스 접근 권한, viewer 역할과 flink-plugin.anomaly.read 권한을 확인 |
| 모든 규칙이 평가되지 않음 | Prometheus 리소스 연결, Flink reporter target, 논리 메트릭 매핑과 범위 조회를 확인 |
| 작업 관련 규칙만 평가되지 않음 | 실행 중인 작업, checkpoint 사용 여부, job/task 라벨과 기본 scope를 확인 |
| MCP 메트릭 해석 실패 | logical_metric_key가 현재 카탈로그에 있는지 확인하고 query_mode를 instant 또는 range로 지정 |
| MCP에서 진단 새로고침을 찾을 수 없음 | 의도된 읽기 전용 경계이며 Konduo의 RBAC 보호 진단 실행 경로를 사용 |
에디션 경계
이상징후 규칙 팩, MCP 설명자와 Enterprise 현지화는 EE overlay에 있습니다. Flink 작업 변경과 진단 실행은 읽기 전용 MCP를 우회하지 않으며 기존 Konduo API, RBAC와 감사 경계를 따릅니다.