Ogg 포맷

Ogg 포맷

Oracle GoldenGate(일명 ogg)은 실시간 데이터 메시 플랫폼을 제공하는 관리형 서비스입니다. 복제를 사용해 데이터를 고가용성으로 유지하고 실시간 분석을 가능하게 합니다. 고객은 컴퓨팅 환경을 할당하거나 관리할 필요 없이 데이터 복제와 스트림 데이터 처리 솔루션을 설계, 실행, 모니터링할 수 있습니다. Ogg는 changelog를 위한 포맷 스키마를 제공하며 JSON으로 메시지를 직렬화하는 것을 지원합니다.

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

  • 데이터베이스에서 다른 시스템으로 증분 데이터 동기화
  • 감사 로그
  • 데이터베이스의 실시간 materialized view
  • 데이터베이스 테이블 변경 이력의 temporal join 등

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

출처: 문서

본문

의존성

Ogg Json

Ogg를 사용하려면 빌드 자동화 도구(예: 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를 Kafka 토픽으로 동기화하기 위해 Ogg Kafka handler를 설정하는 방법은 Ogg Kafka Handler documentation을 참조하세요.

Ogg 포맷 사용 방법

Ogg는 changelog를 위한 통합 포맷을 제공합니다. JSON 포맷으로 Oracle PRODUCTS 테이블에서 캡처된 update 작업의 간단한 예시입니다.

{
  "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
  },
  "op_type": "U",
  "op_ts": "2020-05-13 15:40:06.000000",
  "current_ts": "2020-05-13 15:40:07.000000",
  "primary_keys": [
    "id"
  ],
  "pos": "00000000000000000000143",
  "table": "PRODUCTS"
}

참고: 각 필드의 의미는 Debezium documentation을 참조하세요.

Oracle PRODUCTS 테이블은 4개의 열(id, name, description, weight)을 가집니다. 위 JSON 메시지는 id = 111 행의 weight 값이 5.18에서 5.15로 변경된 PRODUCTS 테이블의 update 변경 이벤트입니다. 이 메시지가 Kafka 토픽 products_ogg에 동기화되었다고 가정하면, 다음 DDL로 이 토픽을 소비하고 변경 이벤트를 해석할 수 있습니다.

CREATE TABLE topic_products (
  -- schema is totally the same to the Oracle "products" table
  id BIGINT,
  name STRING,
  description STRING,
  weight DECIMAL(10, 2)
) WITH (
  'connector' = 'kafka',
  'topic' = 'products_ogg',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id' = 'testGroup',
  'format' = 'ogg-json'
)

토픽을 Flink 테이블로 등록한 후 Ogg 메시지를 changelog 소스로 소비할 수 있습니다.

-- a real-time materialized view on the Oracle "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 Oracle "PRODUCTS" table to
-- Elasticsearch "products" index for future searching
INSERT INTO elasticsearch_products
SELECT *
FROM topic_products;

사용 가능한 메타데이터

다음 포맷 메타데이터는 테이블 정의에서 읽기 전용(VIRTUAL) 열로 노출될 수 있습니다.

주의: 포맷 메타데이터 필드는 해당 커넥터가 포맷 메타데이터를 전달하는 경우에만 사용할 수 있습니다. 현재는 Kafka 커넥터만 value 포맷의 메타데이터 필드를 노출할 수 있습니다.

Key Data Type Description
table STRING NULL 정규화된 테이블 이름을 포함합니다. 정규화된 테이블 이름의 형식은 다음과 같습니다: CATALOG NAME.SCHEMA NAME.TABLE NAME
primary-keys ARRAY<STRING> NULL 원본 테이블의 기본 키 열 이름을 보유하는 배열 변수입니다. primary-keys 필드는 includePrimaryKeys 구성 속성이 true로 설정된 경우에만 JSON 출력에 포함됩니다.
ingestion-timestamp TIMESTAMP_LTZ(6) NULL 커넥터가 이벤트를 처리한 타임스탬프입니다. Ogg 레코드의 current_ts 필드에 해당합니다.
event-timestamp TIMESTAMP_LTZ(6) NULL 소스 시스템이 이벤트를 생성한 타임스탬프입니다. Ogg 레코드의 op_ts 필드에 해당합니다.

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

CREATE TABLE KafkaTable (
  origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL,
  event_time TIMESTAMP(3) METADATA FROM 'value.event-timestamp' VIRTUAL,
  origin_table STRING METADATA FROM 'value.table' VIRTUAL,
  primary_keys ARRAY<STRING> METADATA FROM 'value.primary-keys' 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' = 'ogg-json'
);

포맷 옵션

Option Required Default Type Description
format required (none) String 사용할 포맷을 지정합니다. 여기서는 'ogg-json'이어야 합니다.
ogg-json.ignore-parse-errors optional false Boolean 실패하는 대신 파싱 오류가 있는 필드와 행을 건너뜁니다. 오류가 있는 경우 필드는 null로 설정됩니다.
ogg-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' 그리고 출력 타임스탬프도 같은 형식입니다.
ogg-json.map-null-key.mode optional 'FAIL' String map 데이터의 null 키를 직렬화할 때 처리 모드를 지정합니다. 현재 지원되는 값은 'FAIL', 'DROP', 'LITERAL'입니다. - 'FAIL' 옵션은 null 키를 가진 map을 만나면 예외를 던집니다. - 'DROP' 옵션은 map 데이터의 null 키 항목을 버립니다. - 'LITERAL' 옵션은 null 키를 문자열 리터럴로 대체합니다. 문자열 리터럴은 ogg-json.map-null-key.literal 옵션으로 정의됩니다.
ogg-json.map-null-key.literal optional 'null' String 'ogg-json.map-null-key.mode'가 LITERAL일 때 null 키를 대체할 문자열 리터럴을 지정합니다.
ogg-json.encode.ignore-null-fields optional false Boolean null이 아닌 필드만 인코딩합니다. 기본적으로 모든 필드가 포함됩니다.

데이터 타입 매핑

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

더 알아보기 (Learn more)