MongoDB SQL 커넥터

MongoDB SQL 커넥터 (MongoDB SQL Connector)

스캔 소스: 유한(Bounded) / 룩업 소스: 동기 모드 / 싱크: 배치(Batch) / 싱크: 스트리밍 Append & Upsert 모드

MongoDB 커넥터는 MongoDB에서 데이터를 읽고 MongoDB로 데이터를 쓸 수 있게 해줍니다. 이 문서는 MongoDB에 대해 SQL 쿼리를 실행하도록 MongoDB 커넥터를 설정하는 방법을 설명합니다.

커넥터는 DDL에 정의된 기본 키를 사용해 외부 시스템과 UPDATE/DELETE 메시지를 교환하는 upsert 모드로 동작할 수 있습니다.

DDL에 기본 키가 정의되지 않으면 커넥터는 외부 시스템과 INSERT 전용 메시지만 교환하는 append 모드로만 동작할 수 있습니다.

출처: 문서

본문

의존성 (Dependencies)

MongoDB 커넥터를 사용하려면 빌드 자동화 도구(Maven, SBT 등)를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 모두에 다음 의존성이 필요합니다.

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

MongoDB 커넥터는 바이너리 배포의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기를 참고하세요.

MongoDB 테이블 생성 방법 (How to create a MongoDB table)

MongoDB 테이블은 다음과 같이 정의할 수 있습니다:

-- register a MongoDB table 'users' in Flink SQL
CREATE TABLE MyUserTable (
  _id STRING,
  name STRING,
  age INT,
  status BOOLEAN,
  PRIMARY KEY (_id) NOT ENFORCED
) WITH (
   'connector' = 'mongodb',
   'uri' = 'mongodb://user:***@127.0.0.1:27017',
   'database' = 'my_db',
   'collection' = 'users'
);

-- write data into the MongoDB table from the other table "T"
INSERT INTO MyUserTable
SELECT _id, name, age, status FROM T;

-- scan data from the MongoDB table
SELECT id, name, age, status FROM MyUserTable;

-- temporal join the MongoDB table as a dimension table
SELECT * FROM myTopic
LEFT JOIN MyUserTable FOR SYSTEM_TIME AS OF myTopic.proctime
ON myTopic.key = MyUserTable._id;

커넥터 옵션 (Connector Options)

옵션 필수 Forwarded 기본값 타입 설명
connector required no (none) String 사용할 커넥터. 여기서는 'mongodb'여야 합니다.
uri required yes (none) String MongoDB 연결 uri.
database required yes (none) String 읽거나 쓸 MongoDB 데이터베이스 이름.
collection required yes (none) String 읽거나 쓸 MongoDB 컬렉션 이름.
scan.fetch-size optional yes 2048 Integer 읽을 때 왕복당 데이터베이스에서 가져와야 하는 문서 수에 대한 힌트를 리더에게 줍니다.
scan.cursor.no-timeout optional yes true Boolean MongoDB 서버는 보통 과도한 메모리 사용을 막기 위해 비활동 기간(10분) 후 유휴 커서에 시간 초과를 합니다. 이를 막으려면 이 옵션을 true로 설정하세요. 그러나 애플리케이션이 현재 문서 배치를 처리하는 데 30분 이상 걸리면 세션이 만료된 것으로 표시되고 닫힙니다.
scan.partition.strategy optional no default String 파티션 전략을 지정합니다. 사용 가능한 전략은 single, sample, split-vector, sharded, default입니다. 자세한 내용은 아래 Partitioned Scan 섹션을 참고하세요.
scan.partition.size optional no 64mb MemorySize 파티션 메모리 크기를 지정합니다.
scan.partition.samples optional no 10 Integer 파티션당 샘플 수를 지정합니다. 파티션 전략이 sample일 때만 적용됩니다. 샘플 파티셔너는 컬렉션을 샘플링하고 파티션 필드로 프로젝션·정렬합니다. 그런 다음 매 scan.partition.samples마다 파티션 경계 계산에 사용할 값으로 사용합니다. 총 샘플 수는 samples per partition * (count of documents / number of documents per partition)로 계산됩니다.
lookup.cache optional no NONE Enum (NONE, PARTIAL) 룩업 테이블의 캐시 전략. 현재 NONE(캐시 없음)과 PARTIAL(외부 데이터베이스에서 룩업 연산 시 항목 캐싱)을 지원합니다.
lookup.partial-cache.max-rows optional no (none) Long 룩업 캐시의 최대 행 수. 이 값을 초과하면 가장 오래된 행이 만료됩니다. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"로 설정되어야 합니다. 자세한 내용은 아래 Lookup Cache 섹션을 참고하세요.
lookup.partial-cache.expire-after-write optional no (none) Duration 캐시에 쓴 후 룩업 캐시의 각 행의 최대 수명. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.partial-cache.expire-after-access optional no (none) Duration 캐시의 항목에 접근한 후 룩업 캐시의 각 행의 최대 수명. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.partial-cache.caching-missing-key optional no true Boolean 룩업 키가 테이블의 어떤 행과도 일치하지 않을 때 빈 값을 캐시에 저장할지 여부. 이 옵션을 사용하려면 "lookup.cache"가 "PARTIAL"이어야 합니다.
lookup.max-retries optional no 3 Integer 룩업 데이터베이스가 실패한 경우 최대 재시도 횟수.
lookup.retry.interval optional no 1s Duration 데이터베이스에서 룩업 레코드를 가져오지 못한 경우 재시도 시간 간격을 지정합니다.
filter.handling.policy optional no always Enum (always, never) 필터 푸시다운을 제어하는 세밀한 구성. 지원 정책:
  • always: 지원되는 필터를 항상 MongoDB에 푸시.
  • never: 어떤 필터도 MongoDB에 푸시하지 않음. | | sink.buffer-flush.max-rows | optional | yes | 1000 | Integer | 배치 요청당 버퍼링되는 최대 행 수를 지정합니다. | | sink.buffer-flush.interval | optional | yes | 1s | Duration | 배치 플러시 간격을 지정합니다. | | sink.max-retries | optional | yes | 3 | Integer | 데이터베이스에 레코드 쓰기가 실패한 경우 최대 재시도 횟수. | | sink.retry.interval | optional | yes | 1s | Duration | 데이터베이스에 레코드 쓰기가 실패한 경우 재시도 시간 간격을 지정합니다. | | sink.parallelism | optional | no | (none) | Integer | MongoDB 싱크 연산자의 병렬도를 정의합니다. 기본적으로 병렬도는 업스트림 체인 연산자의 병렬도와 동일하게 프레임워크가 결정합니다. | | sink.delivery-guarantee | optional | no | at-least-once | Enum (none, at-least-once) | 커밋 시 선택적 전달 보장. exactly-once 보장은 아직 지원되지 않습니다. |

기능 (Features)

키 처리 (Key handling)

MongoDB 싱크는 기본 키가 정의되었는지 여부에 따라 upsert 모드 또는 append 모드로 동작할 수 있습니다. 기본 키가 정의되면 MongoDB 싱크는 UPDATE/DELETE 메시지를 포함하는 쿼리를 소비할 수 있는 upsert 모드로 동작합니다. 기본 키가 정의되지 않으면 MongoDB 싱크는 INSERT 전용 메시지를 포함하는 쿼리만 소비할 수 있는 append 모드로 동작합니다.

MongoDB에서 기본 키는 MongoDB 문서의 _id를 계산하는 데 사용됩니다. 그 값은 컬렉션에서 고유하고 불변이어야 하며, Array가 아닌 모든 BSON 타입일 수 있습니다. _id가 하위 필드를 포함하면 하위 필드 이름은 ($) 기호로 시작할 수 없습니다.

기본 키 인덱스에도 몇 가지 제약이 있습니다. MongoDB 4.2 이전에는 BSON 타입에 따라 구조적 오버헤드를 포함할 수 있는 인덱스 항목의 총 크기가 1024바이트보다 작아야 합니다. 4.2 버전부터 MongoDB는 인덱스 키 제한(Index Key Limit)을 제거합니다. 더 자세한 소개는 Index Key Limit를 참고할 수 있습니다.

MongoDB 커넥터는 DDL에 정의된 순서대로 모든 기본 키 필드를 합성해 각 행의 문서 _id를 생성합니다.

  • 지정된 기본 키에 단일 필드만 있으면 해당 필드 데이터를 해당 문서의 _id로 bson 값으로 변환합니다.
  • 지정된 기본 키에 여러 필드가 있으면 이 필드들을 해당 문서의 _id로 bson 문서로 변환·합성합니다. 예를 들어 PRIMARY KEY (f1, f2) NOT ENFORCED라는 기본 키 문이 있으면 추출된 _id는 _id: {f1: v1, f2: v2} 형태가 됩니다.

DDL에 _id 필드가 있지만 기본 키가 _id로 선언되지 않으면 모호해짐에 유의하세요. _id 컬럼을 키로 사용하거나 _id 컬럼 이름을 변경하세요.

PRIMARY KEY 문법에 대한 자세한 내용은 CREATE TABLE DDL을 참고하세요.

파티션 스캔 (Partitioned Scan)

병렬 Source 작업 인스턴스에서 데이터 읽기를 가속화하기 위해 Flink는 MongoDB 컬렉션에 대한 파티션 스캔 기능을 제공합니다. 다음 파티션 전략이 제공됩니다:

  • single: 전체 컬렉션을 단일 파티션으로 취급합니다.
  • sample: 컬렉션을 샘플링해 파티션을 생성합니다. 빠르지만 불균등할 수 있습니다.
  • split-vector: splitVector 명령을 사용해 비-샤딩 컬렉션에 대한 파티션을 생성합니다. 빠르고 균등합니다. splitVector 권한이 필요합니다.
  • sharded: config.chunks(MongoDB는 샤딩 컬렉션을 청크로 분할하며, 청크의 범위는 컬렉션 내에 저장됩니다)를 파티션으로 직접 읽습니다. sharded 전략은 샤딩 컬렉션에만 사용되며 빠르고 균등합니다. config 데이터베이스의 읽기 권한이 필요합니다.
  • default: 샤딩 컬렉션에는 sharded 전략을, 그 외에는 split vector 전략을 사용합니다.

룩업 캐시 (Lookup Cache)

MongoDB 커넥터는 시간 조인(temporal join)에서 룩업 소스(일명 차원 테이블)로 사용될 수 있습니다. 현재 동기 룩업 모드만 지원됩니다.

기본적으로 룩업 캐시는 활성화되지 않습니다. lookup.cachePARTIAL로 설정해 활성화할 수 있습니다.

룩업 캐시는 MongoDB 커넥터의 시간 조인 성능을 향상시키는 데 사용됩니다. 기본적으로 룩업 캐시가 활성화되지 않으므로 모든 요청이 외부 데이터베이스로 전송됩니다. 룩업 캐시가 활성화되면 각 프로세스(즉 TaskManager)가 캐시를 보유합니다. Flink는 먼저 캐시를 조회하고, 캐시 미스일 때만 외부 데이터베이스에 요청을 보내며, 반환된 행으로 캐시를 갱신합니다. 캐시가 최대 캐시 행 수 lookup.partial-cache.max-rows에 도달하거나 행이 lookup.partial-cache.expire-after-write 또는 lookup.partial-cache.expire-after-access가 지정한 최대 수명을 초과하면 캐시의 가장 오래된 행이 만료됩니다. 캐시된 행은 최신이 아닐 수 있으며, 사용자는 만료 옵션을 더 작은 값으로 조정해 더 신선한 데이터를 얻을 수 있지만 데이터베이스에 보내는 요청 수가 늘어날 수 있습니다. 따라서 이는 처리량과 정확성 사이의 균형입니다.

기본적으로 Flink는 기본 키에 대한 빈 쿼리 결과를 캐시합니다. lookup.partial-cache.caching-missing-key를 false로 설정해 이 동작을 전환할 수 있습니다.

멱등 쓰기 (Idempotent Writes)

MongoDB 커넥터는 DDL에 기본 키가 정의되어 있으면 삽입 쓰기 모드 db.connection.insert()가 아닌 upsert 쓰기 모드 db.connection.update(<query>, <update>, { upsert: true })를 사용합니다. 기본 키 필드를 MongoDB의 예약된 기본 키인 문서 _id로 합성합니다. upsert 모드로 행을 MongoDB에 쓰며, 이를 통해 멱등성(idempotence)을 제공합니다.

INSERT OVERWRITE 문으로 MongoDB 테이블에 쓸 때는 MongoDB에 쓰기 위해 upsert 모드 사용을 강제합니다. 따라서 MongoDB 테이블의 기본 키가 DDL에 정의되지 않으면 쓰기 연산이 거부됩니다.

실패가 있으면 Flink 작업은 복구되어 마지막 성공한 체크포인트부터 다시 처리하며, 이로 인해 복구 중 메시지가 재처리될 수 있습니다. 레코드를 재처리해야 하는 경우 제약 위반이나 중복 데이터를 피하는 데 도움이 되므로 upsert 모드를 강력히 권장합니다.

샤딩 컬렉션에서의 Upsert (Upsert on a sharded collection)

Mongo 레퍼런스가 말하듯이:

샤딩 컬렉션에서 db.collection.updateOne()을 사용하려면:

  • upsert: true를 지정하지 않으면 _id 필드의 정확한 일치를 포함하거나 단일 샤드를 대상으로 해야 합니다(필터에 샤드 키 포함 등).
  • upsert: true를 지정하면 필터가 샤드 키를 포함해야 합니다.

그러나 샤딩 컬렉션의 문서에는 샤드 키 필드가 없을 수 있습니다. 샤드 키가 없는 문서를 대상으로 하려면 _id 필드 같은 다른 필터 조건과 함께 null 동등 일치를 사용할 수 있습니다.

샤딩 컬렉션에 upsert할 때 샤드 키 값을 필터에 추가해야 합니다. 예를 들어:

db.collection.updateOne(
    {
        _id: ObjectId('<value>'),
        shardKey0: '<value>',
        shardKey1: '<value>'
    },
    { $set: { status: "D" }},
    { upsert: true }
);

Flink SQL에서 싱크 테이블을 생성할 때 샤드 키는 PARTITIONED BY 문법으로 선언해야 합니다. 샤드 키의 값은 런타임에 각 개별 레코드에서 얻어 필터에 추가됩니다.

CREATE TABLE MySinkTable (
    _id       BIGINT,
    shardKey0 STRING,
    shardKey1 STRING,
    status    STRING,
    PRIMARY KEY (_id) NOT ENFORCED
) PARTITIONED BY (shardKey0, shardKey1) WITH (
    'connector' = 'mongodb',
    'uri' = 'mongodb://user:***@127.0.0.1:27017',
    'database' = 'my_db',
    'collection' = 'users'
);

-- Insert with dynamic partition
INSERT INTO MySinkTable SELECT _id, shardKey0, shardKey1, status FROM T;

-- Insert with static partition
INSERT INTO MySinkTable PARTITION(shardKey0 = 'value0', shardKey1 = 'value1') SELECT 1, 'INIT';

-- Insert with static(shardKey0) and dynamic(shardKey1) partition
INSERT INTO MySinkTable PARTITION(shardKey0 = 'value0') SELECT 1, 'value1' 'INIT';

한계 (LIMITATION): 샤드 키 값이 MongoDB 4.2 이후에는 더 이상 불변이 아니지만, 샤드 키가 불변하도록 보장하는 것이 필요합니다.

Flink SQL upsert 모드로 샤딩 컬렉션에 쓸 때 갱신된 샤드 키 값만 얻을 수 있고 필터에 원래 샤드 키 값을 제공할 수 없어 중복 레코드 오류가 발생할 수 있습니다.

필터 푸시다운 (Filters Pushdown)

MongoDB는 간단한 비교와 논리 필터를 푸시다운해 쿼리를 최적화하는 것을 지원합니다. Flink SQL 필터에서 MongoDB 쿼리 연산자로의 매핑은 다음 표에 나열됩니다.

Flink SQL 필터 MongoDB 쿼리 연산자
= $eq
<> $ne
> $gt
>= $gte
< $lt
<= $lte
IS NULL $eq : null
IS NOT NULL $ne : null
OR $or
AND $and

데이터 타입 매핑 (Data Type Mapping)

MongoDB BSON 타입에서 Flink SQL 데이터 타입으로의 필드 데이터 타입 매핑은 다음 표에 나열됩니다.

MongoDB BSON 타입 Flink SQL 타입
ObjectId STRING
String STRING
Boolean BOOLEAN
Binary BINARY, VARBINARY
Int32 INTEGER
- TINYINT, SMALLINT, FLOAT
Int64 BIGINT
Double DOUBLE
Decimal128 DECIMAL
DateTime TIMESTAMP_LTZ(3)
Timestamp TIMESTAMP_LTZ(0)
Object ROW
Array ARRAY

MongoDB의 특정 타입에 대해서는 Extended JSON 형식을 사용해 Flink SQL STRING 타입으로 매핑합니다.

MongoDB BSON 타입 Flink SQL STRING
Symbol {"_value": {"$symbol": "12"}}
RegularExpression {"_value": {"$regularExpression": {"pattern": "^9$", "options": "i"}}}
JavaScript {"_value": {"$code": "function() { return 10; }"}}
DbPointer {"_value": {"$dbPointer": {"$ref": "db.coll", "$id": {"$oid": "63932a00da01604af329e33c"}}}}

더 알아보기 (Learn more)