全部產品
Search
文件中心

DataHub:建立同步Elasticsearch

更新時間:Aug 26, 2026

準備工作

  1. 建立ES index

    DataHub⽀持將資料同步到Elasticsearch對應的index中,目前支援ES5、ES6和ES7的執行個體。

    ⽬前DataHub僅⽀持將TUPLE類型Topic的資料同步到Elasticsearch中。開始同步任務之前請保證已經在ES中建立index或者允許自動建立index,否則同步任務會失敗。

  2. 準備同步任務帳號並授權

    建立同步ES任務時,需要使用者手動填寫ES的endpoint、index和帳號密碼等資訊。請確保所填入資訊帳號資訊真實有效,否則會導致建立同步任務失敗。

建立同步任務

  1. 進入專案列表 > project詳情 > Topic詳情頁面

  2. 點擊右上方的+同步

  3. 選擇Elasticsearch類型作業。

配置說明

Endpoint

Elasticsearch服務地址,需要填寫內網地址和內網連接埠,格式為內網地址:內網連接埠

例如:內網地址為 es-cn-xxx.elasticsearch.aliyuncs.com,內網連接埠為9200,則填入es-cn-xxx.elasticsearch.aliyuncs.com:9200

Index

目前支援兩種index的指定方式:靜態index和動態index。

靜態index

使用者預先建立好一個index或者允許自動建立index,所有的資料都會寫入該index。

動態index

使用者需要允許自動建立index,否則會寫入失敗。使用者可以指定一個時間周期或者指定某列作為index的產生方式。如果採用資料列產生index,那麼最多可以選擇1列。

支援配置的時間格式:

%Y

%m

%d

%U

  • 樣本一:每天淩晨產生一個新的index配置index為test_${%Y-%m-%d},如果當前日期為2021年3月31日,那麼最終寫入的index為test_2021-03-31

  • 樣本二:根據資料列產生新的index資料列中包含有一列col1,配置index為test_${col1},如果有兩條資料,這兩條資料的col1分別為AAABBB,那麼這兩條資料寫入的index分別為test_AAAtest_BBB

當ES中的index數量增多時,寫入資料會變慢,過多可能會導致DataHub寫入逾時,因此使用者使用動態index時,需要盡量避免產生的index數量過多

User/Password

訪問ES的使用者名稱密碼。

Type屬性列

針對不同版本,DataHub同步ES的產生的Type也不一樣。在ES5中,使用者可以在一個index中建立多個type,但是ES6中,使用者只能在一個index中建立一個type,因此,DataHub同步ES的行為也有所改變。Type不可以為空白。

  • 對於ES5,DataHub同步資料時,將會以使用者選擇作為Type的列的值作為一條資料的type,如果選擇多列,則多列的值會以 “|” 分割作為一條資料的type。選擇作為Type屬性列的欄位不能為null

  • 對於ES6,DataHub同步資料時,將會以使用者選擇的列的列名作為一條資料的type,如果選擇多列,則多列的列名會以“|”分割作為一條資料的type,並且ES6支援以任意名稱作為type。

  • 目前在頁面上建立同步ES6任務時無法自訂type名稱,如果想要自訂使用者可以使用SDK來建立。

  • ES7中所有的type均採用預設type,所以 ES7不需要選擇type屬性列。

例如:

DataHub Schema : f1 string, f2 string, f3 string, f4 string
資料 : ["test1","test2","test3",null]

type屬性列

ES5 type

ES6 type

f1

test1

f1

f1,f3

test1

test3

ff

建立失敗

ff

f1,ff

建立失敗

f1|ff

f4

建立成功,但同步失敗,髒資料

建立成功,並成功同步

ID屬性列

使用者可以根據寫入DataHub的資料來產生寫入ES的資料id,也可以不選擇任何列,由ES將會為每條資料產生一個唯一的id。DataHub同步ES時,將會以使用者選擇的列的值作為一條資料的id,如果選擇多列,則多列的值會以 “|” 分割作為一條資料的id。選擇作為ID屬性列的欄位不能為null

例如:

DataHub Schema : f1 string, f2 string, f3 string, f4 string
資料 : ["test1","test2","test3",null]

ID屬性列

資料id

ES自動產生唯一ID

f1

test1

f1,f3

test1|test3

ff

建立失敗

f4

建立成功,但是同步失敗,髒資料

Router屬性列

根據寫入DataHub的資料產生寫入ES的router,也可以不選擇任何列,不選擇時將不使用ES的Router功能。DataHub同步ES時,將會以使用者選擇的列的值作為一條資料的router,如果選擇多列,則多列的值會以 “|” 分割作為一條資料的id。選擇作為Router屬性列的欄位不能為null

樣本可參考ID屬性列

匯入欄位

DataHub需要匯入到ES的欄位,對於未選擇的欄位,DataHub不會同步到ES中。對於作為ID的欄位和ES5中作為Type的欄位,不會再放到資料中。因為匯入欄位並不能完全決定最後產生的資料,因此不再給出樣本,下文會給出完整的資料同步樣本。

網路類型

根據ES樣本的類型選擇,目前公用雲端上的ES執行個體均為VPC執行個體,所以DataHub公用雲端同步ES任務只支援VPC網路類型。使用VPC網路類型時,需要填寫VPC ID和執行個體ID等資訊。

填入執行個體ID時需要注意加上-worker,例如執行個體ID為es-cn-xxx,則執行個體ID填寫es-cn-xxx-worker

寫入資料樣本

這裡的樣本是指建立同步ES任務成功之後,如果建立同步ES任務失敗,則參考上述的配置進行修改。

DataHub Schema為:

欄位名稱

欄位類型

f1

BIGINT

f2

STRING

f3

BOOLEAN

f4

DOUBLE

f5

TIMESTAMP

f6

DECIMAL

  • 樣本1:

    • Type屬性列 = f1(ES7無Type屬性列)

    • ID屬性列 = f2

    • 匯入欄位 = f1,f2,f3,f4,f5,f6

      • 資料 = v1,v2,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        v1

        v2

        {f3:v3,f4:v4,f5:v5,f6:v6}

        ES6

        f1

        v2

        {f1:v1,f3:v3,f4:v4,f5:v5,f6:v6}

        ES7

        -

        v2

        {f1:v1,f3:v3,f4:v4,f5:v5,f6:v6}

      • 資料 = null,v2,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        type屬性列為null,髒資料

        ES6

        f1

        v2

        {f1:v1,f3:v3,f4:v4,f5:v5,f6:v6}

        ES7

        -

        v2

        {f1:v1,f3:v3,f4:v4,f5:v5,f6:v6}

      • 資料 = v1,null,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        id屬性列為null,髒資料

        ES6

        -

        -

        id屬性列為null,髒資料

        ES7

        -

        -

        id屬性列為null,髒資料

  • 樣本2:

    • Type屬性列 = f1,f2(ES7無Type屬性列)

    • ID屬性列 = f3,f4

    • 匯入欄位 = f1,f2,f3,f4,f5,f6

      • 資料 = v1,v2,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        v1|v2

        v3|v4

        {f5:v5,f6,v6}

        ES6

        f1|f2

        v3|v4

        {f1:v1,f2:v2,f5:v5,f6:v6}

        ES7

        -

        v3|v4

        {f1:v1,f2:v2,f5:v5,f6:v6}

      • 資料 = v1,null,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        type屬性列為null,髒資料

        ES6

        f1|f2

        v3|v4

        {f1:v1,f2:v2,f5:v5,f6:v6}

        ES7

        -

        v3|v4

        {f1:v1,f2:v2,f5:v5,f6:v6}

      • 資料 = v1,v2,null,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        id屬性列為null,髒資料

        ES6

        -

        -

        id屬性列為null,髒資料

        ES7

        -

        -

        id屬性列為null,髒資料

  • 樣本3:

    • Type屬性列 = f1(ES7無Type屬性列)

    • ID屬性列 = f2

    • Router屬性列 = f3

    • 匯入欄位 = f1,f2,f3,f4,f5,f6

      • 資料 = v1,v2,v3,v4,v5,v6

        ES版本

        type

        id

        router

        data

        ES5

        v1

        v2

        v3

        {f4:v4,f5:v5,f6:v6}

        ES6

        f1

        v2

        v3

        {f1:v1,f4:v4,f5:v5,f6:v6}

        ES7

        -

        v2

        v3

        {f1:v1,f4:v4,f5:v5,f6:v6}

      • 資料 = null,v2,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        type屬性列為null,髒資料

        ES6

        f1

        v2

        {f1:v1,f4:v4,f5:v5,f6:v6}

        ES7

        -

        v2

        {f1:v1,f4:v4,f5:v5,f6:v6}

      • 資料 = v1,null,v3,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        id屬性列為null,髒資料

        ES6

        -

        -

        id屬性列為null,髒資料

        ES7

        -

        -

        id屬性列為null,髒資料

      • 資料 = v1,v2,null,v4,v5,v6

        ES版本

        type

        id

        data

        ES5

        -

        -

        router屬性列為null,髒資料

        ES6

        -

        -

        router屬性列為null,髒資料

        ES7

        -

        -

        router屬性列為null,髒資料

查看同步任務

可以點擊對應connector的詳情⻚⾯查看同步任務的運⾏狀態和點位等資訊, 包含同步點位、同步狀態以及重啟和停⽌等操作。需要先停止任務才可以進行儲值點位操作。

同步樣本

本樣本以阿里雲ES6.7為例,展示DataHub同步到ES的完整操作。其中ES相關的操作均使用Kibana Dev Tools執行,其他動作方法請參考ES官方文檔

  1. 建立ES index

    一般情況下ES預設自動建立index,因此該步驟可忽略。如果設定不允許自動建立,需要手動建立index,具體建立命令可參考ES官方文檔

  2. 建立DataHub Topic

    前僅支援Tuple類型的Topic建立同步ES任務。建立DataHub的過程可以參考Topic操作

  3. 建立同步ES任務

    本樣本中建立同步ES任務時,將f1和f2作為Type屬性列,將f3和f4作為ID屬性列,並將所有欄位作為匯入欄位。

  4. 向DataHub Topic寫⼊資料

    使⽤DataHub-SDK或者DataHub外掛程式進⾏資料寫⼊。寫入一條資料後,在頁面上抽樣,查看寫入的資料。

  5. 確認同步資料

    首先查看ES同步任務的點位。若同步任務的點位和同步時間發生了變化,同步時間即資料寫入DataHub的時間,同步點位為1(點位從0開始計算,點位為0表示第一條資料已經寫入)。

    在ES查看一下資料的同步情況,通過Kibana可以看到資料已經同步成功。