Google Cloud PubSub

Google Cloud PubSub

이 커넥터는 Google Cloud PubSub 에서 읽고 쓸 수 있는 Source 와 Sink 를 제공합니다.

출처: 문서

본문

이 커넥터는 Google Cloud PubSub 에서 읽고 쓸 수 있는 Source 와 Sink 를 제공합니다. 이 커넥터를 사용하려면 프로젝트에 다음 의존성을 추가하세요:

Flink 2.3 버전용 커넥터는 (아직) 사용할 수 없습니다.

참고: 이 커넥터는 최근에 Flink 에 추가되었습니다. 아직 광범위한 테스트를 받지 못했습니다.

스트리밍 커넥터는 현재 바이너리 배포에 포함되어 있지 않습니다. 클러스터 실행을 위해 라이브러리와 함께 프로그램을 패키징하는 방법은 여기 를 참고하세요.

PubSubMessages 소비 또는 생성

이 커넥터는 Google PubSub 에서 메시지를 받고 보내기 위한 커넥터를 제공합니다. Google PubSub 은 at-least-once 보장을 가지며, 커넥터도 동일한 보장을 제공합니다.

PubSub SourceFunction

PubSubSource 클래스는 PubSubsource 를 만드는 빌더를 가지고 있습니다: PubSubSource.newBuilder(...).

PubSubSource 가 생성되는 방식을 변경하는 여러 선택 메서드가 있으며, 최소한 Google 프로젝트, Pubsub 구독, PubSubMessages 를 역직렬화하는 방법을 제공해야 합니다.

예제:

Java:

StreamExecutionEnvironment streamExecEnv = StreamExecutionEnvironment.getExecutionEnvironment();

DeserializationSchema<SomeObject> deserializer = (...);
SourceFunction<SomeObject> pubsubSource = PubSubSource.newBuilder()
                                                      .withDeserializationSchema(deserializer)
                                                      .withProjectName("project")
                                                      .withSubscriptionName("subscription")
                                                      .build();

streamExecEnv.addSource(pubsubSource);

현재 소스 함수는 PubSub 에서 메시지를 pull 하며, push endpoints 는 지원되지 않습니다.

PubSub Sink

PubSubSink 클래스는 PubSubSink 를 만드는 빌더를 가지고 있습니다. PubSubSink.newBuilder(...).

이 빌더는 PubSubSource 와 비슷한 방식으로 동작합니다.

예제:

DataStream<SomeObject> dataStream = (...);

SerializationSchema<SomeObject> serializationSchema = (...);
SinkFunction<SomeObject> pubsubSink = PubSubSink.newBuilder()
                                                .withSerializationSchema(serializationSchema)
                                                .withProjectName("project")
                                                .withSubscriptionName("subscription")
                                                .build()

dataStream.addSink(pubsubSink);

Google 자격 증명

Google 은 애플리케이션이 Google Cloud Platform 리소스(예: PubSub)를 사용할 수 있도록 인증하고 승인하는 데 Credentials 를 사용합니다.

두 빌더 모두 이러한 자격 증명을 제공할 수 있게 하지만, 기본적으로 커넥터는 자격 증명이 포함된 파일을 가리켜야 하는 환경 변수 GOOGLE_APPLICATION_CREDENTIALS 를 찾습니다.

자격 증명을 수동으로 제공하려면, 예를 들어 자격 증명을 외부 시스템에서 직접 읽는다면 PubSubSource.newBuilder(...).withCredentials(...) 를 사용할 수 있습니다.

통합 테스트

통합 테스트를 실행할 때 PubSub 에 직접 연결하고 싶지 않을 수 있으며 docker 컨테이너를 사용해 읽고 쓸 수 있습니다. (참고: PubSub locally 테스트)

다음 예제는 에뮬레이터에서 메시지를 읽어 다시 보내는 소스를 만드는 방법을 보여줍니다:

String hostAndPort = "localhost:1234";
DeserializationSchema<SomeObject> deserializationSchema = (...);
SourceFunction<SomeObject> pubsubSource = PubSubSource.newBuilder()
                                                      .withDeserializationSchema(deserializationSchema)
                                                      .withProjectName("my-fake-project")
                                                      .withSubscriptionName("subscription")
                                                      .withPubSubSubscriberFactory(new PubSubSubscriberFactoryForEmulator(hostAndPort, "my-fake-project", "subscription", 10, Duration.ofSeconds(15), 100))
                                                      .build();
SerializationSchema<SomeObject> serializationSchema = (...);
SinkFunction<SomeObject> pubsubSink = PubSubSink.newBuilder()
                                                .withSerializationSchema(serializationSchema)
                                                .withProjectName("my-fake-project")
                                                .withSubscriptionName("subscription")
                                                .withHostAndPortForEmulator(hostAndPort)
                                                .build();

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(pubsubSource)
   .addSink(pubsubSink);

최소 한 번(at-least-once) 보장

SourceFunction

메시지가 여러 번 전송될 수 있는 이유는 여러 가지가 있습니다. 예를 들어 Google PubSub 쪽의 실패 시나리오가 그렇습니다.

또 다른 이유는 ack(승인) 데드라인이 지났을 때입니다. 이는 메시지를 받은 시점과 메시지를 승인하는 시점 사이의 시간입니다. PubSubSource 는 at-least-once 를 보장하기 위해 성공적인 체크포인트에서만 메시지를 승인합니다. 이는 성공적인 체크포인트 사이의 시간이 구독의 ack 데드라인보다 크면 메시지가 대부분 여러 번 처리될 것임을 의미합니다.

이러한 이유로 체크포인트 간격을 ack 데드라인보다 (훨씬) 낮게 설정하는 것이 권장됩니다.

구독의 ack 데드라인을 늘리는 방법은 PubSub 를 참고하세요.

참고: PubSubMessagesProcessedNotAcked 메트릭은 승인되기 전에 다음 체크포인트를 기다리는 메시지 수를 보여줍니다.

SinkFunction

sink 함수는 성능상 이유로 PubSub 에 보낼 메시지를 짧은 시간 동안 버퍼링합니다. 각 체크포인트 전에 이 버퍼는 플러시되며, 메시지가 PubSub 에 전달되지 않으면 체크포인트는 성공하지 않습니다.

더 알아보기 (Learn more)