데보션앱 소개페이지 바로가기
로그인 선택

신고하기

CLOSE
신고사유 (대표 사유 1개)
상세내용 (선택)
0/200
  • 신고한 게시글은 더 이상 보이지 않습니다.
  • 이용약관과 운영정책에 따라 신고사유에 해당하는지 검토 후 조치됩니다.
  • 허위 신고인 경우, 신고자의 서비스 이용이 제한될 수 있으니 유의하시어 신중하게 신고해 주세요.
(이 회원이 작성한 모든 댓글과 커뮤니티 게시물이 보이지 않고, 알림도 오지 않습니다.)

미리보기

커뮤니티

      1,234

      badge 23.06.15

      글 등록

      카테고리를 선택해주세요.

      DEVOTEE를 활성화 시키면
      지금 작성한 커뮤니티 글에 대해 1개의 댓글을 달아줍니다.

      버튼을 누르면 글 수정 시 ChatGPT가 작성한 댓글이 수정됩니다.

      임시저장함에 저장되었습니다. 저장일시 : 2022.5.17 14:29:08

      임시저장함

      제목을 선택하시면 이어서 작성이 가능하며,
      최대 20건까지 저장합니다.
      컨텐츠 유형, 제목, 저장일시, 삭제로 이뤄진 임시저장 목록
      컨텐츠 유형 제목 저장일 삭제

      데보션 블로그 게재 요청

      CLOSE
      • *
      • *

      본인인증

      효율적인 데보션 서비스 이용 및
      고객님의 소중한 개인정보보호를 위해
      본인인증을 진행해주세요. 본인인증 미 진행 시 로그인이 제한됩니다.
      본인인증 실패

      본인인증 로그인에 실패하였습니다.
      회원이 아니시거나 본인인증 등록이
      완료되지 않은 사용자입니다.

      회원정보 연결

      ksqlDB를 활용한 CDC (Change Data Capture) 효과 내기

      jseung21 23.05.12
      1,871 15 1

      개요

      먼저 ksqlDB의 스트림(Stream)과 테이블(Table) 객체에 대한 간단한 설명입니다

      스트림과 테이블은 모두 연속 데이터가 있는 Kafka topic의 래퍼입니다.


      스트림은 세상에서 발생하는 이벤트를 캡처한 데이터를 나타내며 다음과 같은 특징이 있습니다.

      • 무제한 : 데이터의 끝없는 연속 흐름을 저장하므로 스트림은 제한이 없습니다.

      • Immutable : 들어오는 모든 새 데이터는 현재 스트림에 추가하고 기존 레코드를 수정하지 않으며, 변경할 수 없습니다.

      테이블은 키를 사용하여 데이터 저장 또는 해당 이벤트 스트림의 구체화된 뷰를 나타내며 다음과 같은 특징을 가지고 있습니다.

      • 제한 : 스트림의 스냅샷을 나타냅니다.

      • 변경 가능 : 테이블에 동일한 키를 가진 데이터가 들어오는 경우 새 데이터(<Key, Value> 쌍)를 테이블에 추가됩니다. 동일한 키의 데이터가 존재하면 해당 키에 대한 최신 값으로 변경됩니다.

      참고: Kafka 주제의 모든 레코드는 키와 값의 쌍으로 표시됩니다. 따라서 테이블에는 항상 주어진 키의 최신값이 표시됩니다.


      스트림과 테이블에 대한 설명 예제

      거래 및 계정잔액을 활용하여 자금 이체를 관리하려는 은행 시스템을 가정합니다.

      Alice와 Bob 두 사용자의 초기 잔액이 각각 200$와 100$라고 가정합니다.

      다음은 이 두 사용자 사이에서 발생하는 일련의 자금이체 트랜잭션입니다.


      거래 1: 앨리스는 밥에게 100$를 줍니다.

      거래 2: 밥은 앨리스에게 50$를 줍니다.

      거래 3: 밥은 앨리스에게 100$를 줍니다.

      stream_table.jpg

      여기서 스트림은 한 계정에서 다른 계정으로 이체된 돈을 기록하는 거래 이벤트의 변경 불가능한 데이터를 나타냅니다.

      반면에 테이블은 사용자당 계정의 최신 상태를 반영합니다. 예를 들어 테이블은 현재 잔액을 저장합니다.


      테이블은 계정의 최신 상태를 저장하고, 스트림은 트랜잭션 레코드를 저장합니다.


      본론 - CDC (Change Data Capture) 연계 방법

      이제 본격적으로 원천 Database의 테이블로 부터 ksqlDB의 Table 객체로 CDC 연계하는 방법을 살펴 보겠습니다.

      전체 흐름은 아래와 같습니다.

      원천 DB (Oracle)의 Table A를 ksqlDB의 Table 객체로 CDC 로 연계하려고 합니다.

      map.jpg

      1. Kafka Connector로 원천 Database (여기선 Oracle)로 해당 테이블에 대하여 연결
      curl --location --request POST 'http://localhost:8083/connectors' \
      -H 'Content-Type: application/json' \
      -d '{
          "name": "oracle_hub_source_connector"  ,
          "config": {
          "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
          "connection.url": "jdbc:oracle:thin:@xx.xx.xx.xx:1521/xxx",
          "connection.user": "xxx",
          "connection.password": "xxx",
          "mode": "timestamp",
          "timestamp.column.name": "xxx",
          "topic.prefix": "ORACLE_",
          "table.whitelist": "Table A",
          "poll.interval.ms": 5000,
          "db.timezone": "Asia/Seoul"
          }}'


      2. Stream A 생성
      CREATE STREAM STREAM A (
          abc VARCHAR )
      WITH (KAFKA_TOPIC='Topic A', VALUE_FORMAT='JSON');


      3. Stream A_key 생성
      CREATE STREAM STREAM A_key
      AS SELECT
         a, b, c
        FROM  STREAM A
          PARTITION BY c
         EMIT CHANGES;


      4. Table A_key 생성
      CREATE TABLE TABLE A_key(
       a VARCHAR,
       b VARCHAR,
       c VARCHAR)
      WITH (KAFKA_TOPIC='Topic A_key', VALUE_FORMAT='JSON', KEY='c');

      ksqlDB에 만들어진 Table A_key는 원천 DB의 Table A에 변경 데이터 발생시 즉각적으로 반영되게 된다.

      댓글 0

      DEVOTEE를 활성화 시키면
      지금 작성한 댓글에 AI가 댓글을 달아줍니다.

      jseung21 님의 최신 블로그

      더보기
      동영상 기고하기