Maxwell 포맷

Maxwell 포맷 (Maxwell Format)

Maxwell 는 MySQL 에서 Kafka, Kinesis 및 기타 스트리밍 커넥터로 변경 사항을 실시간으로 스트리밍할 수 있는 CDC(Changelog Data Capture) 도구입니다. Flink 는 Maxwell JSON 메시지를 Flink SQL 시스템의 INSERT/UPDATE/DELETE 메시지로 해석할 수 있습니다.

출처: 문서

본문

이 포맷은 Changelog-Data-Capture 포맷 으로, 직렬화 스키마(Serialization Schema)역직렬화 스키마(Deserialization Schema) 로 사용할 수 있습니다.

Maxwell 는 MySQL 에서 Kafka, Kinesis 및 기타 스트리밍 커넥터로 변경 사항을 실시간으로 스트리밍할 수 있는 CDC(Changelog Data Capture) 도구입니다. Maxwell 은 changelog 용 통합 포맷 스키마를 제공하며 JSON 을 사용해 메시지를 직렬화하는 것을 지원합니다.

Flink 는 Maxwell JSON 메시지를 Flink SQL 시스템의 INSERT/UPDATE/DELETE 메시지로 해석하는 것을 지원합니다. 이 기능을 활용하는 것은 여러 경우에 유용합니다. 예를 들면:

  • 데이터베이스에서 다른 시스템으로 증분 데이터를 동기화
  • 감사 로그(auditing logs)
  • 데이터베이스의 실시간 구체화 뷰(materialized views)
  • 데이터베이스 테이블의 변경 이력을 시간적 조인(temporal join) 하는 등

Flink 는 또한 Flink SQL 의 INSERT/UPDATE/DELETE 메시지를 Maxwell JSON 메시지로 인코딩해 Kafka 같은 외부 시스템으로 내보내는 것을 지원합니다. 그러나 현재 Flink 는 UPDATE_BEFORE 와 UPDATE_AFTER 를 단일 UPDATE 메시지로 결합할 수 없습니다. 따라서 Flink 는 UPDATE_BEFORE 와 UPDATE_AFTER 를 DELETE 및 INSERT Maxwell 메시지로 인코딩합니다.

의존성

Maxwell 포맷을 사용하려면 빌드 자동화 도구(예: Maven 또는 SBT)를 사용하는 프로젝트와 SQL JAR 번들을 사용하는 SQL Client 모두에 다음 의존성이 필요합니다.

Maven 의존성:

<dependency>
  <groupId>org.apache.flink</groupId>
  <artifactId>flink-json</artifactId>
  <version>2.3.0</version>
</dependency>

SQL Client: 내장(Built-in)

참고: Maxwell JSON 으로 changelog 를 Kafka 토픽에 동기화하는 방법은 Maxwell 문서 를 참고하세요.

Maxwell 포맷 사용 방법

Maxwell 은 changelog 용 통합 포맷을 제공합니다. 다음은 JSON 형식의 MySQL products 테이블에서 캡처한 update 연산의 간단한 예입니다:

{
   "database":"test",
   "table":"e",
   "type":"insert",
   "ts":1477053217,
   "xid":23396,
   "commit":true,
   "position":"master.000006:800911",
   "server_id":23042,
   "thread_id":108,
   "primary_key": [1, "2016-10-21 05:33:37.523000"],
   "primary_key_columns": ["id", "c"],
   "data":{
     "id":111,
     "name":"scooter",
     "description":"Big 2-wheel scooter",
     "weight":5.15
   },
   "old":{
     "weight":5.18,
   }
}

참고: 각 필드의 의미는 Maxwell 문서 를 참고하세요.

MySQL products 테이블은 4개의 열(id, name, description, weight)이 있습니다. 위의 JSON 메시지는 products 테이블에 대한 update 변경 이벤트로, id = 111 행의 weight 값이 5.18 에서 5.15 로 변경된 것입니다. 이 메시지가 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' = 'maxwell-json'
)

토픽을 Flink 테이블로 등록한 뒤에는 Maxwell 메시지를 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 format 의 메타데이터 필드를 노출할 수 있습니다.

데이터 타입 설명
database STRING NULL 원본 데이터베이스. Maxwell 레코드의 database 필드에 해당합니다(가능한 경우).
table STRING NULL 원본 데이터베이스 테이블. Maxwell 레코드의 table 필드에 해당합니다(가능한 경우).
primary-key-columns ARRAY<STRING> NULL 기본 키 이름의 배열. Maxwell 레코드의 primary_key_columns 필드에 해당합니다(가능한 경우).
ingestion-timestamp TIMESTAMP_LTZ(3) NULL 커넥터가 이벤트를 처리한 타임스탬프. Maxwell 레코드의 ts 필드에 해당합니다.

다음 예제는 Kafka 에서 Maxwell 메타데이터 필드에 접근하는 방법을 보여줍니다:

CREATE TABLE KafkaTable (
  origin_database STRING METADATA FROM 'value.database' VIRTUAL,
  origin_table STRING METADATA FROM 'value.table' VIRTUAL,
  origin_primary_key_columns ARRAY<STRING> METADATA FROM 'value.primary-key-columns' VIRTUAL,
  origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' 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' = 'maxwell-json'
);

포맷 옵션 (Format Options)

옵션 필수 기본값 유형 설명
format required (none) String 사용할 포맷을 지정합니다. 여기서는 'maxwell-json' 여야 합니다.
maxwell-json.ignore-parse-errors optional false Boolean 실패 대신 파싱 오류가 있는 필드와 행을 건너뜁니다. 오류 시 필드는 null 로 설정됩니다.
maxwell-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')으로 입력 타임스탬프를 파싱하고 같은 형식으로 출력합니다.
maxwell-json.map-null-key.mode optional 'FAIL' String map 데이터에 대한 null 키 직렬화 시 처리 모드를 지정합니다. 현재 지원되는 값은 'FAIL', 'DROP', 'LITERAL' 입니다: 'FAIL' 옵션은 null 키가 있는 map 을 만나면 예외를 던집니다. 'DROP' 옵션은 map 데이터의 null 키 항목을 버립니다. 'LITERAL' 옵션은 null 키를 문자열 리터럴로 대체합니다. 문자열 리터럴은 maxwell-json.map-null-key.literal 옵션으로 정의됩니다.
maxwell-json.map-null-key.literal optional 'null' String 'maxwell-json.map-null-key.mode' 가 LITERAL 일 때 null 키를 대체할 문자열 리터럴을 지정합니다.
maxwell-json.encode.decimal-as-plain-number optional false Boolean 가능한 과학적 표기법 대신 모든 소수를 일반 숫자로 인코딩합니다. 기본적으로 소수는 과학적 표기법으로 작성될 수 있습니다. 예를 들어 0.000000027 는 기본적으로 2.7E-8 로 인코딩되며, 이 옵션을 true 로 설정하면 0.000000027 로 작성됩니다.
maxwell-json.encode.ignore-null-fields optional false Boolean null 이 아닌 필드만 인코딩합니다. 기본적으로 모든 필드가 포함됩니다.

주의사항 (Caveats)

중복 변경 이벤트

Maxwell 애플리케이션은 모든 변경 이벤트를 정확히 한 번(exactly-once) 전달할 수 있습니다. Flink 는 Maxwell 이 생성한 이벤트를 이 상황에서 소비할 때 아주 잘 동작합니다. Maxwell 애플리케이션이 최소 한 번(at-least-once) 전달로 동작하면 Kafka 에 중복 변경 이벤트를 전달할 수 있고 Flink 는 중복 이벤트를 얻게 됩니다. 이는 Flink 쿼리가 잘못된 결과나 예상치 못한 예외를 얻는 원인이 될 수 있습니다. 따라서 이 상황에서는 작업 구성 table.exec.source.cdc-events-duplicatetrue 로 설정하고 소스에 PRIMARY KEY 를 정의할 것을 권장합니다. 프레임워크는 추가 상태 기반 operator 를 생성하고 기본 키를 사용해 변경 이벤트를 중복 제거해 정규화된 changelog 스트림을 생성합니다.

데이터 타입 매핑

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

더 알아보기 (Learn more)