準備工作
建立MaxCompute表
DataHub支援將資料同步到MaxCompute對應的資料表中,同時支援分區表和非分區表,一般情況下推薦使用者使用分區表進行資料同步以方便MaxCompute資料處理。
目前DataHub支援將TUPLE和BLOB的資料同步到MaxCompute資料表中。
針對TUPLE類型topic,MaxCompute目標表資料類型需要和DataHub資料類型相匹配,具體的資料類型映射關係如下:
MaxCompute
DataHub
BIGINT
BIGINT
STRING
STRING
BOOLEAN
BOOLEAN
DOUBLE
DOUBLE
DATETIME
TIMESTAMP
DECIMAL
DECIMAL
TINYINT
TINIINT
SMALLINT
SMALLINT
INT
INTEGER
FLOAT
FLOAT
MAP
不支援
ARRAY
不支援
由於目前DataHub並不能完全支援MaxCompute所有的資料類型,請使用者盡量根據DataHub資料類型建立MaxCompute表結構。
針對BLOB資料類型,要求MaxCompute表結構僅需要包含一列STRING類型的column即可,DataHub預設會將資料同步到該column中。
DataHub
MaxCompute
BLOB
STRING
同時為了方便資料追蹤和問題排查,建議使用者在建立MaxCompute表結構時,增加一列
__rowkey__ STRING欄位,DataHub會自動將DataHub對應資料的trace資訊同步到該列中,以方便後續資料排查。
準備同步任務帳號並授權
建立同步MaxCompute任務時,需要使用者手動填寫訪問MaxCompute表的帳號資訊,請使用者確保填入有效帳號資訊(一般情況下採用MaxCompute子帳號即可)。
需要給該帳號授予訪問MaxCompute表的相應許可權,具體許可權包括
CreateInstance、Describe、Alter以及Update許可權。使用者可以使用DataWorks管控台進行MaxCompute對應表的許可權管理,參考配置MaxCompute引擎許可權配置MaxCompute引擎許可權,也可以選擇使用MaxCompute的命令列工具進行授權,參考MaxCompute使用及授權管理。
確認TimestampUnit單位
Connector中TimestampUnit的作用,就是將資料中TIMESTAMP類型的資料(如果有),以TimestampUnit為單位進行轉換後寫入到下遊系統的日期類型(如datetime類型)。
如果TIMESTAMP列寫入的是以秒為單位的值,那建立Connector的時候TimestampUnit就選擇“SECOND”;如果寫入的是以毫秒為單位的值,那就選擇“MILLISECOND”;如果寫入的是以微秒為單位的值,那就選擇“MICROSECOND”。
根據MaxCompute目前的寫入標準,分區數越多就會導致DataHub同步資料越慢。因此,在建立MaxCompute同步任務時,請儘可能的控制分區數,尤其是USER_DEFINE同步模式。
同一分區的資料越連續越好,不要頻繁的分區跳變。
同步模式控制建立分區時,請不要建立過多的分區數。
當MaxCompute專案開啟白名單功能時,僅允許白名單內的裝置訪問專案空間;開啟MaxCompute ip白名單後需設定服務白名單才能保證同步服務正常訪問,具體設定方式請參考概述。
同步模式
Append模式
資料以追加的方式寫入目標表中,這種模式適用於資料僅需追加、無需更新的情境。
Upsert模式
Upsert 是 Update + Insert 的組合操作,其核心邏輯是:
如果目標表中存在與目前記錄主鍵相同的記錄,則更新該記錄;
如果目標表中不存在與目前記錄主鍵相同的記錄,則插入該記錄。
通過 Upsert 模式,使用者可以更靈活地處理資料更新和插入操作,確保目標表中的資料始終保持最新狀態。
關於更多 MaxCompute Upsert功能支援文檔參考 基本概念
應用情境
資料需要根據主鍵進行更新:資料可能隨時間變化,需要根據主鍵更新已有記錄;
保持目標表中資料的唯一性:需要確保目標表中每條記錄的唯一性,避免重複資料;
處理重複資料:需要對海量資料根據某個主鍵進行去重。
配置說明
Datahub Topic 類型:必須是TUPLE類型的Topic;
Datahub Topic Schema:目前支援以下兩種類型
dts同步到datahub的schema類型(以下統稱為Dts格式);
使用者自主建立的schema,需要選擇一列作為操作列,且必須是String類型(以下統稱為自訂格式);
Odps 目標表:必須是 Transaction2.0表。
同步規則
DTS格式
針對Dts同步Datahub的兩種資料格式,Datahub根據Schema中的固定列operation_flag, before_flag, after_flag以下面的規則決定如何將資料同步到 ODPS 目標表中:
operation_flag | before_flag | after_flag | OerationType | 同步目標表 |
I | * | * | UPSERT | 根據主鍵更新目標表中的記錄 |
U | Y | N | DELETE | 根據主鍵刪除目標表中的記錄 |
U | N | Y | UPSERT | 根據主鍵更新目標表中的記錄 |
D | * | * | DELETE | 根據主鍵刪除目標表中的記錄 |
自訂格式
針對使用者自主建立的資料,Datahub將根據使用者選擇的固定操作列來決定如何將資料同步到ODPS目標表中;
ddddd | OerationType | 同步目標表 |
U | UPSERT | 根據主鍵更新目標表中的記錄 |
D | DELETE | 根據主鍵刪除目標表中的記錄 |
建立同步任務
單擊DataHub中已建立的Topic,進入Topic詳情頁。
單擊Topic詳情頁右上方的同步按鈕,建立同步任務。
選擇MaxCompute類型作業,進入建立Connector頁面。
配置項說明:
參數
選項
是否必填
說明
Project名稱
/
是
MaxCompute project名稱,支援下拉框擷取,如無許可權擷取Project列表請手動輸入
Schema
/
否
MaxCompute schema名稱
說明使用Schema功能需開啟Schema文法開發 ,關於開啟檔案和更多Schema說明資訊請參考Schema操作
Table
/
是
MaxCompute Table名稱,支援下拉框選擇,如無許可權擷取Table列表請手動輸入
說明使用Upsert模式同步目標表必須是 Transaction2.0表
同步模式
Append
是
採用追加方式同步到MaxCompute 目標表中
Upsert
主鍵更新或刪除方式同步到MaxCompute DeltaTable
Upsert模式詳情參考上文同步模式章節
認證方式
AK
通過AK進行認證
DataHub預設角色
選擇該項後會在該project角色授權自動新增datahub__access__role 角色,具體權限原則為:
{ "Statement": [ { "Action": [ "odps:CreateInstance", "odps:CreateTable", "odps:Describe", "odps:Alter", "odps:Update" ], "Effect": "Allow", "Resource": [ "*" ] } ], "Version": "1" }自訂角色
使用者自訂角色,可通過ram控制台進行查看或者建立
Upsert方式
SYNC_CUSTOM
選擇同步模式為Upsert 該配置項為必填,同步方式為Append則不涉及該配置項
自訂Upsert操作欄位
SYNC_NONE
全部以Upsert形式寫入目標表
SYNC_DTS
適用於由dts寫入DataHub ,啟用dts新的附件列規則情境
SYNC_DTS_OLD
適用於由dts寫入DataHub,啟用dts新的附件列規則情境
主鍵欄位
/
Upsert同步方式下建立下遊表時指定的主鍵列
Upsert操作欄位
/
選擇SYNC_CUSTOM 方式同步該配置為必選項
選擇任意一個String類型的列作為Operation,用於表示當前資料是以Upsert或Delete的形式同步到下遊表中
關於Upsert 模式介紹請查看本文Upsert模式介紹章節。
匯入欄位:DataHub可以根據使用者佈建將部分column內容同步到MaxCompute表中。
分區模式:分區模式決定了將資料寫入到MaxCompute哪個分區中,目前DataHub支援以下分區方式:
分區模式
分區依據
支援Topic類型
說明
USER_DEFINE
Record中的分區列(和MaxCompute的分區欄位同名)的value值
TUPLE
DataHub schema中必須包含MaxCompute分區欄位
該列值必須為
UTF8字串欄位值可以為空白,表示不分區
SYSTEM_TIME
Record寫入DataHub的時間
TUPLE / BLOB
分區配置中設定MaxCompute分區的時間轉換Format格式
設定時區資訊
EVENT_TIME
Record中的
event_time(TIMESTAMP)列的value值TUPLE
分區配置中設定MaxCompute分區的時間轉換Format格式
設定時區資訊
META_TIME
Record的屬性欄位
__dh_meta_time__的value值TUPLE / BLOB
分區配置中設定MaxCompute分區的時間轉換Format格式
設定時區資訊
其中
SYSTEM_TIME、EVENT_TIME和META_TIME均是根據時間Timestamp和時區配置來進行MaxCompute分區的轉換過程,單位預設為微秒。分區配置決定了根據時間戳記轉換MaxCompute分區時的相關配置。目前管控台預設固定的MaxCompute分區格式,分區配置對應為:
分區
時間Format
說明
ds
%Y%m%d
day
hh
%H
hour
mm
%M
minute
分區間隔決定了根據時間戳記轉換MaxCompute分區時所採用的時間間隔。時間範圍是
15分鐘 ~ 1440分鐘(1天),跳變間隔15分鐘。時區資訊(TimeZone)時區資訊決定了根據時間戳記轉換MaxCompute分區時所採用的轉換時區。
分隔字元BLOB資料同步時,可以指定16進位分隔字元來決定是否對BLOB資料分割後再同步MaxCompute,比如
0A表示\n(分行符號)。Base64編碼DataHub BLOB預設儲存位元據,而MaxCompute對應的同步列為STRING類型,因此管控台建立同步任務時,預設採用base64編碼後進行同步,更多定製化需求請參考SDK實現。
查看同步任務
可以點擊對應connector的詳情頁面查看同步任務的運行狀態和點位等資訊, 包含同步點位、同步狀態以及重啟和停止等操作。
編輯同步任務
點擊同步任務頁面,選擇編輯,即可修改認證方式、匯入欄位、分區間隔、時區、TimestampUnit值。
同步樣本
USER_DEFINE同步模式
建立DataHub Topic
topic schema中必須需要包含MaxCompute分區欄位,類型為STRING。
向DataHub Topic寫入資料,可以使用datahub-sdk寫入資料。
測試過程中使用SDK寫入幾條資料,其中[ds,hh,mm]分別為:[20210304,01,15]和[20210304,02,15]。
建立同步任務
USER_DEFINE分區模式可以通過在同步中設定分區配置欄位,如果MaxCompute沒有對應的表,可自動建立。匯入欄位中設定匯入f1、f2欄位,不同步f3欄位。
確認同步資料
可以從DataHub管控台查看對應同步任務的同步資訊, 查詢MaxCompute資料結果。
在USER_DEFINE模式下,DataHub會根據
MaxCompute分組欄位所對應的value將DataHub中的資料同步到對應的分區中。
SYSTEM_TIME同步模式
建立DataHub Topic
由於分區是根據寫入DataHub時間來計算的,因此topic schema只需包含資料欄位,不需要包含分區欄位。
向DataHub Topic寫入資料,可以使用datahub-sdk進行資料寫入。
測試過程中使用SDK寫入幾條資料,DataHub目前對應的寫入時間為
2021-03-04 14:02:45。建立同步任務
請注意分區配置需要和MaxCompute表分區一致。
確認同步資料
可以從DataHub管控台查看對應同步任務的同步資訊,如DoneTime, 查詢MaxCompute資料結果。
在SYSTEM_TIME模式下,DataHub會根據
資料寫入DataHub的時間將DataHub中的資料同步到對應的分區中。
常見問題
同步到MaxCompute timestamp欄位時間變為1970-01-19
原因:DataHub同步MaxCompute預設時間戳記單位為微秒,使用者寫入時間戳記為毫秒。
解決方案:寫入DataHub時間戳記以微秒方式寫入。