MongoDB 커넥터

MongoDB 커넥터

Flink는 at-least-once 보장으로 MongoDB 컬렉션에서 데이터를 읽고 쓰기 위한 MongoDB 커넥터를 제공해요.

출처: MongoDB Connector

본문

이 커넥터를 사용하려면 프로젝트에 다음 의존성 중 하나를 추가하세요.

Flink 버전 2.3용 커넥터는 아직 제공되지 않아요.

MongoDB Source

아래 예시는 source를 구성하고 생성하는 방법을 보여줘요:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.connector.mongodb.source.MongoSource;
import org.apache.flink.connector.mongodb.source.reader.deserializer.MongoDeserializationSchema;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

import org.bson.BsonDocument;

MongoSource<String> source = MongoSource.<String>builder()
        .setUri("mongodb://user:***@127.0.0.1:27017")
        .setDatabase("my_db")
        .setCollection("my_coll")
        .setProjectedFields("_id", "f0", "f1")
        .setFetchSize(2048)
        .setLimit(10000)
        .setNoCursorTimeout(true)
        .setPartitionStrategy(PartitionStrategy.SAMPLE)
        .setPartitionSize(MemorySize.ofMebiBytes(64))
        .setSamplesPerPartition(10)
        .setDeserializationSchema(new MongoDeserializationSchema<String>() {
            @Override
            public String deserialize(BsonDocument document) {
                return document.toJson();
            }

            @Override
            public TypeInformation<String> getProducedType() {
                return BasicTypeInfo.STRING_TYPE_INFO;
            }
        })
        .build();

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

env.fromSource(source, WatermarkStrategy.noWatermarks(), "MongoDB-Source")
        .setParallelism(2)
        .print()
        .setParallelism(1);

Configurations

Flink의 MongoDB source는 정적 빌더 MongoSource.<OutputType>builder()로 만들어져요.

  • setUri(String uri) — 필수. MongoDB의 연결 문자열을 설정해요.
  • setDatabase(String database) — 필수. 읽을 데이터베이스의 이름.
  • setCollection(String collection) — 필수. 읽을 컬렉션의 이름.
  • setFetchSize(int fetchSize) — 선택. 기본값: 2048. 읽기 시 왕복당 가져와야 하는 문서 수를 설정해요.
  • setNoCursorTimeout(boolean noCursorTimeout) — 선택. 기본값: true. MongoDB 서버는 과도한 메모리 사용을 막기 위해 보통 비활성 기간(10분) 후 유휴 커서를 타임아웃시켜요. 이 옵션을 설정하면 그걸 방지해요. 세션이 30분 이상 유휴 상태면 MongoDB 서버는 그 세션을 만료된 것으로 표시하고 언제든 닫을 수 있어요. MongoDB 서버가 세션을 닫으면 세션과 관련된 진행 중인 작업과 열린 커서도 함께 종료해요. 여기에는 noCursorTimeout() 또는 30분보다 큰 maxTimeMS()로 구성된 커서도 포함돼요.
  • setPartitionStrategy(PartitionStrategy partitionStrategy) — 선택. 기본값: PartitionStrategy.DEFAULT. 파티션 전략을 설정해요. 사용 가능한 파티션 전략은 SINGLE, SAMPLE, SPLIT_VECTOR, SHARDED, DEFAULT예요. 자세한 내용은 Partition Strategies 섹션 참고.
  • setPartitionSize(MemorySize partitionSize) — 선택. 기본값: 64mb. MongoDB split의 파티션 메모리 크기를 설정해요. 파티션 메모리 크기에 따라 MongoDB 컬렉션을 여러 파티션으로 분할해요. 여러 reader가 파티션을 병렬로 읽어 전체 읽기 시간을 높일 수 있어요.
  • setSamplesPerPartition(int samplesPerPartition) — 선택. 기본값: 10. 파티션당 가져올 샘플 수를 설정하며 sample 파티션 전략 SAMPLE에만 사용돼요. sample 파티셔너는 컬렉션을 샘플링하고 파티션 필드로 프로젝션·정렬해요. 그런 다음 매 samplesPerPartition마다 값을 사용해 파티션 경계를 계산해요. 가져온 총 샘플 수는: 샘플 수/파티션 * (문서 수 / 파티션당 문서 수)예요.
  • setLimit(int limit) — 선택. 기본값: -1. 각 reader가 읽을 문서 한도를 설정해요. limit이 설정되지 않았거나 -1이면 전체 컬렉션의 문서를 읽어요. 읽기 병렬도를 1보다 크게 설정하면 읽을 최대 문서 수는 병렬도 * limit과 같아요.
  • setProjectedFields(String… projectedFields) — 선택. 읽을 문서의 프로젝션 필드를 설정해요. 설정하지 않으면 컬렉션의 모든 필드를 읽어요.
  • setDeserializationSchema(MongoDeserializationSchema deserializationSchema) — 필수. MongoDB BSON 문서를 파싱하려면 MongoDeserializationSchema가 필요해요.

Partition Strategies

파티션은 여러 reader가 병렬로 읽어 전체 읽기 시간을 단축할 수 있어요. 다음 파티션 전략이 제공돼요:

  • SINGLE: 전체 컬렉션을 단일 파티션으로 취급해요.
  • SAMPLE: 컬렉션을 샘플링해 파티션을 생성하며, 빠르지만 불균등할 수 있어요.
  • SPLIT_VECTOR: splitVector 명령을 사용해 비샤딩 컬렉션의 파티션을 생성하며, 빠르고 균등해요. splitVector 권한이 필요해요.
  • SHARDED: 파티션으로 config.chunks를 직접 읽어요(MongoDB는 샤딩된 컬렉션을 chunks로 분할하고, chunks의 범위는 컬렉션 안에 저장돼요). sharded 전략은 샤딩된 컬렉션에만 사용되며 빠르고 균등해요. config 데이터베이스의 읽기 권한이 필요해요.
  • DEFAULT: 샤딩된 컬렉션에는 sharded 전략을, 그렇지 않으면 split vector 전략을 사용해요.

MongoDB Sink

아래 예시는 sink를 구성하고 생성하는 방법을 보여줘요:

import org.apache.flink.connector.base.DeliveryGuarantee;
import org.apache.flink.connector.mongodb.sink.MongoSink;
import org.apache.flink.streaming.api.datastream.DataStream;

import com.mongodb.client.model.InsertOneModel;
import org.bson.BsonDocument;

DataStream<String> stream = ...;

MongoSink<String> sink = MongoSink.<String>builder()
        .setUri("mongodb://user:***@127.0.0.1:27017")
        .setDatabase("my_db")
        .setCollection("my_coll")
        .setBatchSize(1000)
        .setBatchIntervalMs(1000)
        .setMaxRetries(3)
        .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
        .setSerializationSchema(
                (input, context) -> new InsertOneModel<>(BsonDocument.parse(input)))
        .build();

stream.sinkTo(sink);

Configurations

Flink의 MongoDB sink는 정적 빌더 MongoSink.<InputType>builder()로 만들어져요.

  • setUri(String uri) — 필수. MongoDB의 연결 문자열을 설정해요.
  • setDatabase(String database) — 필수. 싱크할 데이터베이스의 이름.
  • setCollection(String collection) — 필수. 싱크할 컬렉션의 이름.
  • setBatchSize(int batchSize) — 선택. 기본값: 1000. 각 배치 요청에 버퍼링할 최대 작업 수를 설정해요. 배칭을 비활성화하려면 -1을 전달할 수 있어요.
  • setBatchIntervalMs(long batchIntervalMs) — 선택. 기본값: 1000. 배치 플러시 간격(밀리초)을 설정해요. 비활성화하려면 -1을 전달할 수 있어요.
  • setMaxRetries(int maxRetries) — 선택. 기본값: 3. 레코드 기록이 실패할 때 최대 재시도 횟수를 설정해요.
  • setDeliveryGuarantee(DeliveryGuarantee deliveryGuarantee) — 선택. 기본값: DeliveryGuarantee.AT_LEAST_ONCE. 원하는 DeliveryGuarantee를 설정해요. EXACTLY_ONCE 보장은 아직 지원되지 않아요.
  • setSerializationSchema(MongoSerializationSchema serializationSchema) — 필수. 입력 레코드를 MongoDB WriteModel로 파싱하려면 MongoSerializationSchema가 필요해요.

Fault Tolerance

Flink의 체크포팅이 활성화되면 Flink MongoDB Sink는 MongoDB 클러스터에 write 작업의 at-least-once 전달을 보장해요. 체크포인트 시점에 MongoWriter의 모든 대기 중인 write 작업을 기다리는 방식으로 동작해요. 이는 체크포인트가 트리거되기 전의 모든 요청이 더 많은 레코드를 sink에 보내기 전에 MongoDB로부터 성공적으로 승인되었음을 효과적으로 보장해요.

체크포인트와 내결함성에 대한 자세한 내용은 fault tolerance 문서에 있어요.

내결함성 있는 MongoDB Sink를 사용하려면 실행 환경에서 토폴로지의 체크포팅이 활성화되어야 해요:

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000); // checkpoint every 5000 msecs
val env = StreamExecutionEnvironment.getExecutionEnvironment()
env.enableCheckpointing(5000) // checkpoint every 5000 msecs
env = StreamExecutionEnvironment.get_execution_environment()
# checkpoint every 5000 msecs
env.enable_checkpointing(5000)

중요: 체크포팅은 기본적으로 활성화되지 않지만 기본 delivery guarantee는 AT_LEAST_ONCE예요. 이는 sink가 끝나거나 MongoWriter가 자동으로 플러시될 때까지 요청을 버퍼링하게 해요. 기본적으로 MongoWriter는 1000개의 write 작업이 추가된 후 플러시해요. writer가 더 자주 플러시하도록 구성하려면 MongoWriter configuration 섹션을 참고하세요.

결정적 id와 upsert 메서드를 가진 WriteModel을 사용하면 커넥터에 AT_LEAST_ONCE delivery가 구성되어 있어도 MongoDB에서 exactly-once 의미론을 달성할 수 있어요.

Configuring the Internal Mongo Writer

내부 MongoWriter는 write 작업이 플러시되는 방식을 MongoSinkBuilder의 다음 메서드로 더 구성할 수 있어요:

  • setBatchSize(int batchSize): 플러시 전에 버퍼링할 최대 write 작업 수. 비활성화하려면 -1을 전달할 수 있어요.
  • setBatchIntervalMs(long batchIntervalMs): 버퍼링된 write 작업 크기와 무관하게 플러시하는 간격. 비활성화하려면 -1을 전달할 수 있어요.

다음과 같이 설정하면 다음과 같은 쓰기 동작이 있어요:

  • 시간 간격이나 배치 크기가 한도를 초과하면 플러시. batchSize > 1이고 batchInterval > 0
  • 체크포인트에서만 플러시. batchSize == -1이고 batchInterval == -1
  • 모든 단일 write 작업마다 플러시. batchSize == 1 또는 batchInterval == 0

더 알아보기 (Learn more)