Debezium 포맷
Debezium 포맷
Debezium은 CDC(Changelog Data Capture) 도구로, MySQL, PostgreSQL, Oracle, Microsoft SQL Server 등 여러 데이터베이스에서 변경 사항을 실시간으로 Kafka로 스트리밍할 수 있어요. Debezium은 changelog를 위한 통합 포맷 스키마를 제공하며, JSON과 Apache Avro를 사용해 메시지를 직렬화하는 것을 지원해요.
출처: 문서
본문
Flink는 Debezium JSON과 Avro 메시지를 Flink SQL 시스템의 INSERT/UPDATE/DELETE 메시지로 해석하는 것을 지원해요. 이 기능을 활용하면 많은 경우에 유용한데, 예를 들어
- 데이터베이스에서 다른 시스템으로 증분 데이터를 동기화
- 로그 감사(auditing)
- 데이터베이스의 실시간 머티어리얼라이즈드 뷰(materialized views)
- 데이터베이스 테이블의 변경 이력을 사용한 temporal join 등이 있어요.
Flink는 또한 Flink SQL의 INSERT/UPDATE/DELETE 메시지를 Debezium JSON 또는 Avro 메시지로 인코딩해 Kafka 같은 외부 시스템으로 방출하는 것을 지원해요. 그러나 현재 Flink는 UPDATE_BEFORE와 UPDATE_AFTER를 단일 UPDATE 메시지로 결합할 수 없어요. 따라서 Flink는 UPDATE_BEFORE와 UPDATE_AFTER를 각각 DELETE와 INSERT Debezium 메시지로 인코딩해요.
의존성 (Dependencies)
Debezium Confluent Avro
Debezium 포맷을 사용하려면 빌드 자동화 도구(Maven이나 SBT 같은)를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 양쪽 모두에서 다음 의존성이 필요해요.
| Maven dependency | SQL Client |
|---|---|
org.apache.flink : flink-avro-confluent-registry : 2.3.0 |
Download |
Debezium Json
Debezium 포맷을 사용하려면 빌드 자동화 도구를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 양쪽 모두에서 다음 의존성이 필요해요.
| Maven dependency | SQL Client |
|---|---|
org.apache.flink : flink-json : 2.3.0 |
Built-in |
참고: changelog를 Kafka 토픽으로 동기화하도록 Debezium Kafka Connect를 설정하는 방법은 Debezium 문서를 참조하세요.
Debezium 포맷 사용 방법 (How to use Debezium format)
Debezium은 changelog를 위한 통합 포맷을 제공해요. 다음은 MySQL products 테이블에서 캡처된 update 연산의 간단한 JSON 형식 예시예요:
{
"before": {
"id": 111,
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": 5.18
},
"after": {
"id": 111,
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": 5.15
},
"source": {...},
"op": "u",
"ts_ms": 1589362330904,
"transaction": null
}
참고: 각 필드의 의미에 대해서는 Debezium 문서를 참조하세요.
MySQL products 테이블은 4개의 컬럼(id, name, description, weight)을 가져요. 위 JSON 메시지는 id = 111인 행의 weight 값이 5.18에서 5.15로 변경된 products 테이블의 update 변경 이벤트예요.
이 메시지가 Kafka 토픽 products_binlog로 동기화된다고 가정하면, 다음 DDL(De bezium JSON 및 Debezium Confluent Avro용)을 사용해 이 토픽을 소비하고 변경 이벤트를 해석할 수 있어요.
Debezium JSON DDL
CREATE TABLE topic_products (
-- schema is totally the same to the MySQL "products" table
id BIGINT,
name STRING,
description STRING,
weight DECIMAL(10, 2)
) WITH (
'connector' = 'kafka',
'topic' = 'products_binlog',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
-- using 'debezium-json' as the format to interpret Debezium JSON messages
'format' = 'debezium-json'
)
어떤 경우 사용자는 Kafka 구성 'value.converter.schemas.enable'을 활성화한 상태로 Debezium Kafka Connect를 설정해 메시지에 스키마를 포함시킬 수 있어요. 그러면 Debezium JSON 메시지가 다음과 같이 보일 수 있어요:
{
"schema": {...},
"payload": {
"before": {
"id": 111,
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": 5.18
},
"after": {
"id": 111,
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": 5.15
},
"source": {...},
"op": "u",
"ts_ms": 1589362330904,
"transaction": null
}
}
이런 메시지를 해석하려면 위 DDL WITH 절에 'debezium-json.schema-include' = 'true' 옵션을 추가해야 해요(기본값 false). 보통 스키마를 포함하는 것은 메시지를 매우 장황하게 만들고 파싱 성능을 떨어뜨리므로 권장되지 않아요.
Debezium Confluent Avro DDL
CREATE TABLE topic_products (
-- schema is totally the same to the MySQL "products" table
id BIGINT,
name STRING,
description STRING,
weight DECIMAL(10, 2)
) WITH (
'connector' = 'kafka',
'topic' = 'products_binlog',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
-- using 'debezium-avro-confluent' as the format to interpret Debezium Avro messages
'format' = 'debezium-avro-confluent',
-- the URL to the schema registry for Kafka
'debezium-avro-confluent.url' = 'http://localhost:8081'
)
결과 생성 (Producing Results)
모든 데이터 포맷에 대해, 토픽을 Flink 테이블로 등록한 후 Debezium 메시지를 changelog 소스로 소비할 수 있어요.
-- a real-time materialized view on the MySQL "products"
-- which calculate the latest average of weight for the same products
SELECT name, AVG(weight) FROM topic_products GROUP BY name;
-- synchronize all the data and incremental changes of MySQL "products" table to
-- Elasticsearch "products" index for future searching
INSERT INTO elasticsearch_products
SELECT * FROM topic_products;
사용 가능한 메타데이터 (Available Metadata)
다음 포맷 메타데이터는 테이블 정의에서 읽기 전용(VIRTUAL) 컬럼으로 노출될 수 있어요.
주의 포맷 메타데이터 필드는 해당 커넥터가 포맷 메타데이터를 전달할 때만 사용할 수 있어요. 현재 오직 Kafka 커넥터만 value 포맷에 대한 메타데이터 필드를 노출할 수 있어요.
| 키(Key) | 데이터 타입(Data Type) | 설명(Description) |
|---|---|---|
schema |
STRING NULL |
페이로드의 스키마를 설명하는 JSON 문자열. Debezium 레코드에 스키마가 포함되지 않으면 null. |
ingestion-timestamp |
TIMESTAMP_LTZ(3) NULL |
커넥터가 이벤트를 처리한 타임스탬프. Debezium 레코드의 ts_ms 필드에 해당. |
source.timestamp |
TIMESTAMP_LTZ(3) NULL |
소스 시스템이 이벤트를 생성한 타임스탬프. Debezium 레코드의 source.ts_ms 필드에 해당. |
source.database |
STRING NULL |
원본 데이터베이스. 사용 가능하면 Debezium 레코드의 source.db 필드에 해당. |
source.schema |
STRING NULL |
원본 데이터베이스 스키마. 사용 가능하면 Debezium 레코드의 source.schema 필드에 해당. |
source.table |
STRING NULL |
원본 데이터베이스 테이블. 사용 가능하면 Debezium 레코드의 source.table 또는 source.collection 필드에 해당. |
source.properties |
MAP NULL |
다양한 소스 속성의 맵. Debezium 레코드의 source 필드에 해당. |
다음 예시는 Kafka에서 Debezium 메타데이터 필드에 접근하는 방법을 보여줘요:
CREATE TABLE KafkaTable (
origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL,
event_time TIMESTAMP(3) METADATA FROM 'value.source.timestamp' VIRTUAL,
origin_database STRING METADATA FROM 'value.source.database' VIRTUAL,
origin_schema STRING METADATA FROM 'value.source.schema' VIRTUAL,
origin_table STRING METADATA FROM 'value.source.table' VIRTUAL,
origin_properties MAP<STRING, STRING> METADATA FROM 'value.source.properties' VIRTUAL,
user_id BIGINT,
item_id BIGINT,
behavior STRING
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'value.format' = 'debezium-json'
);
포맷 옵션 (Format Options)
Flink는 Debezium이 생성한 Avro 또는 JSON 메시지를 해석하기 위해 debezium-avro-confluent와 debezium-json 포맷을 제공해요.
Debezium Avro 메시지를 해석하려면 debezium-avro-confluent 포맷을, Debezium JSON 메시지를 해석하려면 debezium-json 포맷을 사용해요.
Debezium Avro
| 옵션(Option) | 필수 | 기본값(Default) | 타입(Type) | 설명(Description) |
|---|---|---|---|---|
| format | required | (none) | String | 어떤 포맷을 사용할지 지정해요. 여기서는 'debezium-avro-confluent'를 사용해야 해요. |
| debezium-avro-confluent.basic-auth.credentials-source | optional | (none) | String | Schema Registry를 위한 기본 인증(basic auth) 크레덴셜 소스 |
| debezium-avro-confluent.basic-auth.user-info | optional | (none) | String | schema registry를 위한 기본 인증 사용자 정보 |
| debezium-avro-confluent.bearer-auth.credentials-source | optional | (none) | String | Schema Registry를 위한 Bearer 인증 크레덴셜 소스 |
| debezium-avro-confluent.bearer-auth.token | optional | (none) | String | Schema Registry를 위한 Bearer 인증 토큰 |
| debezium-avro-confluent.properties | optional | (none) | Map | 기본 Schema Registry로 전달되는 속성 맵. Flink 구성 옵션으로 공식 노출되지 않는 옵션에 유용해요. 다만 Flink 옵션이 더 높은 우선순위를 가진다는 점에 주의하세요. |
| debezium-avro-confluent.ssl.keystore.location | optional | (none) | String | SSL keystore의 위치/파일 |
| debezium-avro-confluent.ssl.keystore.password | optional | (none) | String | SSL keystore의 비밀번호 |
| debezium-avro-confluent.ssl.truststore.location | optional | (none) | String | SSL truststore의 위치/파일 |
| debezium-avro-confluent.ssl.truststore.password | optional | (none) | String | SSL truststore의 비밀번호 |
| debezium-avro-confluent.schema | optional | (none) | String | Confluent Schema Registry에 등록되거나 등록될 스키마. 스키마가 제공되지 않으면 Flink는 테이블 스키마를 avro 스키마로 변환해요. 제공된 스키마는 'before', 'after', 'op' 필드를 포함하는 nullable 레코드 타입인 Debezium 스키마와 일치해야 해요. |
| debezium-avro-confluent.subject | optional | (none) | String | 직렬화 중 이 포맷이 사용하는 스키마를 등록할 Confluent Schema Registry subject. 기본적으로 이 포맷이 value 또는 key 포맷으로 사용되면 'kafka'와 'upsert-kafka' 커넥터는 '-value' 또는 '-key'를 기본 subject 이름으로 사용해요. 그러나 다른 커넥터(예: 'filesystem')의 경우, 싱크로 사용될 때 subject 옵션이 필요해요. |
| debezium-avro-confluent.url | required | (none) | String | 스키마를 가져오고/등록할 Confluent Schema Registry의 URL. |
Debezium Json
| 옵션(Option) | 필수 | 기본값(Default) | 타입(Type) | 설명(Description) |
|---|---|---|---|---|
| format | required | (none) | String | 어떤 포맷을 사용할지 지정해요. 여기서는 'debezium-json'을 사용해야 해요. |
| debezium-json.schema-include | optional | false | Boolean | Debezium Kafka Connect를 설정할 때 사용자가 Kafka 구성 'value.converter.schemas.enable'을 활성화해 메시지에 스키마를 포함할 수 있어요. 이 옵션은 Debezium JSON 메시지가 스키마를 포함하는지 여부를 나타내요. |
| debezium-json.ignore-parse-errors | optional | false | Boolean | 실패하는 대신 파싱 오류가 있는 필드와 행을 건너뛰어요. 오류 시 필드는 null로 설정돼요. |
| debezium-json.timestamp-format.standard | optional | 'SQL' |
String | 입력 및 출력 타임스탬프 포맷을 지정해요. 현재 지원되는 값은 'SQL'과 'ISO-8601'이에요: - 옵션 'SQL'은 입력 타임스탬프를 "yyyy-MM-dd HH:mm:ss.s{precision}" 형식(예: '2020-12-30 12:13:14.123')으로 파싱하고 동일한 형식으로 출력해요. - 옵션 'ISO-8601'은 입력 타임스탬프를 "yyyy-MM-ddTHH:mm:ss.s{precision}" 형식(예: '2020-12-30T12:13:14.123')으로 파싱하고 동일한 형식으로 출력해요. |
| debezium-json.map-null-key.mode | optional | 'FAIL' |
String | map 데이터의 null 키를 직렬화할 때의 처리 모드를 지정해요. 현재 지원되는 값은 'FAIL', 'DROP', 'LITERAL'이에요: - 옵션 'FAIL'은 null 키 map을 만나면 예외를 던져요. - 옵션 'DROP'은 map 데이터의 null 키 항목을 버려요. - 옵션 'LITERAL'은 null 키를 문자열 리터럴로 대체해요. 문자열 리터럴은 debezium-json.map-null-key.literal 옵션으로 정의돼요. |
| debezium-json.map-null-key.literal | optional | 'null' | String | 'debezium-json.map-null-key.mode'가 LITERAL일 때 null 키를 대체할 문자열 리터럴을 지정해요. |
| debezium-json.encode.decimal-as-plain-number | optional | false | Boolean | 모든 decimal을 가능한 과학적 표기법 대신 일반 숫자(plain numbers)로 인코딩해요. 기본적으로 decimal은 과학적 표기법으로 작성될 수 있어요. 예를 들어 0.000000027은 기본적으로 2.7E-8로 인코딩되며, 이 옵션을 true로 설정하면 0.000000027로 작성돼요. |
| debezium-json.encode.ignore-null-fields | optional | false | Boolean | null이 아닌 필드만 인코딩해요. 기본적으로 모든 필드가 포함돼요. |
주의 사항 (Caveats)
중복 변경 이벤트 (Duplicate change events)
정상적인 운영 시나리오에서 Debezium 애플리케이션은 모든 변경 이벤트를 정확히 한 번(exactly-once) 전달해요. Flink는 이런 상황에서 Debezium이 생성한 이벤트를 소비할 때 아주 잘 동작해요.
그러나 장애 조치(failover)가 발생하면 Debezium 애플리케이션은 최소 한 번(at-least-once) 전달로 동작해요. 전달 보장에 대한 자세한 내용은 Debezium 문서를 참조하세요.
즉 비정상 상황에서 Debezium이 Kafka에 중복 변경 이벤트를 전달할 수 있고 Flink는 중복 이벤트를 받게 된다는 뜻이에요.
이것은 Flink 쿼리가 잘못된 결과나 예상치 못한 예외를 얻는 원인이 될 수 있어요. 따라서 이런 상황에서는 작업 구성 table.exec.source.cdc-events-duplicate을 true로 설정하고 소스에 PRIMARY KEY를 정의하는 것을 권장해요.
프레임워크는 추가 상태 연산자를 생성하고, 기본 키를 사용해 변경 이벤트를 중복 제거하고 정규화된 changelog 스트림을 생성해요.
Debezium Postgres 커넥터가 생성한 데이터 소비 (Consuming data produced by Debezium Postgres Connector)
Debezium Connector for PostgreSQL을 사용해 변경 사항을 Kafka로 캡처한다면, 모니터링되는 PostgreSQL 테이블의 REPLICA IDENTITY 구성이 FULL로 설정되었는지 확인하세요(기본값은 DEFAULT).
그렇지 않으면 Flink SQL은 현재 Debezium 데이터를 해석하는 데 실패할 거예요.
FULL 전략에서 UPDATE와 DELETE 이벤트는 테이블의 모든 컬럼의 이전 값을 포함해요. 다른 전략에서는 UPDATE와 DELETE 이벤트의 "before" 필드에 기본 키 컬럼만 포함되거나, 기본 키가 없으면 null이 돼요.
ALTER TABLE <table_name> REPLICA IDENTITY FULL을 실행해 REPLICA IDENTITY를 변경할 수 있어요.
자세한 내용은 PostgreSQL REPLICA IDENTITY에 대한 Debezium 문서를 참조하세요.
데이터 타입 매핑 (Data Type Mapping)
현재 Debezium 포맷은 직렬화와 역직렬화에 JSON과 Avro 포맷을 사용해요. 데이터 타입 매핑에 대한 자세한 내용은 JSON Format 문서와 Confluent Avro Format 문서를 참조하세요.