NPSaaS und Kafka
Übersicht
Apache Kafka ist ein Publish/Subscribe-Messaging-System, mit dem Sie Daten zwischen gängigen Anwendungen verschieben können.
Nach der Integration Ihrer IBM® Netezza® Performance Server for IBM Cloud Pak® for Data as a Service-Instanz mit Kafka über den Kafka JDBC-Connector können Sie NPSaaS wie folgt verwenden:
- Eine Datenquelle, die Daten in Kafkabringt.
- Eine Datensenke, die Daten aus Kafkaliest.
Verwendung von NPSaaS als Datenquelle
Ein E-Commerce-Unternehmen speichert seine Produktlisten in einer NPSaaS-Datenbank. Um die App-interne Suchfunktionalität zu optimieren und auf Echtzeitanalysen zuzugreifen, haben Konsumenten-Apps (z. B. Elasticsearch und Apache Flink) Zugriff auf die Listen.
In diesem Fall werden Daten aus NPSaaS über den Kafka JDBC-Quellenconnector gelesen und Kafka streamt die Daten. Die Konsumenten-Apps lesen aus dem Datenstrom und verarbeiten die Daten weiter.
Die folgende Abbildung zeigt den Datenfluss NPSaaS als Datenquelle.
NPSaaS als Datensenke verwenden
Um die Patientenergebnisse zu verbessern, Risikofaktoren effizient zu identifizieren und schnellere Interventionszeiten zu ermöglichen, extrahiert ein Krankenhaus aussagekräftige Erkenntnisse aus physiologischen Daten. Unterschiedliche Datensätze aus verschiedenen Kanälen werden analysiert, sobald sie ankommen.
In diesem Fall werden die eingehenden Daten über Kafka gestreamt und anschließend berechnet. Die Produzenten sind die Quellen physiologischer Daten, die aus verschiedenen Kanälen stammen.
Nach der Verarbeitung der Daten werden sie über den JDBC-Sink-Connector von Kafka in NPSaaS gespeichert.
Die folgende Abbildung zeigt den Datenfluss für NPSaaS als Datensenke.
NPSaaS und Kafka integrieren
Wenn Sie Ihre NPSaaS-Instanz in Kafkaintegrieren möchten, müssen Sie den Connector Kafka JDBC verwenden.
Der Connector Kafka JDBC bietet Unterstützung für JDBC-Connectors für Quellen und Senken. Mit dem Quellenconnector können Sie Daten aus NPSaaS in Kafka-Topics übertragen. Mit dem Sink-Connector können Sie Daten aus Kafka-Topics in NPSaaSübertragen, indem Sie den Kafka JDBC-Connector verwenden.
Kafka-Connector für JDBC einrichten
Sie müssen den Treiber in der Bibliothek Kafkainstallieren, indem Sie plugin.pathbearbeiten.
-
Konfigurieren Sie Java.
Damit Kafka funktioniert, benötigen Sie Java 8 oder höher; für die JDBC-Connectors benötigen Sie Java 11.
sudo yum install -y java-11-openjdk-headless -
Kopieren Sie den Treiber Netezza JDBC in Kafka
libs.cp path/to/nzjdbc3.jar kafka/libs -
Richten Sie Ihren Connector ein.
Im Beispiel wird der Aiven-Connector verwendet.
a) Laden Sie den Connector herunter.
curl -SLO https://github.com/aiven/jdbc-connector-for-apache-kafka/releases/download/v6.6.2/jdbc-connector-for-apache-kafka-6.6.2.tarb) Entpacken Sie das Paket.
tar xvf jdbc-connector-for-apache-kafka-6.6.2.tarc) Bearbeiten Sie
plugin.pathin Kafkaconfig/connect-[distributed/standalone].propertiesso, dass auf den extrahierten Ordner verwiesen wird. Damit kann Kafka das Plug-in und die Connector-JAR-Dateien suchen und laden. -
Starten Sie die verteilte Verbindung.
connect-distributed.sh config/connect-distributed.propertiesAusgabe:
[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) -
Registrieren Sie die Connectors.
Die Datenbank und die Tabelle müssen in NPSaaS vorhanden sein, bevor Sie versuchen, die Connectors zu registrieren.
-
Führen Sie für den Quellkonnektor den folgenden Befehl aus.
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" } -
Führen Sie für den Senkenconnector den folgenden Befehl aus.
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" }
-
-
Überprüfen Sie, ob die Registrierung erfolgreich war.
a) Die Connectors auflisten.
curl -s http://localhost:8083/connectors/ | jq [ "test-source-jdbc", "test-sink" ]b) Überprüfen Sie den Status der Connectors.
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" }
Sie können überprüfen, ob die Connectors ordnungsgemäß funktionieren, indem Sie die Protokolle von NPSaaS überprüfen.
-
Für den Quellenconnector:
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" -
Für den Senkenconnector:
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)