您可以使用社區版Flink或阿里雲Realtime Compute版Flink訪問Lindorm寬表。本文介紹同時適用於阿里雲Flink和社區版Flink訪問Lindorm寬表的方法。
背景資訊
您可以將Lindorm寬表作為Flink中的維表或者結果表,通過Flink SQL或者Flink DataStream訪問Lindorm寬表。
選擇串連方式
在Lindorm寬表中,表根據建立方式的不同分成HBase表(使用HBase API建立並寫入資料的表)和SQL表(使用Lindorm SQL建立並寫入資料的表)兩種形態,訪問它們需要用不同的介面。因此在開始建立Flink任務之前,需要先確定連接器的種類,並進而根據連接器種類確定寬表的串連地址。
確定要訪問的寬表形態
在Lindorm中,可以使用Lindorm SQL來確定開發中的Flink任務所要訪問的寬表種類:
使用lindorm-cli(或者直接使用 Lindorm Insight 或 DMS)串連Lindorm寬表引擎。
執行下述SQL語句查看錶的
IS_HBASE_LIKE屬性。如果查得的屬性值為TRUE,則表示該表是一張HBase表。
如果查得的屬性值為FALSE,則表示該表是一張SQL表。
SHOW TABLE VARIABLES FROM table_name LIKE 'IS_HBASE_LIKE';
使用
lindorm-cli串連寬表引擎的方法,可參考文檔 通過Lindorm-cli串連並使用寬表引擎。關於
SHOW TABLE VARIABLE的詳細文法,可參考文檔 SHOW VARIABLES。
明確需要使用的連接器種類
使用者需要基於選定的Flink的產品形態,決定具體使用哪種連接器來訪問Lindorm寬表。
寬表形態 | 社區版Flink | 阿里雲Realtime ComputeFlink |
HBase表 | (支援作為維表、結果表) | (支援作為維表、結果表) |
SQL表 | (支援作為結果表) | (支援作為維表、結果表) |
阿里雲Realtime ComputeFlink指的是阿里雲上託管的Flink產品形態,詳情可參考文檔Realtime ComputeFlink版。請注意,在阿里雲ECS上使用開源Flink搭建的Flink叢集,仍然歸屬於上圖中的社區版Flink。
擷取寬表串連資訊
情境一:使用開源HBase連接器或雲原生多模資料庫Lindorm連接器
在此情境下,用於串連的地址需要使用寬表引擎的HBase Java API訪問地址(專用網路)。
在 Lindorm 執行個體詳情頁,單擊左側菜單資料庫連接,選擇寬表引擎頁簽。在通過 HBase 相容地址串連地區,擷取專用網路地址,格式為ld-<執行個體ID>-proxy-lindorm.lindorm.rds.aliyuncs.com:30020。在通過 MySql 相容地址串連地區,擷取 MySQL 相容地址,格式為ld-<執行個體ID>-proxy-sql-lindorm.lindorm.rds.aliyuncs.com:33060。情境二:使用JDBC連接器
在此情境下,用於串連的地址需要使用寬表引擎的MySQL相容地址(專用網路)。
在 Lindorm 執行個體詳情頁左側導覽列單擊資料庫連接,選擇寬表引擎頁簽。在通過 HBase 相容地址串連地區可查看 HBase Java API 的專用網路地址(格式為ld-<執行個體ID>-proxy-lindorm.rds.aliyuncs.com:30020)。
如果Flink任務中使用新建立的Lindorm使用者來訪問寬表,請確保該使用者有訪問Flink表的讀寫權限,賦予許可權的具體操作請參見為指定使用者賦予許可權。
Lindorm寬表的各種串連地址的詳細說明,可參考文檔查看寬表引擎串連地址。
Flink作業訪問寬表的方法
使用者可根據所選擇的Realtime Compute架構的通用方式進行開發。對於作業中訪問Lindorm寬表的需求,則可基於上文所選擇的連接器對照下面的文檔實現計算作業對Lindorm寬表的訪問。
前提條件
使用社區版Flink開發訪問Lindorm寬表的作業
如使用開源HBase連接器訪問寬表,則需保證寬表引擎的版本為2.4.3及以上版本。
如使用開源JDBC連接器訪問寬表,則需保證寬表引擎的版本為2.6.5.2及以上版本,且已開通MySQL協議相容功能。
使用阿里雲Realtime ComputeFlink開發訪問Lindorm寬表的作業,則對寬表引擎版本無限制。
確保 Flink 叢集所屬環境已與Lindorm執行個體實現了網路打通,且已將用戶端IP地址添加至Lindorm白名單。如何添加,請參見設定白名單。
社區版Flink作業開發
開源HBase連接器
使用社區版Flink開發訪問HBase表的作業情境下,如果想要通過公網訪問或訪問目標的Lindorm執行個體類型為Lindorm單節點,那麼在執行後續操作前,必須先升級SDK並更改配置。
具體操作,請參見通過HBase Java API串連並使用寬表引擎章節中的步驟1。
使用開源HBase連接器建立維表與結果表的具體操作方法請參見HBase連接器使用文檔。
開源JDBC連接器
使用開源JDBC連接器訪問Lindorm寬表時,目前只支援將寬表作為結果表。整體的使用方法可以參考JDBC連接器官方文檔,不過其中存在一些需要特別注意的地方,如下所示:
依賴要求
相對官方文檔中比較寬泛的依賴包版本,使用JDBC連接器訪問Lindorm寬表目前對於相關依賴包額版本限定在以下列表:flink-connector-jdbc-core-4.0.0-2.0.jarflink-connector-jdbc-mysql-4.0.0-2.0.jarmysql-connector-j-8.3.0.jar
MySQL JDBC驅動的依賴包可從社區官方下載。
JDBC連接器參數
由於不支援作為源表和維表,因此與源表、維表功能相關的參數都不支援(如類似於scan.fetch-size之類以 scan 作為首碼的參數以及lookup.cache之類以lookup作為首碼的參數等)。
JDBC的部分連接器參數使用建議如下:url:建議遵循文檔基於Java JDBC介面的應用開發進行配置。
username:使用 Lindorm 執行個體中建立的使用者名稱。
password:上述使用者名稱對應的密碼。
connector,table-name:遵循 JDBC 連接器社區的建議。
sink 相關參數:結合作業實際情況微調。
資料類型映射
Flink資料類型與Lindorm資料類型的映射原則上可參照MySQL的資料類型進行(可參照JDBC連接器章節資料類型映射),但有些資料類型Lindorm並沒有與MySQL對齊。比如在JDBC連接器中宣稱支援映射的MySQL類型中,下述類型在Lindorm中並不支援:MEDIUMINT類型
DATETIME類型
除BIGINT UNSIGNED以外的UNSIGNED類型
以下樣本通過Flink SQL定義了一個基於JDBC連接器訪問Lindorm寬表的作業。其中,假定在Lindorm寬表引擎中已經定義了一張表名為testflink的表。
# Flink建表和啟動任務
CREATE TABLE source_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '2',
'fields.c2.length' = '5',
'fields.c1.min' = '1',
'fields.c1.max' = '100'
);
CREATE TABLE sink_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://ld-xxxxx-proxy-lindorm.lindorm.rds.aliyuncs.com:33060/default?sslMode=disabled&allowPublicKeyRetrieval=true&useServerPrepStmts=true&useLocalSessionState=true&rewriteBatchedStatements=true&cachePrepStmts=true&prepStmtCacheSize=300&prepStmtCacheSqlLimit=50000000',
'username' = 'root',
'password' = 'root',
'table-name' = 'testflink'
);
INSERT INTO sink_table SELECT * FROM source_table;阿里雲Realtime ComputeFlink作業開發
雲原生多模資料庫Lindorm連接器
在Realtime ComputeFlink作業開發中,支援使用Flink SQL的方式開發訪問Lindorm寬表引擎的作業。關於Realtime ComputeFlink的作業開發詳情,可參考作業開發地圖。
使用連接器建立維表與結果表的具體操作方法請參見雲原生多模資料庫Lindorm連接器使用文檔。