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-duplicatetrue로 설정하고 소스에 PRIMARY KEY를 정의하는 것이 좋습니다. 프레임워크는 추가 상태 기반 연산자를 생성하고 기본 키를 사용해 변경 이벤트를 중복 제거한 후 정규화된 changelog 스트림을 생성합니다.

데이터 타입 매핑

현재 Canal 포맷은 직렬화와 역직렬화에 JSON 포맷을 사용합니다. 데이터 타입 매핑에 대한 자세한 내용은 JSON 포맷 문서를 참조하세요.

더 알아보기 (Learn more)