REST API
REST API
Flink에는 실행 중인 작업의 상태와 통계, 그리고 최근 완료된 작업을 질의하는 데 사용할 수 있는 모니터링 API가 있어요. 이 모니터링 API는 Flink 자체 대시보드에서 사용되지만, 커스텀 모니터링 도구에서도 사용하도록 설계됐어요.
출처: 문서
본문
모니터링 API는 HTTP 요청을 받고 JSON 데이터로 응답하는 REST-ful API예요.
개요 (Overview)
모니터링 API는 JobManager의 일부로 실행되는 웹 서버에 의해 뒷받침돼요. 기본적으로 이 서버는 포트 8081에서 수신 대기하며, 이는 Flink configuration file에서 rest.port로 구성할 수 있어요. 모니터링 API 웹 서버와 웹 대시보드 웹 서버는 현재 동일하며 같은 포트에서 함께 실행된다는 점에 유의하세요. 다만 서로 다른 HTTP URL에 응답해요.
여러 JobManager가 있는 경우(고가용성을 위해), 각 JobManager는 자체 모니터링 API 인스턴스를 실행하며, 해당 JobManager가 클러스터 리더로 선출된 동안 완료되고 실행 중인 작업에 대한 정보를 제공해요.
개발 (Developing)
REST API 백엔드는 flink-runtime 프로젝트에 있어요. 핵심 클래스는 org.apache.flink.runtime.webmonitor.WebMonitorEndpoint로, 서버와 요청 라우팅을 설정해요.
REST 요청을 처리하고 URL을 번역하기 위해 Netty와 Netty Router 라이브러리를 사용해요. 이 선택은 이 조합이 가벼운 의존성을 가지고, Netty HTTP의 성능이 매우 좋기 때문이에요.
새 요청을 추가하려면 다음을 해야 해요:
- 새 요청의 인터페이스 역할을 하는 새
MessageHeaders클래스 추가, - 추가된
MessageHeaders클래스에 따라 요청을 처리하는 새AbstractRestHandler클래스 추가, org.apache.flink.runtime.webmonitor.WebMonitorEndpoint#initializeHandlers()에 핸들러 추가.
좋은 예는 org.apache.flink.runtime.rest.messages.JobExceptionsHeaders를 사용하는 org.apache.flink.runtime.rest.handler.job.JobExceptionsHandler예요.
API
REST API는 버전 관리되며, URL에 버전 접두사를 붙여 특정 버전을 질의할 수 있어요. 접두사는 항상 v[version_number] 형태예요.
예를 들어 /foo/bar의 버전 1에 접근하려면 /v1/foo/bar를 질의해요.
버전이 지정되지 않으면 Flink는 요청을 지원하는 가장 오래된 버전을 기본값으로 사용해요.
지원되지 않는/존재하지 않는 버전을 질의하면 404 오류가 반환돼요.
이 API 중에는 trigger savepoint, rescale a job 같은 여러 비동기 연산이 존재해요. 이들은 방금 POST한 연산을 식별하는 triggerid를 반환하고, 그 triggerid를 사용해 연산의 상태를 질의해야 해요.
(stop-with-)savepoint 연산의 경우, 연산을 트리거하는 요청의 본문에 설정해 이 triggerId를 제어할 수 있어요.
이를 통해 여러 savepoint를 트리거하지 않고 그러한 연산을 안전하게* 재시도할 수 있어요.
이 재시도는 async operation store duration이 경과할 때까지 안전해요.
JobManager
OpenAPI 스펙은 여전히 실험적이에요.
API 참조 (API reference)
v1
| /applications/overview | |
Verb: GET |
Response code: 200 OK |
| 모든 애플리케이션에 대한 개요를 반환해요. | |
Request {} |
|
Response { "type" : "object", "id" : "urn:jsonschema:org:apache:flink:runtime:messages:webmonitor:MultipleApplicationsDetails", "properties" : { "applications" : { "type" : "array", "items" : { "type" : "object", "id" : "urn:jsonschema:org:apache:flink:runtime:messages:webmonitor:ApplicationDetails", "properties" : { "duration" : { "type" : "integer" }, "end-time" : { "type" : "integer" }, "id" : { "type" : "any" }, "jobs" : { "type" : "object", "additionalProperties" : { "type" : "integer" } }, "name" : { "type" : "string" }, "start-time" : { "type" : "integer" }, "status" : { "type" : "string" } } } } } } |
REST API는 실행 중인 작업의 상태와 통계를 질의하기 위한 많은 엔드포인트를 제공해요. JobManager, TaskManager, 작업, 태스크, 체크포인트, savepoint, 메트릭 및 로그와 관련된 엔드포인트들이 포함돼요.
주요 엔드포인트 범주는 다음과 같아요:
- JobManager:
/overview,/config,/jobmanager,/jobmanager/metrics등 - 작업 (Jobs):
/jobs,/jobs/<jobid>,/jobs/<jobid>/exceptions,/jobs/<jobid>/checkpoints,/jobs/<jobid>/vertices/<vertexid>등 - TaskManager:
/taskmanagers,/taskmanagers/<taskmanagerid>등 - Metrics: 메트릭 단위/집계 질의
- Savepoints/Checkpoints:
POST /jobs/<jobid>/savepoints등 - 클러스터 연산:
/cluster리사이즈 등
전체 엔드포인트 목록과 각 엔드포인트의 요청/응답 스키마는 위 OpenAPI 스펙에서 확인할 수 있어요.
각 엔드포인트는 HTTP 메서드(GET/POST 등), 응답 코드, 요청 본문 스키마와 응답 본문 스키마를 정의해요. 예를 들어 GET /jobs/<jobid>는 특정 작업의 세부 정보를, POST /jobs/<jobid>/savepoints는 savepoint 트리거를, GET /jobs/<jobid>/metrics는 작업 메트릭을 반환해요.
자세한 요청/응답 스키마, 데이터 모델과 모든 엔드포인트에 대한 참조는 OpenAPI 스펙을 참고하세요.