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-duplicate 을 true 로 설정하고 소스에 PRIMARY KEY 를 정의할 것을 권장합니다. 프레임워크는 추가 상태 기반 operator 를 생성하고 기본 키를 사용해 변경 이벤트를 중복 제거해 정규화된 changelog 스트림을 생성합니다.
데이터 타입 매핑
현재 Maxwell 포맷은 직렬화와 역직렬화에 JSON 을 사용합니다. 데이터 타입 매핑에 대한 자세한 내용은 JSON Format 문서 를 참고하세요.