Prometheus 싱크

Prometheus 싱크 (Prometheus Sink)

이 싱크 커넥터는 Remote Write Prometheus 인터페이스를 사용하여 데이터를 Prometheus 호환 저장소에 쓸 수 있습니다. PrometheusSink 빌더를 사용하여 구성하며, PrometheusTimeSeries 레코드를 입력으로 받습니다.

출처: 문서

본문

이 싱크 커넥터는 Remote Write Prometheus 인터페이스를 사용하여 데이터를 Prometheus 호환 저장소에 쓰는 데 사용할 수 있습니다.

Prometheus 호환 백엔드는 Remote Write 1.0 표준 API를 지원해야 하며, Remote Write 엔드포인트가 활성화되어 있어야 합니다.

이 커넥터는 내부 Flink 메트릭을 Prometheus에 보내기 위한 것이 아닙니다. Flink 클러스터의 상태와 운영을 모니터링하기 위해 Flink 메트릭을 게시하려면 Metric Reporters를 사용해야 합니다.

커넥터를 사용하려면 프로젝트에 다음 Maven 의존성을 추가하세요:

Flink 버전 2.3용 커넥터는 아직 없습니다.

사용법 (Usage)

Prometheus 싱크는 PrometheusSink 인스턴스를 빌드하는 빌더 클래스를 제공합니다. 아래 코드 스니펫은 기본 구성과 선택적 요청 서명자(request signer)로 PrometheusSink를 빌드하는 방법을 보여줍니다.

PrometheusSink sink = PrometheusSink.builder()
    .setPrometheusRemoteWriteUrl(prometheusRemoteWriteUrl)
    .setRequestSigner(new AmazonManagedPrometheusWriteRequestSigner(prometheusRemoteWriteUrl, prometheusRegion)) // Optional
    .build();

유일한 필수 구성은 prometheusRemoteWriteUrl입니다. 다른 모든 구성은 선택 사항입니다.

싱크의 병렬도가 1보다 크면 PrometheusTimeSeriesLabelsAndMetricNameKeySelector 키 선택기로 스트림을 키잉(keyed)하여 동일한 시계열의 모든 샘플이 같은 파티션에 있고 순서가 유지되도록 해야 합니다. 자세한 내용은 Sink parallelism과 keyed streams를 참고하세요.

입력 데이터 객체 (Input data objects)

싱크는 입력으로 PrometheusTimeSeries 레코드를 기대합니다. 입력 데이터를 싱크에 보내기 전에 map 또는 flatMap 연산자를 사용해 PrometheusTimeSeries로 변환해야 합니다.

PrometheusTimeSeries 인스턴스는 불변(immutable)이며 재사용할 수 없습니다. 빌더를 사용하여 인스턴스를 만들고 채울 수 있습니다.

PrometheusTimeSeries는 Remote Write 인터페이스로 보낼 때 단일 시계열 레코드를 나타냅니다. 각 시계열 레코드는 여러 샘플을 포함할 수 있습니다.

Prometheus에서 "time-series"라는 용어는 과적재되어 있습니다. 고유한 라벨 집합을 가진 샘플의 시리즈(기본 시계열 데이터베이스의 시계열)와 Remote Write 인터페이스로 보내는 레코드를 모두 의미합니다. PrometheusTimeSeries 인스턴스는 인터페이스로 보내는 레코드를 나타냅니다.

이 두 개념은 관련되어 있습니다. 같은 라벨 집합을 가진 시계열 "레코드"는 같은 "데이터베이스 시계열"로 보내지기 때문입니다.

PrometheusTimeSeries 레코드는 다음을 포함합니다:

  • 하나의 metricName. __name__ 라벨의 값으로 변환되는 문자열.
  • 0개 이상의 Label 항목. 각 라벨은 keyvalue(둘 다 String)를 가집니다. 라벨은 샘플의 추가 차원을 나타냅니다. 중복 라벨 키는 허용되지 않습니다.
  • 하나 이상의 Sample. 각 샘플은 측정값을 나타내는 value(double)와 Epoch로부터의 밀리초 단위 측정 시간을 나타내는 timestamp(long)를 가집니다. 같은 레코드의 중복 타임스탬프는 허용되지 않습니다.

다음 의사코드는 PrometheusTimeSeries 레코드의 구조를 나타냅니다:

 PrometheusTimeSeries
 +--> metricName <String>
 +--> Label [0..*] + name <String>
                   + value <String>
 +--> Sample [1..*] + timestamp <long>
                    + value <double>

라벨 집합과 metricName은 데이터베이스 시계열의 고유 식별자입니다. 모든 라벨과 metricName의 조합은 시계열당 순서를 보장하기 위해 Flink 애플리케이션 안과 업스트림 모두에서 데이터를 파티션하는 키이기도 합니다.

PrometheusTimeSeries 채우기 (Populating a PrometheusTimeSeries)

PrometheusTimeSeries는 빌더 인터페이스를 제공합니다.

PrometheusTimeSeries inputRecord = PrometheusTimeSeries.builder()
    .withMetricName(metricName)
    .addLabel("DeviceID", instanceId)
    .addLabel("RoomID", roomId)
    .addSample(measurement1, time1)
    .addSample(measurement2, time2)
    .build();

PrometheusTimeSeries 인스턴스는 여러 샘플을 포함할 수 있습니다. 각각에 대해 .addSample(...)을 호출하세요. 샘플이 추가되는 순서는 유지됩니다. 레코드당 최대 샘플 수는 maxBatchSizeInSamples 구성에 의해 제한됩니다.

여러 샘플을 단일 PrometheusTimeSeries 레코드로 집계하면 쓰기 성능이 향상될 수 있습니다.

Prometheus remote-write 제약 (Prometheus remote-write constraints)

Prometheus는 데이터 형식과 순서에 대해 엄격한 제약을 부과합니다. 이러한 제약을 위반하는 레코드를 포함하는 쓰기 요청은 거부됩니다.

이러한 제약에 대한 자세한 내용은 Remote Write specification을 참고하세요.

실제로 Prometheus 호환 백엔드에 데이터를 쓸 때의 동작은 Prometheus 구현과 구성에 따라 다릅니다. 어떤 경우에는 제약이 완화되어 Remote Write 사양을 위반하는 쓰기가 수락될 수도 있습니다.

이러한 이유로 이 커넥터는 데이터 제약을 직접 강제하지 않습니다. 사용자는 Prometheus 구현의 실제 제약을 위반하지 않는 데이터를 싱크로 보낼 책임이 있습니다. 자세한 내용은 User responsibilities 참고.

순서 제약 (Ordering constraints)

Remote Write 사양은 여러 순서 제약을 요구합니다:

  1. PrometheusTimeSeries 레코드 내의 라벨key 기준 사전순(lexicographical) 순서여야 합니다.
  2. PrometheusTimeSeries 레코드 내의 샘플은 오래된 것부터 새로운 것 순으로 timestamp 순서여야 합니다.
  3. 같은 시계열(고유한 라벨과 metricName 집합)에 속한 모든 샘플timestamp 순서로 쓰여야 합니다.
  4. 같은 시계열 내에서 같은 타임스탬프를 가진 중복 샘플은 허용되지 않습니다.

Prometheus 호환 백엔드 구현이 out-of-order time windows를 지원하고 그 옵션이 활성화되어 있으면 샘플 순서 제약이 완화됩니다. 구성된 윈도우 내에서 순서가 뒤바뀐 데이터를 보낼 수 있습니다.

형식 제약 (Format constraints)

싱크로 보내는 PrometheusTimeSeries 레코드도 다음 제약을 준수해야 합니다:

  • **metricName**은 정의되고 비어 있지 않아야 합니다. 커넥터는 이 속성을 __name__ 라벨의 값으로 변환합니다.
  • 라벨 이름은 정규식 [a-zA-Z:_]([a-zA-Z0-9_:])을 따라야 합니다. 특히 @, $, !, .(점), 또는 (콜론 :과 하이픈 -을 제외한) 다른 구두점을 포함하는 라벨 이름은 유효하지 않습니다.
  • 라벨 이름__(이중 밑줄)로 시작해서는 안 됩니다. 이 라벨 이름은 예약되어 있습니다.
  • 중복 라벨 이름은 허용되지 않습니다.
  • 라벨 metricName은 모든 UTF-8 문자를 포함할 수 있습니다.
  • 라벨 은 비어 있을 수 없습니다(null 또는 빈 문자열).

PrometheusTimeSeries 빌더는 이러한 제약을 강제하지 않습니다.

사용자 책임 (User responsibilities)

사용자는 Prometheus 구현이 요구하는 형식 및 순서 제약을 준수하는 레코드(PrometheusTimeSeries)를 싱크로 보낼 책임이 있습니다. 커넥터는 검증이나 재정렬을 수행하지 않습니다.

타임스탬프 순서로 샘플 정렬은 특히 중요합니다. 같은 시계열(같은 라벨 집합과 메트릭 이름)에 속한 샘플은 타임스탬프 순서로 쓰여야 합니다. 소스 데이터는 순서대로 생성되어야 합니다. 싱크 전에도 순서가 유지되어야 합니다. 데이터를 파티션할 때 순서를 유지하려면 같은 라벨과 메트릭 이름 집합을 가진 레코드를 같은 파티션으로 보내야 합니다.

Remote Write 엔드포인트에 쓰여진 형식이 잘못되었거나 순서가 잘못된 레코드는 싱크에 의해 거부되고 버려집니다. 이는 데이터 손실을 일으킬 수 있습니다.

싱크로 보내진 순서를 위반하는 모든 레코드는 버려지고 같은 쓰기 요청에 일괄 처리된 다른 레코드가 버려질 수 있습니다. 자세한 내용은 Connector guarantees 참고.

싱크 병렬도와 키드 스트림 (Sink parallelism and keyed streams)

각 싱크 연산자 서브태스크는 단일 스레드를 사용하여 Remote Write 엔드포인트에 쓰기 요청을 보내며, PrometheusTimeSeries 레코드는 서브태스크가 받은 순서 그대로 쓰입니다.

같은 시계열(즉, LabelmetricName 목록이 동일한 PrometheusTimeSeries)에 속한 모든 레코드가 같은 싱크 서브태스크에 의해 쓰여지도록 보장하려면 PrometheusTimeSeries 스트림을 PrometheusTimeSeriesLabelsAndMetricNameKeySelector로 키잉해야 합니다.

DataStream<MyRecord> inputRecords;
// ...
KeyedStream<PrometheusTimeSeries> timeSeries = inputRecords
    .map(new MyRecordToTimeSeriesMapper())
    .keyBy(new PrometheusTimeSeriesLabelsAndMetricNameKeySelector());
timeSeries.sinkTo(prometheusSink);

이 키 선택기를 사용하면 싱크 연산자 앞의 재파티셔닝으로 인한 우발적인 무순서를 방지할 수 있습니다. 그러나 사용자는 레코드를 올바르게 파티션하여 이 시점까지 순서를 유지해야 합니다.

오류 처리 (Error handling)

이 단락은 Remote Write 엔드포인트에 데이터를 쓸 때 오류 조건의 처리를 다룹니다.

네 가지 유형의 오류 조건이 있습니다:

  1. Remote-Write 서버의 일시적 오류 조건 또는 스로틀링으로 인한 재시도 가능 오류: 5xx 또는 429 http 응답, 연결 문제.
  2. 제약을 위반하는 데이터, 형식이 잘못된 데이터 또는 무순서 샘플로 인한 재시도 불가 오류: 4xx http 응답(429, 403, 404 제외).
  3. 치명적 오류 응답: 인증 실패(403 http 응답) 또는 잘못된 엔드포인트 경로(404 http 응답).
  4. Prometheus 엔드포인트에 쓰는 동안의 예외로 인한 기타 예기치 않은 실패.

오류 시 동작 (On-error behaviors)

위 오류 중 하나가 발생하면 커넥터는 다음 두 동작 중 하나를 구현합니다:

  1. FAIL: 처리되지 않은 예외를 던져 작업이 실패합니다
  2. DISCARD_AND_CONTINUE: 오류를 일으킨 요청을 버리고 다음 레코드로 계속합니다

쓰기 요청이 DISCARD_AND_CONTINUE로 버려지면 다음이 모두 발생합니다:

  1. 오류의 원인과 함께 WARN 수준 메시지를 로그합니다. 엔드포인트의 응답으로 인한 오류라면 엔드포인트 응답의 페이로드가 포함됩니다.
  2. 거부된 샘플과 쓰기 요청 수를 세는 카운터 메트릭을 증가시킵니다.
  3. 전체 쓰기 요청을 버립니다. 배칭으로 인해 쓰기 요청이 여러 PrometheusTimeSeries를 포함할 수 있음에 유의하세요.
  4. 다음 입력 레코드로 계속합니다.

Prometheus Remote Write는 부분 실패를 지원하지 않습니다. 배칭으로 인해 단일 쓰기 요청이 여러 입력 레코드(PrometheusTimeSeries)를 포함할 수 있습니다. 요청에 단일 위반 레코드라도 있으면 전체 쓰기 요청(전체 배치)을 버려야 합니다.

재시도 가능 오류 응답 (Retryable error responses)

전형적인 재시도 가능 오류 조건은 엔드포인트 스로틀링이며 429, Too Many Requests 응답과 함께 발생합니다.

재시도 가능 오류에서 커넥터는 구성 가능한 백오프 전략으로 재시도합니다. 최대 재시도 횟수를 초과하면 쓰기 요청이 실패합니다. 이 시점에 무엇이 일어나는지는 onMaxRetryExceeded 오류 처리 구성에 달려 있습니다.

  • onMaxRetryExceededFAIL(기본값)이면 작업이 실패하고 체크포인트에서 재시작합니다.
PrometheusSink sink = PrometheusSink.builder()
    // ...
    .setRequestSigner(new AmazonManagedPrometheusWriteRequestSigner(prometheusRemoteWriteUrl, prometheusRegion))
    .build();

HTTP 클라이언트 구성 (HTTP client configuration)

Remote Write 엔드포인트에 쓰기 요청을 보내는 HTTP 클라이언트를 구성할 수 있습니다.

  • socketTimeoutMs: (기본값: 5000 밀리초) HTTP 클라이언트 소켓 타임아웃
  • httpUserAgent: (기본값: Flink-Prometheus) User-Agent 헤더
PrometheusSink sink = PrometheusSink.builder()
    // ...
    .setSocketTimeoutMs(5000)
    .setHttpUserAgent(USER_AGENT)
    .build();

커넥터 메트릭 (Connector metrics)

커넥터는 엔드포인트에 성공적으로 쓰여진 데이터와 DISCARD_AND_CONTINUE로 인해 버려진 데이터를 세는 사용자 지정 메트릭을 노출합니다.

메트릭 이름 설명
numSamplesOut Prometheus에 성공적으로 쓰여진 샘플
numWriteRequestsOut Prometheus에 성공적으로 쓰여진 쓰기 요청
numWriteRequestsRetries 재시도 가능 오류(예: 스로틀링)로 인한 쓰기 요청 재시도 수
numSamplesDropped DISCARD_AND_CONTINUE로 인해 버려진(데이터 손실!) 샘플
numSamplesNonRetryableDropped onPrometheusNonRetryableErrorDISCARD_AND_CONTINUE(기본값)로 설정되어 버려진(데이터 손실!) 샘플
numSamplesRetryLimitDropped onMaxRetryExceededDISCARD_AND_CONTINUE로 설정되었을 때 재시도 한도 초과로 버려진(데이터 손실!) 샘플
numWriteRequestsPermanentlyFailed 어떤 이유로든 영구적으로 실패한 쓰기 요청

numByteSend 메트릭은 무시해야 합니다. 이 커넥터가 기반하는 AsyncSink의 제한으로 인해 이 메트릭은 실제 바이트를 측정하지 않습니다. 싱크의 실제 출력을 모니터링하려면 numSamplesOutnumWriteRequestsOut을 사용하세요.

메트릭 그룹 이름은 기본적으로 "Prometheus"입니다. 변경할 수 있습니다:

PrometheusSink sink = PrometheusSink.builder()
    // ...
    .setMetricGroupName("my-metric-group")
    .build();

커넥터 보장 (Connector guarantees)

이 커넥터는 at-most-once 보장을 제공합니다. 특히 입력 데이터가 형식이 잘못되었거나 순서가 없으면 데이터 손실이 발생할 수 있습니다.

DISCARD_AND_CONTINUE 오류 시 동작으로 인해 데이터가 손실될 수 있습니다. 이 동작은 최대 재시도 횟수 초과 시 선택적으로 활성화될 수 있지만, 재시도 불가 오류 조건이 발생하면 항상 활성화됩니다.

이 동작은 같은 시계열에서 무순서 쓰기를 허용하지 않는 Prometheus Remote Write 인터페이스의 설계 때문입니다. 무순서 데이터가 발생할 때 작업을 체크포인트 재시작의 무한 루프에 빠뜨리지 않도록 하기 위해 버리고 계속하는 것입니다. 무순서 쓰기는 작업이 체크포인트에서 재시작할 때도 발생합니다. discard and continue가 없으면 싱크가 작업이 전혀 체크포인트에서 복구되는 것을 허용하지 않을 것입니다.

동시에 Prometheus는 시계열별 타임스탬프 순서를 부과합니다. 싱크는 파티션별로 순서가 유지됨을 보장합니다. PrometheusTimeSeriesLabelsAndMetricNameKeySelector로 키잉하면 입력이 시계열별로 파티션되고 쓰기 중 우발적 재정렬이 발생하지 않음을 보장합니다. 사용자는 싱크 전에 데이터를 파티션하여 순서가 유지되도록 할 책임이 있습니다.

예제 애플리케이션 (Example application)

이 싱크의 구성과 사용법을 보여주는 완전한 애플리케이션은 커넥터의 테스트에서 찾을 수 있습니다.

org.apache.flink.connector.prometheus.sink.examples.DataStreamExample 소스를 확인하세요.

이 클래스는 내부에서 임의 데이터를 생성하고 Prometheus에 쓰는 완전한 애플리케이션을 포함합니다.

더 알아보기 (Learn more)