NPSaaS et Kafka

Présentation

Apache Kafka est un système de messagerie de publication / abonnement que vous pouvez utiliser pour déplacer des données entre des applications populaires.

Après avoir intégré votre instance IBM® Netezza® Performance Server for IBM Cloud Pak® for Data as a Service à Kafka via le connecteur Kafka JDBC, vous pouvez utiliser NPSaaS comme suit:

  • Une source de données, qui fournit des données à Kafka.
  • Un collecteur de données, qui lit les données de Kafka.

Utiliser NPSaaS comme source de données

Une société de commerce électronique stocke ses listes de produits dans une base de données NPSaaS. Pour rationaliser l'expérience de recherche dans l'application et accéder à l'analyse en temps réel, les applications client (par exemple, Elasticsearch et Apache Flink) ont accès aux listes.

Dans ce cas, les données sont lues à partir de NPSaaS via le connecteur source Kafka JDBC et Kafka diffuse les données. Les applications client lisent à partir du flux et traitent les données.

L'image suivante illustre le flux de données NPSaaS en tant que source de données.

NPSaaS en tant que source de données
Image 1. Le diagramme illustre la façon dont Kafka lit les données de Netezza via le connecteur source JDBC et permet aux applications client d'y accéder.

Utilisation de NPSaaS comme collecteur de données

Pour améliorer les résultats des patients, identifier efficacement les facteurs de risque et accélérer les temps d'intervention, un hôpital extrait des informations utiles des données physiologiques. Différents ensembles de données provenant de différents canaux sont analysés à mesure qu'ils arrivent.

Dans ce cas, les données entrantes sont diffusées via Kafka, puis calculées. Les producteurs sont les sources de données physiologiques qui proviennent de différents canaux.

Une fois les données traitées, elles sont stockées dans NPSaaS à des fins d'enregistrement de l'historique du patient via le connecteur JDBC JDBC Kafka.

L'image suivante illustre le flux de données pour NPSaaS en tant que collecteur de données.

NPSaaS en tant que collecteur de données
Image 2. Le diagramme illustre la façon dont les données entrantes provenant de différents producteurs sont diffusées en flux et calculées par Kafka via le pilote JDBC et stockées sur Netezza.

Intégration de NPSaaS et Kafka

Si vous souhaitez intégrer votre instance NPSaaS à Kafka, vous devez utiliser le connecteur Kafka JDBC.

Le connecteur Kafka JDBC prend en charge les connecteurs JDBC source et récepteur. Avec le connecteur source, vous pouvez transférer des données de NPSaaS vers des rubriques Kafka. Avec le connecteur récepteur, vous pouvez transférer des données des rubriques Kafka vers NPSaaSà l'aide du connecteur Kafka JDBC.

Configuration du connecteur JDBC Kafka

Vous devez installer le pilote dans la bibliothèque de Kafkaen éditant plugin.path.

  1. Configurez Java.

    Pour que Kafka fonctionne, vous avez besoin de Java 8 ou version ultérieure ; pour les connecteurs JDBC, vous avez besoin de Java 11.

    sudo yum install -y java-11-openjdk-headless
    
  2. Copiez le pilote Netezza JDBC dans Kafka libs

    cp path/to/nzjdbc3.jar kafka/libs
    
  3. Configurez votre connecteur.

    Dans l'exemple, le connecteur Aiven est utilisé.

    a) Téléchargez le connecteur.

    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) Décompressez le paquet.

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

    c) Editez plugin.path dans Kafka config/connect-[distributed/standalone].properties pour qu'il pointe vers le dossier extrait. Avec cela, Kafka peut trouver et charger les fichiers JAR de plug-in et de connecteur.

  4. Démarrez la connexion répartie.

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

    Sortie :

    [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. Enregistrez les connecteurs.

    La base de données et la table doivent exister sur NPSaaS avant de tenter d'enregistrer les connecteurs.

    • Pour le connecteur source, exécutez la commande suivante.

      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"
      }
      
    • Pour le connecteur de collecteur, exécutez la commande suivante.

      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. Vérifiez si l'enregistrement a réussi.

    a) Lister les connecteurs.

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

    b) Vérifier l'état des connecteurs.

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

Vous pouvez vérifier si les connecteurs fonctionnent correctement en vérifiant les journaux NPSaaS.

  • Pour le connecteur source:

    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"
    
  • Pour le connecteur de collecteur:

    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)