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을 번역하기 위해 NettyNetty 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 스펙

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 스펙을 참고하세요.

더 알아보기 (Learn more)