Upsert Kafka SQL 커넥터

Upsert Kafka SQL 커넥터

Upsert Kafka 커넥터를 사용하면 upsert 방식으로 Kafka 토픽에서 데이터를 읽고 쓸 수 있어요. scan source는 Unbounded, sink는 Streaming Upsert Mode로 동작합니다.

출처: Upsert Kafka SQL Connector

본문

Upsert Kafka 커넥터는 upsert 방식으로 Kafka 토픽에서 데이터를 읽고 쓰는 것을 허용해요.

Source로서 upsert-kafka 커넥터는 changelog 스트림을 생성하는데, 각 데이터 레코드는 update 또는 delete 이벤트를 나타내요. 더 정확히는, 데이터 레코드의 value는 같은 키의 마지막 값에 대한 UPDATE로 해석돼요(해당 키가 아직 없으면 그 update는 INSERT로 간주돼요). 테이블 비유로 보면 changelog 스트림의 데이터 레코드는 UPSERT 즉 INSERT/UPDATE로 해석되는데, 같은 키를 가진 기존 행이 덮어써지기 때문이에요. 또한 null 값은 특별한 방식으로 해석돼요: null value를 가진 레코드는 "DELETE"를 나타내요.

Sink로서 upsert-kafka 커넥터는 changelog 스트림을 소비할 수 있어요. INSERT/UPDATE_AFTER 데이터는 일반 Kafka 메시지 value로 쓰고, DELETE 데이터는 null value를 가진 Kafka 메시지(키에 대한 tombstone 표시)로 써요. Flink는 기본 키 컬럼의 값을 기준으로 데이터를 파티셔닝해 기본 키에 대한 메시지 순서를 보장하므로, 같은 키의 update/delete 메시지는 같은 파티션에 속해요.

Dependencies

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

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

Full Example

아래 예시는 Upsert Kafka 테이블을 만들고 사용하는 방법을 보여줘요:

CREATE TABLE pageviews_per_region (
  user_region STRING,
  pv BIGINT,
  uv BIGINT,
  PRIMARY KEY (user_region) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'pageviews_per_region',
  'properties.bootstrap.servers' = '...',
  'key.format' = 'avro',
  'value.format' = 'avro'
);

CREATE TABLE pageviews (
  user_id BIGINT,
  page_id BIGINT,
  viewtime TIMESTAMP,
  user_region STRING,
  WATERMARK FOR viewtime AS viewtime - INTERVAL '2' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'pageviews',
  'properties.bootstrap.servers' = '...',
  'format' = 'json'
);

-- calculate the pv, uv and insert into the upsert-kafka sink
INSERT INTO pageviews_per_region
SELECT
  user_region,
  COUNT(*),
  COUNT(DISTINCT user_id)
FROM pageviews
GROUP BY user_region;

주의: DDL에 primary key를 정의해야 해요.

Available Metadata

사용 가능한 모든 메타데이터 필드 목록은 Kafka 커넥터를 참고하세요.

Connector Options

Option Required Default Type Description
connector required (none) String 사용할 커넥터 지정. Upsert Kafka는 'upsert-kafka' 사용.
topic required (none) String 테이블이 source로 사용될 때 읽을 토픽 이름, sink로 사용될 때 쓸 토픽. source의 토픽 목록도 세미콜론('topic-1;topic-2')으로 지원. "topic-pattern"과 "topic" 중 하나만 지정 가능. sink의 경우 토픽 이름은 데이터를 쓸 토픽이에요. sink의 토픽 목록도 지원돼요. 제공된 topic-list는 topic 메타데이터 컬럼의 유효한 값에 대한 허용 목록(allow list)으로 취급돼요. 목록이 제공되면 sink 테이블의 경우 'topic' 메타데이터 컬럼은 쓰기 가능하며 지정해야 해요.
properties.bootstrap.servers required (none) String 쉼표로 구분된 Kafka 브로커 목록.
properties.* optional (none) String 임의의 Kafka 구성을 설정·전달할 수 있어요. 접미사 이름은 Kafka Configuration 문서에 정의된 구성 키와 일치해야 해요. Flink는 "properties." 키 접두사를 제거하고 변환된 키·값을 기본 KafkaClient에 전달해요. 예: 'properties.allow.auto.create.topics' = 'false'로 자동 토픽 생성을 비활성화. 그러나 Flink가 덮어쓰기 때문에 설정할 수 없는 구성도 있어요. 예: 'auto.offset.reset'.
key.format required (none) String Kafka 메시지의 key 부분을 역직렬화·직렬화하는 데 사용하는 포맷. 자세한 내용은 formats 페이지 참고. 주의: 일반 Kafka 커넥터와 달리 key 필드는 PRIMARY KEY 문법으로 지정돼요.
key.fields-prefix optional (none) String value 포맷의 필드와 이름 충돌을 피하기 위해 key 포맷의 모든 필드에 커스텀 접두사를 정의. 기본적으로 접두사는 비어 있음. 커스텀 접두사가 정의되면 테이블 스키마와 'key.fields' 모두 접두사가 붙은 이름으로 작동해요. key 포맷의 데이터 타입을 구성할 때 접두사가 제거되고 비접두사 이름이 key 포맷 안에서 사용돼요. 이 옵션은 'value.fields-include'가 'EXCEPT_KEY'로 설정되어야 함.
value.format required (none) String Kafka 메시지의 value 부분을 역직렬화·직렬화하는 데 사용하는 포맷. 자세한 내용은 formats 페이지 참고.
value.fields-include optional ALL Enum Possible values: [ALL, EXCEPT_KEY] value 포맷의 데이터 타입에서 key 컬럼을 다루는 전략을 정의. 기본적으로 테이블 스키마의 'ALL' 물리적 컬럼이 value 포맷에 포함되어, key 컬럼이 key와 value 포맷 둘 다의 데이터 타입에 나타나게 돼요.
scan.parallelism optional no (none) Integer
sink.parallelism optional (none) Integer upsert-kafka sink 연산자의 병렬도를 정의. 기본적으로 프레임워크가 상류 체인 연산자의 병렬도를 사용해 결정.
sink.buffer-flush.max-rows optional 0 Integer 플러시 전에 버퍼링된 레코드의 최대 크기. sink가 같은 키에 많은 update를 받으면 버퍼는 같은 키의 마지막 레코드를 유지해요. 이는 데이터 셔플을 줄이고 Kafka 토픽에 가능한 tombstone 메시지를 피하는 데 도움. '0'으로 설정해 비활성화 가능. 기본적으로 비활성화. 참고: 'sink.buffer-flush.max-rows'와 'sink.buffer-flush.interval' 둘 다 0보다 크게 설정해야 sink buffer flushing이 활성화돼요.
sink.buffer-flush.interval optional 0 Duration 플러시 간격(밀리초). 이 시간이 지나면 비동기 스레드가 데이터를 플러시해요. sink가 같은 키에 많은 update를 받으면 버퍼는 같은 키의 마지막 레코드를 유지해요. '0'으로 설정해 비활성화 가능. 기본적으로 비활성화. 참고: 'sink.buffer-flush.max-rows'와 'sink.buffer-flush.interval' 둘 다 0보다 크게 설정해야 활성화돼요.
sink.delivery-guarantee optional no at-least-once String
sink.transactional-id-prefix optional yes (none) String

Features

Key and Value Formats

key와 value 포맷에 대한 자세한 설명은 Kafka 커넥터를 참고하세요. 다만 이 커넥터는 key 필드가 PRIMARY KEY 제약 조건에서 파생되는 key와 value 포맷을 모두 요구한다는 점에 유의하세요.

다음 예시는 key와 value 포맷을 지정·구성하는 방법을 보여줘요. 포맷 옵션은 'key' 또는 'value'에 포맷 식별자를 붙인 접두사로 시작해요.

CREATE TABLE KafkaTable (
  `ts` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp',
  `user_id` BIGINT,
  `item_id` BIGINT,
  `behavior` STRING,
  PRIMARY KEY (`user_id`) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  ...

  'key.format' = 'json',
  'key.json.ignore-parse-errors' = 'true',

  'value.format' = 'json',
  'value.json.fail-on-missing-field' = 'false',
  'value.fields-include' = 'EXCEPT_KEY'
)

Primary Key Constraints

Upsert Kafka는 항상 upsert 방식으로 작동하며 DDL에서 primary key를 정의해야 해요. 같은 키를 가진 레코드가 같은 파티션에 정렬되어야 한다는 가정 아래, changelog source의 primary key 의미론은 materialized changelog가 primary key에 대해 고유함을 의미해요. primary key 정의는 Kafka의 key에 어떤 필드가 들어갈지도 제어해요.

Consistency Guarantees

기본적으로 Upsert Kafka sink는 체크포팅이 활성화된 상태에서 쿼리가 실행되면 at-least-once 보장으로 Kafka 토픽에 데이터를 넣어요.

즉 Flink는 같은 키를 가진 중복 레코드를 Kafka 토픽에 쓸 수 있어요. 그러나 커넥터가 upsert 모드로 작동하므로 source로 다시 읽을 때 같은 키의 마지막 레코드가 적용돼요. 따라서 upsert-kafka 커넥터는 HBase sink처럼 멱등(idempotent) 쓰기를 달성해요.

Flink의 체크포팅이 활성화되면 upsert-kafka 커넥터는 exactly-once delivery 보장을 제공할 수 있어요.

Flink의 체크포팅 활성화 외에도 적절한 sink.delivery-guarantee 옵션을 전달해 세 가지 다른 작동 모드를 선택할 수 있어요:

  • none: Flink가 어떤 것도 보장하지 않아요. 생성된 레코드는 손실되거나 중복될 수 있어요.
  • at-least-once (기본 설정): 어떤 레코드도 손실되지 않음을 보장해요(중복될 수는 있음).
  • exactly-once: Kafka 트랜잭션이 exactly-once 의미론을 제공하는 데 사용돼요. 트랜잭션으로 Kafka에 쓸 때는 Kafka에서 레코드를 소비하는 애플리케이션에 대해 원하는 isolation.level(read_uncommitted 또는 read_committed — 후자가 기본값)을 설정하는 것을 잊지 마세요.

delivery guarantee에 대한 더 많은 주의사항은 Kafka 커넥터 문서를 참고하세요.

Source Per-Partition Watermarks

Flink는 Upsert Kafka에 대한 per-partition 워터마크 생성을 지원해요. 워터마크는 Kafka consumer 안에서 생성돼요. per-partition 워터마크는 streaming shuffle 동안 워터마크가 병합되는 것과 같은 방식으로 병합돼요. source의 출력 워터마크는 읽는 파티션들 중 최소 워터마크로 결정돼요. 토픽의 일부 파티션이 idle이면 워터마크 생성기가 진행하지 않아요. 테이블 구성에서 'table.exec.source.idle-timeout' 옵션을 설정해 이 문제를 완화할 수 있어요.

자세한 내용은 Kafka watermark 전략을 참고하세요.

Data Type Mapping

Upsert Kafka는 메시지 key와 value를 바이트로 저장하므로 스키마나 데이터 타입이 없어요. 메시지는 csv, json, avro 같은 포맷으로 직렬화·역직렬화돼요. 따라서 데이터 타입 매핑은 특정 포맷에 의해 결정돼요. 자세한 내용은 Formats 페이지를 참고하세요.

더 알아보기 (Learn more)