MongoDB 커넥터
MongoDB 커넥터
Flink는 at-least-once 보장으로 MongoDB 컬렉션에서 데이터를 읽고 쓰기 위한 MongoDB 커넥터를 제공해요.
본문
이 커넥터를 사용하려면 프로젝트에 다음 의존성 중 하나를 추가하세요.
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)— 필수. 입력 레코드를 MongoDBWriteModel로 파싱하려면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