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 에 전달되지 않으면 체크포인트는 성공하지 않습니다.