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) 可以存取清單。
在此情況下,會透過 Kafka JDBC 來源連接器從 NPSaaS 讀取資料,並 Kafka 對資料進行串流處理。 消費者應用程式從串流讀取並進一步處理資料。
下列影像說明作為資料來源的資料流程 NPSaaS。
使用 NPSaaS 作為資料接收槽
為了改善病患結果,有效識別風險因素,並提供更快速的介入時間,醫院會從生理資料中擷取有意義的洞察。 來自各種通道的不同資料集會在到達時進行分析。
在此情況下,送入的資料會透過 Kafka 進行串流,然後進行計算。 生產者是來自不同渠道的生理資料的來源。
處理資料之後,它會透過 Kafka JDBC 接收槽連接器儲存在 NPSaaS 以用於病患歷程記錄。
下列影像說明 NPSaaS 作為資料接收槽的資料流程。
整合 NPSaaS 與 Kafka
如果您想要整合 NPSaaS 實例與 Kafka,則必須使用 Kafka JDBC 連接器。
Kafka JDBC 連接器支援來源及接收槽 JDBC 連接器。 使用來源連接器,您可以將資料從 NPSaaS 傳送至 Kafka 主題。 使用接收槽連接器,您可以使用 Kafka JDBC 連接器,將資料從 Kafka 主題傳送至 NPSaaS。
設定 JDBC Kafka 連接器
您必須透過編輯 plugin.path,將驅動程式安裝在 Kafka的程式庫中。
-
設定 Java。
若要讓 Kafka 運作,您需要 Java 8 或更新版本; 對於 JDBC 連接器,您需要 Java 11。
sudo yum install -y java-11-openjdk-headless -
將 Netezza JDBC 驅動程式複製到 Kafka
libscp path/to/nzjdbc3.jar kafka/libs -
設定連接器。
在此範例中,使用 Aven 連接器。
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.tarb) 解壓縮套件。
tar xvf jdbc-connector-for-apache-kafka-6.6.2.tarc) 在 Kafka
config/connect-[distributed/standalone].properties中編輯plugin.path,以指向解壓縮的資料夾。 使用此功能,Kafka 可以尋找並載入外掛程式及連接器 JAR。 -
啟動分散式連接。
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) -
登錄連接器。
在嘗試登錄連接器之前,資料庫及表格必須存在於 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" }
-
-
驗證登錄是否成功。
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)