NPSaaS y Kafka
Visión general
Apache Kafka es un sistema de mensajería de publicación/suscripción, que puede utilizar para mover datos entre aplicaciones populares.
Después de integrar la instancia de IBM® Netezza® Performance Server for IBM Cloud Pak® for Data as a Service con Kafka a través del conector Kafka JDBC, puede utilizar NPSaaS como uno de los siguientes:
- Un origen de datos, que aporta datos a Kafka.
- Un receptor de datos, que lee datos de Kafka.
Uso de NPSaaS como fuente de datos
Una empresa de comercio electrónico almacena sus listados de productos en una base de datos NPSaaS. Para agilizar la experiencia de búsqueda en la aplicación y acceder a la analítica en tiempo real, las aplicaciones de consumidor (por ejemplo, Elasticsearch y Apache Flink) tienen acceso a los listados.
En este caso, los datos se leen desde NPSaaS a través del conector de origen Kafka JDBC y Kafka transmite los datos. Las aplicaciones de consumidor leen de la secuencia y procesan más los datos.
La imagen siguiente ilustra el flujo de datos NPSaaS como un origen de datos.
Utilización de NPSaaS como sumidero de datos
Para mejorar los resultados de los pacientes, identificar de forma eficiente los factores de riesgo y proporcionar tiempos de intervención más rápidos, un hospital extrae información significativa de los datos fisiológicos. Los diferentes conjuntos de datos de varios canales se analizan a medida que llegan.
En este caso, los datos entrantes se transmiten a través de Kafka y, a continuación, se calculan. Los productores son las fuentes de datos fisiológicos que provienen de diferentes canales.
Una vez procesados los datos, se almacenan en NPSaaS para fines de registro de historial de pacientes a través del conector de receptor Kafka JDBC.
La imagen siguiente ilustra el flujo de datos para NPSaaS como un receptor de datos.
Integración de NPSaaS y Kafka
Si desea integrar la instancia de NPSaaS con Kafka, debe utilizar el conector Kafka JDBC.
El conector Kafka JDBC tiene soporte para conectores JDBC de origen y receptor. Con el conector de origen, puede transferir datos de NPSaaS a los temas de Kafka. Con el conector sink, puede transferir datos de los temas Kafka a NPSaaSutilizando el conector Kafka JDBC.
Configuración del conector JDBC Kafka
Debe instalar el controlador en la biblioteca de Kafkaeditando plugin.path.
-
Configure Java.
Para que Kafka funcione, necesita Java 8 o posterior; para los conectores JDBC, necesita Java 11.
sudo yum install -y java-11-openjdk-headless -
Copie el controlador Netezza JDBC en Kafka
libscp path/to/nzjdbc3.jar kafka/libs -
Configure el conector.
En el ejemplo, se utiliza el conector Aiven.
a) Descargue el conector.
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) Desempaquete el paquete.
tar xvf jdbc-connector-for-apache-kafka-6.6.2.tarc) Edite
plugin.pathen Kafkaconfig/connect-[distributed/standalone].propertiespara que apunte a la carpeta extraída. Con esto, Kafka puede encontrar y cargar el plugin y los jars de conector. -
Inicie la conexión distribuida.
connect-distributed.sh config/connect-distributed.propertiesSalida:
[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) -
Registre los conectores.
La base de datos y la tabla deben existir en NPSaaS antes de intentar registrar los conectores.
-
Para el conector de origen, ejecute el siguiente 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" } -
Para el conector de receptor, ejecute el mandato siguiente.
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" }
-
-
Compruebe si el registro se ha realizado correctamente.
a) Listar los conectores.
curl -s http://localhost:8083/connectors/ | jq [ "test-source-jdbc", "test-sink" ]b) Compruebe el estado de los conectores.
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" }
Puede verificar si los conectores funcionan correctamente comprobando los registros de NPSaaS.
-
Para el conector de origen:
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" -
Para el conector de sumidero:
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)