NPSaaS 및 Kafka

개요

Apache Kafka 널리 사용되는 애플리케이션 간에 데이터를 이동하는 데 사용할 수 있는 게시-구독 메시징 시스템입니다.

Kafka JDBC 커넥터를 통해 IBM® Netezza® Performance Server for IBM Cloud Pak® for Data as a Service 인스턴스를 Kafka 통합한 후에는 NPSaaS 다음 중 하나로 사용할 수 있습니다:

  • Kafka 데이터를 가져오는 데이터 소스입니다.
  • Kafka 데이터를 읽는 데이터 싱크입니다.

NPSaaS 데이터 소스로 사용

이커머스 회사는 NPSaaS 데이터베이스에 제품 목록을 저장합니다. 인앱 검색 환경을 간소화하고 실시간 분석에 액세스하기 위해 소비자 앱(예: Elasticsearch 및 Apache Flink)은 목록에 액세스할 수 있습니다.

이 경우, NPSaaS Kafka JDBC 소스 커넥터를 통해 데이터를 읽으면 Kafka 데이터를 스트리밍합니다. 소비자 앱은 스트림에서 데이터를 읽고 추가적으로 데이터를 처리합니다.

다음 이미지는 데이터 소스로서의 데이터 흐름 NPSaaS 보여줍니다.

NPSaaS 데이터 소스로
이미지 1. 이 다이어그램은 Kafka JDBC 소스 커넥터를 통해 Netezza 데이터를 읽고 소비자 앱이 액세스할 수 있도록 하는 방법을 보여줍니다.

NPSaaS 데이터 싱크로 사용

환자 치료 결과를 개선하고, 위험 요소를 효율적으로 식별하며, 더 빠른 개입 시간을 제공하기 위해 병원에서는 생리학적 데이터에서 의미 있는 인사이트를 추출합니다. 다양한 채널의 다양한 데이터 세트가 도착하는 대로 분석됩니다.

이 경우 들어오는 데이터는 Kafka 통해 스트리밍된 다음 계산됩니다. 제작자는 다양한 채널에서 나오는 생리적 데이터의 출처입니다.

데이터가 처리된 후에는 Kafka JDBC 싱크 커넥터를 통해 환자 이력 기록을 위해 NPSaaS 저장됩니다.

다음 이미지는 데이터 싱크로서 NPSaaS 데이터 흐름을 보여줍니다.

NPSaaS 데이터 싱크로서
이미지 2. 이 다이어그램은 다양한 생산자로부터 들어오는 데이터를 Kafka JDBC 드라이버를 통해 스트리밍하고 계산하여 Netezza 저장하는 방법을 보여줍니다.

NPSaaS Kafka 통합하기

NPSaaS 인스턴스를 Kafka 통합하려면 Kafka JDBC 커넥터를 사용해야 합니다.

Kafka JDBC 커넥터는 소스 및 싱크 JDBC 커넥터를 지원합니다. 소스 커넥터를 사용하면 NPSaaS 데이터를 Kafka 토픽으로 전송할 수 있습니다. 싱크 커넥터를 사용하면 Kafka JDBC 커넥터를 사용하여 Kafka 토픽의 데이터를 NPSaaS전송할 수 있습니다.

JDBC Kafka 커넥터 설정하기

plugin.path 편집하여 Kafka 라이브러리에 드라이버를 설치해야 합니다.

  1. Java 설정합니다.

    Kafka 작동하려면 Java 8 이상이 필요하고, JDBC 커넥터의 경우 Java 11이 필요합니다.

    sudo yum install -y java-11-openjdk-headless
    
  2. ' Netezza ' JDBC ' 드라이버를 ' Kafka ' libs'로 복사합니다

    cp path/to/nzjdbc3.jar kafka/libs
    
  3. 커넥터를 설정합니다.

    이 예에서는 Aiven 커넥터가 사용되었습니다.

    a) 커넥터를 다운로드합니다.

    curl -SLO https://github.com/aiven/jdbc-connector-for-apache-kafka/releases/download/v6.6.2/jdbc-connector-for-apache-kafka-6.6.2.tar
    

    b) 패키지의 포장을 풉니다.

    tar xvf jdbc-connector-for-apache-kafka-6.6.2.tar
    

    c) 추출된 폴더를 가리키도록 Kafka ' config/connect-[distributed/standalone].properties '에서 ' plugin.path '을 편집합니다. 이를 통해 Kafka 플러그인 및 커넥터 용기를 찾아서 로드할 수 있습니다.

  4. 분산 연결을 시작합니다.

    connect-distributed.sh config/connect-distributed.properties
    

    출력:

    [2022-06-29 02:29:43,356] INFO Started o.e.j.s.ServletContextHandler@5562c2c9{/,null,AVAILABLE} (org.eclipse.jetty.server.handler.ContextHandler:915)
    [2022-06-29 02:29:43,356] INFO REST resources initialized; server is started and ready to handle requests (org.apache.kafka.connect.runtime.rest.RestServer:303)
    [2022-06-29 02:29:43,356] INFO Kafka Connect started (org.apache.kafka.connect.runtime.Connect:57)
    
  5. 커넥터를 등록합니다.

    커넥터 등록을 시도하기 전에 데이터베이스와 테이블이 NPSaaS 있어야 합니다.

    • 소스 커넥터의 경우 다음 명령을 실행합니다.

      curl -s -X POST -H "Content-Type: Application/json" --data '{ "name": "test-source-jdbc", "config": { "connector.class": "io.aiven.connect.jdbc.JdbcSourceConnector", "tasks.max": "1", "connection.url":"jdbc:netezza://localhost:5480/db1;user=user1;password=secret","connection.user":"user1","connection.password":"secret","dialect.name":"GenericDatabaseDialect","topic.prefix":"my","mode": "bulk", "poll.interval.ms":"60000", "table.whitelist":"TEST_TABLE", "batch.max.rows":"10000000" } }' http://localhost:8083/connectors | jq
      {
        "name": "test-source-jdbc",
        "config": {
          "connector.class": "io.aiven.connect.jdbc.JdbcSourceConnector",
          "tasks.max": "1",
          "connection.url": "jdbc:netezza://localhost:5480/db1;user=user1;password=secret,
          "connection.user": ""user1,
          "connection.password": "secret",
          "dialect.name": "GenericDatabaseDialect",
          "topic.prefix": "my",
          "mode": "bulk",
          "poll.interval.ms": "60000",
          "table.whitelist": "TEST_TABLE",
          "batch.max.rows": "10000000",
          "name": "test-source-jdbc"
        },
        "tasks": [],
        "type": "source"
      }
      
    • 싱크 커넥터의 경우 다음 명령을 실행합니다.

      curl -s -X POST -H "Content-Type: Application/json" --data '{ "name": "test-sink", "config": { "connector.class": "io.aiven.connect.jdbc.JdbcSinkConnector", "tasks.max": "2", "connection.url":"jdbc:netezza://localhost:5480/db1;user=user1;password=secret","connection.user":"user1","connection.password":"secret","dialect.name":"GenericDatabaseDialect","topics": "TEST_TABLE", "insert.mode": "insert" } }' http://localhost:8083/connectors | jq
      {
        "name": "test-sink",
        "config": {
          "connector.class": "io.aiven.connect.jdbc.JdbcSinkConnector",
          "tasks.max": "2",
          "connection.url": "jdbc:netezza://localhost:5480/db1;user=user1;password=secret",
          "connection.user": "user1",
          "connection.password": "secret",
          "dialect.name": "GenericDatabaseDialect",
          "topics": "TEST_TABLE",
          "insert.mode": "insert",
          "name": "test-sink"
        },
          "tasks": [],
          "type": "sink"
        }
      
  6. 등록이 성공했는지 확인합니다.

    a) 커넥터를 나열합니다.

    curl -s http://localhost:8083/connectors/ | jq
    [
      "test-source-jdbc",
      "test-sink"
    ]
    

    b) 커넥터의 상태를 확인합니다.

    curl -s http://localhost:8083/connectors/test-source-jdbc/status | jq
    {
         "name": "test-source-jdbc",
         "connector": {
             "state": "RUNNING",
             "worker_id": "10.11.112.15:8083"
      },
        "tasks": [
          {
              "id": 0,
              "state": "RUNNING",
              "worker_id": "10.11.112.15:8083"
           }
         ],
         "type": "source"
    }
    

NPSaaS 로그를 확인하여 커넥터가 올바르게 작동하는지 확인할 수 있습니다.

  • 소스 커넥터의 경우:

    2022-06-29 04:13:09.057046 PDT [27743]  DEBUG:  QUERY: SELECT 1
    2022-06-29 04:13:09.061443 PDT [27743]  DEBUG:  QUERY: SELECT * FROM "DB1"."USER1"."TEST_TABLE"
    ANALYZE
    2022-06-29 04:13:09.064209 PDT [27743]  DEBUG:  QUERY: SELECT * FROM "DB1"."USER1"."TEST_TABLE"
    
  • 싱크 커넥터의 경우:

    2022-06-29 09:08:32.379442 PDT [18976]  DEBUG:  QUERY: CREATE EXTERNAL TABLE bulkETL_18976_0 ( c0 nvarchar(8),c1 nvarchar(10),c2 nvarchar(4),c3 nvarchar(6),c4 nvarchar(7),c5 nvarchar(9),c6 nvarchar(7),c7 nvarchar(9),c8 nvarchar(12),c9 nvarchar(12),c10 nvarchar(25),c11 nvarchar(13),c12 nvarchar(1) )  USING (  DATAOBJECT('/tmp/junk')  REMOTESOURCE 'jdbc'  DELIMITER ' '  ESCAPECHAR '\'  CTRLCHARS 'YES'  CRINSTRING 'YES'  ENCODING 'INTERNAL'  MAXERRORS 1 QUOTEDVALUE 'YES' );
    2022-06-29 09:08:32.525024 PDT [18976]  DEBUG:  QUERY: INSERT INTO "TEST_TABLE"("C1","C2","C3","C4","C5","C6","C7","C8","DATE_PROD","TIME_PROD","TIMESTMP","TIMETZ_PROD","C18") VALUES(bulkETL_18976_0.c0,bulkETL_18976_0.c1,bulkETL_18976_0.c2,bulkETL_18976_0.c3,bulkETL_18976_0.c4,bulkETL_18976_0.c5,bulkETL_18976_0.c6,bulkETL_18976_0.c7,bulkETL_18976_0.c8,bulkETL_18976_0.c9,bulkETL_18976_0.c10,bulkETL_18976_0.c11,bulkETL_18976_0.c12)