Canal 포맷
Canal 포맷
Canal은 MySQL의 변경 사항을 실시간으로 다른 시스템으로 스트리밍할 수 있는 CDC(Changelog Data Capture) 도구입니다. Canal은 changelog를 위한 통합 포맷 스키마를 제공하며 JSON과 protobuf로 메시지를 직렬화하는 것을 지원합니다(protobuf는 Canal의 기본 포맷입니다).
Flink는 Canal JSON 메시지를 Flink SQL 시스템의 INSERT/UPDATE/DELETE 메시지로 해석하는 것을 지원합니다. 이 기능은 다음과 같은 여러 경우에 유용합니다.
- 데이터베이스에서 다른 시스템으로 증분 데이터 동기화
- 감사 로그
- 데이터베이스의 실시간 materialized view
- 데이터베이스 테이블 변경 이력의 temporal join 등
Flink는 또한 Flink SQL의 INSERT/UPDATE/DELETE 메시지를 Canal JSON 메시지로 인코딩하여 Kafka 같은 저장소로 내보내는 것을 지원합니다. 다만 현재 Flink는 UPDATE_BEFORE와 UPDATE_AFTER를 단일 UPDATE 메시지로 결합할 수 없습니다. 따라서 Flink는 UPDATE_BEFORE와 UPDATE_AFTER를 DELETE와 INSERT Canal 메시지로 인코딩합니다.
참고: Canal protobuf 메시지 해석 지원은 로드맵에 있습니다.
출처: 문서
본문
의존성
Canal 포맷을 사용하려면 빌드 자동화 도구(예: Maven 또는 SBT)를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 양쪽 모두에서 다음 의존성이 필요합니다.
| Maven dependency | SQL Client |
|---|---|
xml <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-json</artifactId> <version>2.3.0</version> </dependency> |
Built-in |
참고: changelog를 메시지 큐로 동기화하기 위해 Canal을 배포하는 방법은 Canal documentation을 참조하세요.
Canal 포맷 사용 방법
Canal은 changelog를 위한 통합 포맷을 제공합니다. MySQL products 테이블에서 캡처된 update 작업의 간단한 예시입니다.
{
"data": [
{
"id": "111",
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": "5.18"
}
],
"database": "inventory",
"es": 1589373560000,
"id": 9,
"isDdl": false,
"mysqlType": {
"id": "INTEGER",
"name": "VARCHAR(255)",
"description": "VARCHAR(512)",
"weight": "FLOAT"
},
"old": [
{
"weight": "5.15"
}
],
"pkNames": [
"id"
],
"sql": "",
"sqlType": {
"id": 4,
"name": 12,
"description": 12,
"weight": 7
},
"table": "products",
"ts": 1589373560798,
"type": "UPDATE"
}
참고: 각 필드의 의미는 Canal documentation을 참조하세요.
MySQL products 테이블은 4개의 열(id, name, description, weight)을 가집니다. 위 JSON 메시지는 id = 111 행의 weight 값이 5.18에서 5.15로 변경된 products 테이블의 update 변경 이벤트입니다. 메시지가 Kafka 토픽 products_binlog에 동기화되었다고 가정하면, 다음 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',
'format' = 'canal-json' -- using canal-json as the format
)
토픽을 Flink 테이블로 등록한 후 Canal 메시지를 changelog 소스로 소비할 수 있습니다.
-- a real-time materialized view on the MySQL "products"
-- which calculates 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;
사용 가능한 메타데이터
다음 포맷 메타데이터는 테이블 정의에서 읽기 전용(VIRTUAL) 열로 노출될 수 있습니다.
포맷 메타데이터 필드는 해당 커넥터가 포맷 메타데이터를 전달하는 경우에만 사용할 수 있습니다. 현재는 Kafka 커넥터만 value 포맷의 메타데이터 필드를 노출할 수 있습니다.
| Key | Data Type | Description |
|---|---|---|
database |
STRING NULL |
원본 데이터베이스입니다. 사용 가능하면 Canal 레코드의 database 필드에 해당합니다. |
table |
STRING NULL |
원본 데이터베이스 테이블입니다. 사용 가능하면 Canal 레코드의 table 필드에 해당합니다. |
sql-type |
MAP<STRING, INT> NULL |
다양한 sql type의 맵입니다. 사용 가능하면 Canal 레코드의 sqlType 필드에 해당합니다. |
pk-names |
ARRAY<STRING> NULL |
기본 키 이름의 배열입니다. 사용 가능하면 Canal 레코드의 pkNames 필드에 해당합니다. |
ingestion-timestamp |
TIMESTAMP_LTZ(3) NULL |
커넥터가 이벤트를 처리한 타임스탬프입니다. Canal 레코드의 ts 필드에 해당합니다. |
event-timestamp |
TIMESTAMP(3) WITH LOCAL TIME ZONE NULL |
MySQL 서버에서 해당 변경이 실행된 타임스탬프입니다. Canal 레코드의 es 필드에 해당합니다. |
다음 예시는 Kafka에서 Canal 메타데이터 필드에 접근하는 방법을 보여줍니다.
CREATE TABLE KafkaTable (
origin_database STRING METADATA FROM 'value.database' VIRTUAL,
origin_table STRING METADATA FROM 'value.table' VIRTUAL,
origin_sql_type MAP<STRING, INT> METADATA FROM 'value.sql-type' VIRTUAL,
origin_pk_names ARRAY<STRING> METADATA FROM 'value.pk-names' VIRTUAL,
origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL,
origin_es TIMESTAMP(3) METADATA FROM 'value.event-timestamp' VIRTUAL,
user_id BIGINT,
item_id BIGINT,
behavior STRING,
WATERMARK FOR origin_es AS origin_es - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'value.format' = 'canal-json'
);
포맷 옵션
| Option | Required | Default | Type | Description |
|---|---|---|---|---|
| format | required | (none) | String | 사용할 포맷을 지정합니다. 여기서는 'canal-json'이어야 합니다. |
| canal-json.ignore-parse-errors | optional | false | Boolean | 실패하는 대신 파싱 오류가 있는 필드와 행을 건너뜁니다. 오류가 있는 경우 필드는 null로 설정됩니다. |
| canal-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' 그리고 출력 타임스탬프도 같은 형식입니다. |
| canal-json.map-null-key.mode | optional | 'FAIL' |
String | map 데이터의 null 키를 직렬화할 때 처리 모드를 지정합니다. 현재 지원되는 값은 'FAIL', 'DROP', 'LITERAL'입니다. - 'FAIL' 옵션은 null 키를 가진 map value를 만나면 예외를 던집니다. - 'DROP' 옵션은 map 데이터의 null 키 항목을 버립니다. - 'LITERAL' 옵션은 null 키를 문자열 리터럴로 대체합니다. 문자열 리터럴은 canal-json.map-null-key.literal 옵션으로 정의됩니다. |
| canal-json.map-null-key.literal | optional | 'null' | String | 'canal-json.map-null-key.mode'가 LITERAL일 때 null 키를 대체할 문자열 리터럴을 지정합니다. |
| canal-json.encode.decimal-as-plain-number | optional | false | Boolean | 가능한 과학적 표기법 대신 모든 decimal을 일반 숫자로 인코딩합니다. 기본적으로 decimal은 과학적 표기법으로 쓰일 수 있습니다. 예: 0.000000027은 기본적으로 2.7E-8로 인코딩되며, 이 옵션을 true로 설정하면 0.000000027으로 쓰입니다. |
| canal-json.database.include | optional | (none) | String | Canal 레코드의 "database" 메타 필드를 정규식 매칭하여 특정 데이터베이스의 changelog 행만 읽는 선택적 정규식입니다. 패턴 문자열은 Java의 Pattern과 호환됩니다. |
| canal-json.table.include | optional | (none) | String | Canal 레코드의 "table" 메타 필드를 정규식 매칭하여 특정 테이블의 changelog 행만 읽는 선택적 정규식입니다. 패턴 문자열은 Java의 Pattern과 호환됩니다. |
주의 사항
중복 변경 이벤트
정상 운영 시나리오에서 Canal 애플리케이션은 모든 변경 이벤트를 exactly-once로 전달합니다. 이 상황에서 Flink는 Canal이 생성한 이벤트를 소비하는 데 매우 잘 동작합니다. 그러나 어떤 장애 조치가 발생하면 Canal 애플리케이션은 at-least-once 전달로 동작합니다. 즉, 비정상 상황에서 Canal은 중복 변경 이벤트를 메시지 큐에 전달할 수 있고 Flink는 중복 이벤트를 받게 됩니다. 이로 인해 Flink 쿼리가 잘못된 결과를 얻거나 예상치 못한 예외가 발생할 수 있습니다. 따라서 이 상황에서는 작업 구성 table.exec.source.cdc-events-duplicate를 true로 설정하고 소스에 PRIMARY KEY를 정의하는 것이 좋습니다. 프레임워크는 추가 상태 기반 연산자를 생성하고 기본 키를 사용해 변경 이벤트를 중복 제거한 후 정규화된 changelog 스트림을 생성합니다.
데이터 타입 매핑
현재 Canal 포맷은 직렬화와 역직렬화에 JSON 포맷을 사용합니다. 데이터 타입 매핑에 대한 자세한 내용은 JSON 포맷 문서를 참조하세요.