NPSaaS e Kafka

Panoramica

Apache Kafka è un sistema di messaggistica di pubblicazione - sottoscrizione, che puoi utilizzare per spostare i dati tra le applicazioni più comuni.

Dopo aver integrato la tua istanza IBM® Netezza® Performance Server for IBM Cloud Pak® for Data as a Service con Kafka tramite il connettore Kafka JDBC, puoi utilizzare NPSaaS come uno dei seguenti:

  • Un'origine dati, che porta i dati a Kafka.
  • Un sink di dati, che legge i dati da Kafka.

Utilizzo di NPSaaS come origine dati

Una società di e-commerce memorizza i propri elenchi di prodotti in un database NPSaaS. Per semplificare l'esperienza di ricerca in - app e per accedere all'analisi in tempo reale, le applicazioni consumer (ad esempio, Elasticsearch e Apache Flink) hanno accesso agli elenchi.

In questo caso, i dati vengono letti dal connettore di origine NPSaaS tramite il connettore di origine Kafka JDBC e Kafka trasmette i dati. Le applicazioni consumer leggono dal flusso ed elaborano ulteriormente i dati.

La seguente immagine illustra il flusso di dati NPSaaS come origine dati.

NPSaaS come origine dati
Immagine 1. Il diagramma descrive in che modo Kafka legge i dati da Netezza tramite il connettore di origine JDBC e abilita le app consumer ad accedervi.

Utilizzo di NPSaaS come un data sink

Per migliorare i risultati dei pazienti, identificare in maniera efficiente i fattori di rischio e fornire tempi di intervento più rapidi, un ospedale estrae informazioni significative dai dati fisiologici. Diversi dataset da vari canali vengono analizzati man mano che arrivano.

In questo caso, i dati in entrata vengono trasmessi tramite Kafka e quindi calcolati. I produttori sono le fonti di dati fisiologici che provengono da canali diversi.

Una volta elaborati, i dati vengono archiviati in NPSaaS per scopi di registrazione della cronologia del paziente tramite il connettore sink Kafka JDBC.

La seguente immagine illustra il flusso di dati per NPSaaS come un sink dati.

NPSaaS as a data sink
Immagine 2. Il diagramma illustra il modo in cui i dati in entrata da vari produttori vengono trasmessi e calcolati da Kafka tramite il driver JDBC e memorizzati su Netezza.

Integrazione di NPSaaS e Kafka

Se vuoi integrare la tua istanza NPSaaS con Kafka, devi utilizzare il connettore Kafka JDBC.

Il connettore Kafka JDBC ha il supporto per i connettore JDBC di origine e di destinazione. Con il connettore di origine, puoi trasferire i dati da NPSaaS agli argomenti Kafka. Con il connettore sink, puoi trasferire i dati dagli argomenti Kafka in NPSaaSutilizzando il connettore JDBC JDBC Kafka.

Configurazione del connettore Kafka JDBC

Devi installare il driver nella libreria di Kafkamodificando plugin.path.

  1. Impostare Java.

    Perché Kafka funzioni, hai bisogno di Java 8 o versioni successive; per i connettori JDBC, hai bisogno di Java 11.

    sudo yum install -y java-11-openjdk-headless
    
  2. Copiare il driver Netezza JDBC in Kafka libs

    cp path/to/nzjdbc3.jar kafka/libs
    
  3. Configurare il connettore.

    Nell'esempio, viene utilizzato il connettore Aiven.

    a) Scaricare il connettore.

    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) Smettere il pacchetto.

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

    c) Modifica plugin.path in Kafka config/connect-[distributed/standalone].properties per puntare alla cartella estratta. Con ciò, Kafka può trovare e caricare il plugin e i jar del connettore.

  4. Avviare la connessione distribuita.

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

    Output:

    [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. Registrare i connettori.

    Il database e la tabella devono esistere su NPSaaS prima di provare a registrare i connettori.

    • Per il connettore di origine, eseguire il seguente comando.

      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"
      }
      
    • Per il connettore sink, eseguire il seguente comando.

      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. Verificare se la registrazione è stata eseguita correttamente.

    a) Elencare i connettori.

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

    b) Controllare lo stato dei connettori.

    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"
    }
    

Puoi verificare se i connettori funzionano correttamente controllando i log NPSaaS.

  • Per il connettore di origine:

    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"
    
  • Per il connettore sink:

    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)