Hologres資料來源為您提供讀取和寫入Hologres雙向通道的功能,本文為您介紹DataWorks的Hologres資料同步的能力支援情況。
使用限制
Hologres資料來源僅支援使用Serverless資源群組運行同步任務。
離線讀寫
Hologres Writer不支援將資料寫入Hologres的外部表格。
Hologres資料來源連通性擷取Hologres端點的邏輯:
當前地區的Hologres執行個體,Hologres端點擷取順序:any Tunnel > Single Tunnel > Public(公網)。
跨地區的Hologres執行個體,Hologres端點擷取順序:Public(公網)> Single Tunnel。
單表即時讀
Hologres版本必須在2.1以上。
不支援Hologres分區表的增量同步處理。
不支援Hologres表DDL變更訊息同步。
Hologres增量同步處理支援的資料類型包括以下類型:
INTEGER、BIGINT、TEXT、CHAR(n)、VARCHAR(n)、REAL、JSON、SERIAL、OID、INT4[]、INT8[]、FLOAT8[]、BOOLEAN[]、TEXT[]。
Hologres單表即時同步時,需開啟源端的Hologres資料庫的表Hologres Binlog,詳情可參見訂閱Hologres Binlog。
整庫即時寫
即時資料同步任務暫不支援同步沒有主鍵的表。
MySQL整庫即時同步資料至Hologres時,目前僅支援將資料寫入分區表子表,暫不支援寫入資料至分區表父表。
支援的欄位類型
欄位類型 | 離線讀(Hologres Reader) | 離線寫(Hologres Writer) | 即時寫 |
UUID | 不支援 | 不支援 | 不支援 |
CHAR | 支援 | 支援 | 支援 |
NCHAR | 支援 | 支援 | 支援 |
VARCHAR | 支援 | 支援 | 支援 |
LONGVARCHAR | 支援 | 支援 | 支援 |
NVARCHAR | 支援 | 支援 | 支援 |
LONGNVARCHAR | 支援 | 支援 | 支援 |
CLOB | 支援 | 支援 | 支援 |
NCLOB | 支援 | 支援 | 支援 |
SMALLINT | 支援 | 支援 | 支援 |
TINYINT | 支援 | 支援 | 支援 |
INTEGER | 支援 | 支援 | 支援 |
BIGINT | 支援 | 支援 | 支援 |
NUMERIC | 支援 | 支援 | 支援 |
DECIMAL | 支援 | 支援 | 支援 |
FLOAT | 支援 | 支援 | 支援 |
REAL | 支援 | 支援 | 支援 |
DOUBLE | 支援 | 支援 | 支援 |
TIME | 支援 | 支援 | 支援 |
DATE | 支援 | 支援 | 支援 |
TIMESTAMP | 支援 | 支援 | 支援 |
BINARY | 支援 | 支援 | 支援 |
VARBINARY | 支援 | 支援 | 支援 |
BLOB | 支援 | 支援 | 支援 |
LONGVARBINARY | 支援 | 支援 | 支援 |
BOOLEAN | 支援 | 支援 | 支援 |
BIT | 支援 | 支援 | 支援 |
JSON | 支援 | 支援 | 支援 |
JSONB | 支援 | 支援 | 支援 |
實現原理
離線讀
Hologres Reader支援兩種讀模數式:
JDBC模式(預設)
通過PSQL讀取Hologres表中的資料,根據表的Shard Count發起多個並發,每個Shard對應一個Select並發任務:
Hologres在建立表時,在同一個
CREATE TABLE事務中,通過CALL set_table_property('table_name', 'shard_count', 'xx')配置表的Shard Count。預設情況下,使用資料庫預設的Shard Count,具體數值取決於Hologres執行個體的配置。
Select語句通過表的內建列hg_shard_id的Shard篩選資料。
Arrow模式(useArrow=true)
通過Hologres的COPY OUT Arrow協議讀取資料,支援LZ4壓縮傳輸,效能更高:
需要Hologres版本 >= 4.0.18。
同樣按照Shard Count拆分並發任務。
不支援函數列和常量列(只支援表中的物理列)。
自動啟用壓縮傳輸(版本滿足時)。
離線寫
Hologres Writer通過資料同步架構擷取Reader產生的協議資料,根據writeMode(寫入模式)和conflictMode(衝突策略)的配置決定寫入資料時的通道和衝突解決方案策略。
寫入模式(writeMode)
寫入模式 | 實現方式 | 適用情境 | 版本要求 |
INSERT | 使用Holo Client批量寫入(INSERT ON CONFLICT)。 | 即時情境(包含資料回撤)、通用寫入情境。 | 所有版本 |
FIXED_COPY | 使用Holo Client Fixed Copy流式寫入。 | 即時情境(不包含資料回撤)、也可以用作離線匯入情境。 | Hologres >= 1.1 |
STAGE | 先寫入Hologres Internal Stage(Arrow格式),再從Stage匯入目標表。 | 離線同步推薦模式,大大量匯入,效能最優。 | Hologres >= 4.1.0 |
COPY | 使用JDBC COPY FROM STDIN寫入文本CSV資料。 | 相容舊版本的大量匯入情境。 | 所有版本 |
情境推薦
衝突處理模式(conflictMode)
您可以通過配置conflictMode,決定新匯入的資料和已有資料的主鍵發生衝突時,如何處理新匯入的資料:
conflictMode僅適用於有主鍵的表。具體寫入原理和效能,詳情請參考技術原理。
conflictMode為Replace(整行更新)模式時,新資料覆蓋舊資料,整行所有列全部覆蓋,沒有配置列映射的欄位會強制寫入NULL。
conflictMode為Update(更新)模式時,新資料覆蓋舊資料,只覆蓋配置有列映射的欄位。
conflictMode為Ignore(忽略)模式時,忽略新資料。
建立資料來源
在進行資料同步任務開發時,您需要在DataWorks上建立一個對應的資料來源,操作流程請參見資料來源管理,詳細的配置參數解釋可在配置介面查看對應參數的文案提示。
資料同步任務開發
單表離線
支援資料來源:Data Integration模組資料來源支援的所有資料來源類型。
配置指導:嚮導模式配置、指令碼模式配置。指令碼模式配置的全量參數和指令碼Demo請參見下文的附錄:指令碼Demo與參數說明。
單表即時
支援資料來源:DataHub、Hologres、Kafka、LogHub
配置指導:單表即時同步任務配置。
整庫離線
支援資料來源:AnalyticDB for MySQL 3.0、ClickHouse、Doris、Hologres、Oracle、PolarDB、SQL Server
配置指導:整庫離線同步任務配置
整庫即時
支援資料來源:AnalyticDB for OceanBase、MongoDB、MySQL、Oracle、PolarDB、PolarDB-X 2.0、PostgreSQL
配置指導:整庫即時同步任務配置
Serverless整庫即時
支援資料來源:MySQL
配置指導:Serverless同步任務配置
即時任務同步常見問題:即時同步至Hologres常見問題。
附錄:指令碼Demo與參數說明
離線任務指令碼配置方式
如果您配置離線任務時使用指令碼模式的方式進行配置,您需要按照統一的指令碼格式要求,在任務指令碼中編寫相應的參數,詳情請參見指令碼模式配置,以下為您介紹指令碼模式下資料來源的參數配置詳情。
Reader指令碼Demo
配置非分區表
配置從Hologres非分區表讀取資料至記憶體,如下所示。
{ "transform": false, "type": "job", "version": "2.0", "steps": [ { "stepType": "holo", "parameter": { "datasource": "holo_db", "envType": 1, "column": [ "tag", "id", "title", "body" ], "where": "", "table": "holo_reader_basic_src" }, "name": "Reader", "category": "reader" }, { "stepType": "stream", "parameter": { "print": false, "fieldDelimiter": "," }, "name": "Writer", "category": "writer" } ], "setting": { "executeMode": null, "failoverEnable": null, "errorLimit": { "record": "0" }, "speed": { "concurrent": 2, "throttle": false } }, "order": { "hops": [ { "from": "Reader", "to": "Writer" } ] } }Hologres表的DDL語句,如下所示。
begin; drop table if exists holo_reader_basic_src; create table holo_reader_basic_src( tag text not null, id int not null, title text not null, body text, primary key (tag, id)); call set_table_property('holo_reader_basic_src', 'orientation', 'column'); call set_table_property('holo_reader_basic_src', 'shard_count', '3'); commit;
配置分區表
配置從Hologres分區表的子表讀取資料至記憶體。
說明請注意partition的配置。
{ "transform": false, "type": "job", "version": "2.0", "steps": [ { "stepType": "holo", "parameter": { "selectedDatabase": "public", "partition": "tag=foo", "datasource": "holo_db", "envType": 1, "column": [ "tag", "id", "title", "body" ], "tableComment": "", "where": "", "table": "public.holo_reader_basic_part_src" }, "name": "Reader", "category": "reader" }, { "stepType":"stream", "parameter":{}, "name":"Writer", "category":"writer" } ], "setting":{ "errorLimit":{ "record":"0" }, "speed":{ "throttle":true, "concurrent":1, "mbps":"12" } }, "order":{ "hops":[ { "from":"Reader", "to":"Writer" } ] } }Hologres表的DDL語句,如下所示。
begin; drop table if exists holo_reader_basic_part_src; create table holo_reader_basic_part_src( tag text not null, id int not null, title text not null, body text, primary key (tag, id)) partition by list( tag ); call set_table_property('holo_reader_basic_part_src', 'orientation', 'column'); call set_table_property('holo_reader_basic_part_src', 'shard_count', '3'); commit; create table holo_reader_basic_part_src_1583161774228 partition of holo_reader_basic_part_src for values in ('foo'); # 確保分區表子表已經建立且匯入資料。 postgres=# \d+ holo_reader_basic_part_src Table "public.holo_reader_basic_part_src" Column | Type | Collation | Nullable | Default | Storage | Stats target | Description --------+---------+-----------+----------+---------+----------+--------------+------------- tag | text | | not null | | extended | | id | integer | | not null | | plain | | title | text | | not null | | extended | | body | text | | | | extended | | Partition key: LIST (tag) Indexes: "holo_reader_basic_part_src_pkey" PRIMARY KEY, btree (tag, id) Partitions: holo_reader_basic_part_src_1583161774228 FOR VALUES IN ('foo')
Reader指令碼參數
參數 | 描述 | 是否必選 | 預設值 |
database | Hologres執行個體內部資料庫的名稱。 | 是 | 無 |
table | Hologres的表名稱,支援 | 條件必選 | 無 |
querySql | 自訂查詢SQL,配置後table和column參數將被忽略。與 | 條件必選 | 無 |
column | 定義需要讀取的資料列。 | 是 | 無 |
partition | 針對分區表,表示分區Column以及對應的Value,格式為 重要
| 否 | 空,表示非分區表。 |
where | 過濾條件,會拼接到SELECT語句的WHERE子句中。僅在table模式下生效。 | 否 | 空 |
useArrow | 是否使用Arrow列存格式進行高效能資料同步。啟用後Reader使用COPY OUT Arrow協議讀取資料,以列式格式直通下遊Writer,效能更高。需要Hologres >= 4.0.18,版本不滿足時自動降級為JDBC模式。Arrow模式不支援函數列和常量列。目前支援源端與目標端為MaxCompute、Hologres、Hive/OSS/HDFS(Parquet/ORC)這幾種類型的整庫離線同步和單表離線同步。詳情請參見Arrow列存格式高效能同步。 | 否 | false |
compress | Arrow模式下是否啟用LZ4壓縮傳輸。僅在 | 否 | true(版本滿足時自動化佈建) |
fetchSize | JDBC模式下每次從資料庫擷取的行數。 | 否 | 1000 |
jdbcReadTimeout | JDBC讀取逾時時間,單位為秒。 | 否 | 60 |
enableServerlessComputing | 是否啟用Hologres Serverless Computing加速查詢。注意:此參數控制的是Hologres執行個體的Serverless Computing能力,與DataWorks的Serverless資源群組是不同的概念。 | 否 | false |
serverlessComputingQueryPriority | Hologres Serverless Computing查詢優先順序,範圍1-10,數字越大優先順序越高。僅在 | 否 | 3 |
serverlessComputingRequiredCores | Hologres Serverless Computing請求的核心數。僅在 | 否 | 5 |
Writer指令碼Demo
配置非分區表
配置從MySQL產生的資料匯入至Hologres普通表,樣本為通過INSERT模式匯入的配置。
{ "type": "job", "version": "2.0", "steps": [ { "stepType": "mysql", "parameter": { "envType": 0, "useSpecialSecret": false, "column": [ "<column1>", "<column2>", ......, "<columnN>" ], "tableComment": "", "connection": [ { "datasource": "<mysql_source_name>",//mysql資料來源名 "table": [ "<mysql_table_name>" ] } ], "where": "", "splitPk": "", "encoding": "UTF-8" }, "name": "Reader", "category": "reader" }, { "stepType": "holo", "parameter": { "selectedDatabase":"public", "schema": "public", "writeMode": "FIXED_COPY", "maxConnectionCount": 9, "truncate":true,//清理規則 "datasource": "<holo_sink_name>",//Hologres資料來源名稱 "conflictMode": "ignore", "envType": 0, "column": [ "<column1>", "<column2>", ......, "<columnN>" ], "tableComment": "", "table": "<holo_table_name>", "reShuffleByDistributionKey":false }, "name": "Writer", "category": "writer" } ], "setting": { "executeMode": null, "errorLimit": { "record": "0" }, "locale": "zh_CN", "speed": { "concurrent": 2,//作業並發數 "throttle": false//限流 } }, "order": { "hops": [ { "from": "Reader", "to": "Writer" } ] } }Hologres表的DDL語句,如下所示。
begin; drop table if exists mysql_to_holo_test; create table mysql_to_holo_test( tag text not null, id int not null, body text not null, brrth date, primary key (tag, id)); call set_table_property('mysql_to_holo_test', 'orientation', 'column'); call set_table_property('mysql_to_holo_test', 'distribution_key', 'id'); call set_table_property('mysql_to_holo_test', 'clustering_key', 'birth'); commit;
配置分區表
說明目前Hologres僅支援LIST分區,分區Column僅支援單個Column分區,且僅支援INT4或TEXT類型。
請確認該參數和表DDL的分區配置匹配。
配置從MySQL產生的資料同步至Hologres分區表的子表。
{ "type": "job", "version": "2.0", "steps": [ { "stepType": "mysql", "parameter": { "envType": 0, "useSpecialSecret": false, "column": [ "<column1>", "<column2>", ......, "<columnN>" ], "tableComment": "", "connection": [ { "datasource": "<mysql_source_name>", "table": [ "<mysql_table_name>" ] } ], "where": "", "splitPk": "<mysql_pk>",//mysql的pk欄位 "encoding": "UTF-8" }, "name": "Reader", "category": "reader" }, { "stepType": "holo", "parameter": { "selectedDatabase": "public", "writeMode": "insert", "maxConnectionCount": 9, "partition": "ds=20201215",//Hologres分區鍵 "truncate": "false", "datasource": "<holo_sink_name>",//Hologres資料來源名 "conflictMode": "ignore", "envType": 0, "column": [ "<column1>", "<column2>", ......, "<columnN>" ], "tableComment": "", "table": "<holo_table_name>", "reShuffleByDistributionKey":false }, "name": "Writer", "category": "writer" } ], "setting": { "executeMode": null, "failoverEnable": null, "errorLimit": { "record": "0" }, "speed": { "concurrent": 2,//作業並發數 "throttle": false//限流 } }, "order": { "hops": [ { "from": "Reader", "to": "Writer" } ] } }Hologres表的DDL語句,如下所示。
BEGIN; CREATE TABLE public.hologres_parent_table( a text , b int, c timestamp, d text, ds text, primary key(ds,b) ) PARTITION BY LIST(ds); CALL set_table_property('public.hologres_parent_table', 'orientation', 'column'); CREATE TABLE public.holo_child_1 PARTITION OF public.hologres_parent_table FOR VALUES IN('20201215'); CREATE TABLE public.holo_child_2 PARTITION OF public.hologres_parent_table FOR VALUES IN('20201216'); CREATE TABLE public.holo_child_3 PARTITION OF public.hologres_parent_table FOR VALUES IN('20201217'); COMMIT;
Writer指令碼參數
基礎參數
參數 | 描述 | 是否必選 | 預設值 |
database | Hologres執行個體內部資料庫的名稱。 | 是 | 無 |
table | Hologres的表名稱,目前支援表名稱中包含Schema,例如 | 是 | 無 |
writeMode | 寫入模式,支援 | 是 | 無 |
conflictMode | 衝突處理模式,包括 | 是 | 無 |
column | 定義匯入目標表的資料列,必須包含目標表的主鍵集合(serial自增主鍵和generated column可除外)。 | 是 | 無 |
partition | 針對分區表,表示分區Column以及對應的Value,格式為 說明
| 否 | 空,表示非分區表 |
reShuffleByDistributionKey | 在 Hologres 中,主鍵表的大量匯入預設會觸發表鎖,這限制了多個串連的並發寫入能力。開啟 reShuffle 功能可以在離線同步情境下,允許不同的任務根據資料分區鍵將資料寫入指定的 Holo shard,這樣可以實現並發批量寫入,從而顯著提升寫入效能。與傳統 JDBC 模式的即時寫入相比,啟用該功能不僅能降低 Holo 服務端的負載,還能進一步提升寫入效率。 重要 該功能僅在Serverless資源群組開啟。 | 否 | false |
truncate | 寫入Holo表之前是否需要清空目標表。
| 否 | false |
partitionFormat | 動態分區值格式化規則。對於Date類型來源資料,指定格式如 | 否 | 無 |
preSql | 寫入前執行的SQL列表。不支援在分區父表上執行。 | 否 | 無 |
postSql | 寫入後執行的SQL列表。 | 否 | 無 |
串連與績效參數
參數 | 描述 | 是否必選 | 預設值 |
maxConnectionCount | 寫入並發串連數。該值會根據實際任務並發數自動調整( | 否 | 3 |
maxRetryCount | 失敗重試次數。 | 否 | 10 |
jdbcReadTimeout | JDBC連線逾時時間,單位為秒。 | 否 | 120 |
maxCommitSize | 單次提交的最大位元組數(位元組)。 | 否 | 2097152(2MB) |
maxCommitCount | 單次提交的最大記錄數。 | 否 | 256 |
reShuffleByDistributionKey | 在Hologres中,主鍵表的大量匯入預設會觸發表鎖,這限制了多個串連的並發寫入能力。開啟reShuffle功能可以在離線同步情境下,允許不同的任務根據資料分區鍵將資料寫入指定的Holo shard,這樣可以實現並發批量寫入,從而顯著提升寫入效能。與傳統JDBC模式的即時寫入相比,啟用該功能不僅能降低Holo服務端的負載,還能進一步提升寫入效率。 重要 該功能僅在Serverless資源群組開啟。 | 否 | false |
removeU0000InTextColumnValue | 是否移除文本列值中的 | 否 | true |
Hologres Serverless Computing參數
以下參數用於控制Hologres執行個體的Serverless Computing能力,可加速Reader的查詢讀取以及STAGE模式Writer的INSERT FROM Stage操作。
注意:Hologres Serverless Computing與DataWorks Serverless資源群組是不同的概念。前者是Hologres執行個體內部的彈性計算能力,後者是DataWorks運行同步任務的調度資源。
適用範圍:讀取端(Reader)所有模式均支援;寫入端(Writer)僅STAGE模式支援。
參數 | 描述 | 是否必選 | 預設值 |
enableServerlessComputing | 是否啟用Hologres Serverless Computing加速。 | 否 | false |
serverlessComputingQueryPriority | Serverless Computing查詢優先順序,範圍1-10,數字越大優先順序越高。 | 否 | 3 |
serverlessComputingRequiredCores | Serverless Computing請求的核心數。設為0表示由引擎自動決定。 | 否 | 無(不設定) |
FIXED_COPY模式特有參數
參數 | 描述 | 是否必選 | 預設值 |
isBinaryFormat | 是否使用二進位格式進行資料轉送。設定為true時使用二進位格式,具有更高的傳輸效率;設定為false時使用文本(CSV)格式。 | 否 | true |
checkRecordBeforePut | 寫入前是否對Record進行校正(包括類型檢查、長度檢查等)。開啟後可以在用戶端提前發現髒資料。 | 否 | true |
maxCellBufferSize | 單行資料的最大緩衝區大小,單位為位元組。需要保證能放下一行資料,否則會寫入失敗。 | 否 | 10485760(10MB) |
STAGE模式特有參數
參數 | 描述 | 是否必選 | 預設值 |
stageTTL | Internal Stage的生命週期,單位為秒。Stage超過該時間後會被自動清理。 | 否 | 86400(1天) |
stageFileSizeLimit | 單個Stage檔案大小上限,單位為位元組。超過該限制後會自動建立新檔案。 | 否 | 67108864(64MB) |
stageMaxBatchSize | Record寫入Arrow批次的最大行數。每湊夠該行數後寫入一個Arrow RecordBatch。 | 否 | 8192 |
stageCompress | 是否啟用Arrow LZ4壓縮。啟用後可顯著減少Stage儲存空間和網路傳輸量。需要Hologres >= 4.2.8,版本不滿足時自動關閉。 | 否 | true |
進階參數
參數 | 描述 | 是否必選 | 預設值 |
useArrow | 是否使用Arrow列存寫入鏈路。設定為true時,如果Hologres版本 >= 4.1.0,會自動啟用STAGE寫入模式。Hologres Writer支援接收任意上遊Reader產出的Arrow格式資料(ArrowTabularRecord),也支援接收普通行式Record(自動轉換為Arrow格式寫入Stage)。詳情請參見Arrow列存格式高效能同步。 | 否 | false |
default.enable | 是否為NOT NULL但未設定值的列自動填滿預設值。 | 否 | true |
enableWriteBitTypeWithString | 是否允許以字串形式寫入BIT類型資料。 | 否 | false |
holoClient | Holo Client進階配置項(Map格式),可設定底層HoloConfig的任意參數。例如 | 否 | 無 |