簡介
本文面向 BitSail 的 Connector 開發人員,通過開發者的角度全面的闡述開發一個完整 Connector 的全流程,快速上手 Connector 開發,
目錄結構
首先開發者需要通過 git 下載最新代碼到本地,并匯入到 IDE 中,同時創建自己的作業分支,使用該分支開發自己的 Connector,專案地址:https://github.com/bytedance/bitsail.git,
專案結構如下:
開發流程
BitSail 是一款基于分布式架構的資料集成引擎,Connector 會并發執行,并由 BitSail 框架來負責任務的調度、并發執行、臟資料處理等,開發者只需要實作對應介面即可,具體開發流程如下:
-
工程配置,開發者需要在
bitsail/bitsail-connectors/pom.xml模塊中注冊自己的 Connector,同時在bitsail/bitsail-dist/pom.xml增加自己的 Connector 模塊,同時為你的連接器注冊組態檔,來使得框架可以在運行時動態發現它,
-
Connector 開發,實作 Source、Sink 提供的抽象方法,具體細節參考后續介紹,
-
資料輸出型別,目前支持的資料型別為 BitSail Row 型別,無論是 Source 在 Reader 中傳遞給下游的資料型別,還是 Sink 從上游消費的資料型別,都應該是 BitSail Row 型別,
Architecture
當前 Source API 的設計同時兼容了流批一批的場景,換言之就是同時支持 pull & push 的場景,在此之前,我們需要首先再過一遍傳統流批場景中各組件的互動模型,
Batch Model
傳統批式場景中,資料的讀取一般分為如下幾步:
-
createSplits:一般在 client 端或者中心節點執行,目的是將完整的資料按照指定的規則盡可能拆分為較多的rangeSplits,createSplits在作業生命周期內有且執行一次, -
runWithSplit: 一般在執行節點節點執行,執行節點啟動后會向中心節點請求存在的rangeSplit,然后再本地進行執行;執行完成后會再次向中心節點請求直到所有splits執行完成, -
commit:全部的 split 的執行完成后,一般會在中心節點執行commit的操作,用于將資料對外可見,
Stream Model
傳統流式場景中,資料的讀取一般分為如下幾步:
-
createSplits:一般在 client 端或者中心節點執行,目的是根據滑動視窗或者滾動視窗的策略將資料流劃分為rangeSplits,createSplits在流式作業的生命周期中按斬訓分視窗的會一直執行, -
runWithSplit: 一般在執行節點節點執行,中心節點會向可執行節點發送rangeSplit,然后在可執行節點本地進行執行;執行完成后會將處理完的splits資料向下游發送, -
commit:全部的 split 的執行完成后,一般會向目標資料源發送retract message,實時動態展現結果,
BitSail Model
-
createSplits:BitSail 通過SplitCoordinator模塊劃分rangeSplits,在流式作業中的生命周期中createSplits會周期性執行,而在批式作業中僅僅會執行一次, -
runWithSplit: 在執行節點節點執行,BitSail 中執行節點包括Reader和Writer模塊,中心節點會向可執行節點發送rangeSplit,然后在可執行節點本地進行執行;執行完成后會將處理完的splits資料向下游發送, -
commit:writer在完成資料寫入后,committer來完成提交,在不開啟checkpoint時,commit會在所有writer都結束后執行一次;在開啟checkpoint時,commit會在每次checkpoint的時候都會執行一次,
Source Connector
-
Source: 資料讀取組件的生命周期管理,主要負責和框架的互動,構架作業,不參與作業真正的執行
-
SourceSplit: 資料讀取分片;大資料處理框架的核心目的就是將大規模的資料拆分成為多個合理的 Split
-
State:作業狀態快照,當開啟 checkpoint 之后,會保存當前執行狀態,
-
SplitCoordinator: 既然提到了 Split,就需要有相應的組件去創建、管理 Split;SplitCoordinator 承擔了這樣的角色
-
SourceReader: 真正負責資料讀取的組件,在接收到 Split 后會對其進行資料讀取,然后將資料傳輸給下一個算子
Source Connector 開發流程如下
-
首先需要創建
Source類,需要實作Source和ParallelismComputable介面,主要負責和框架的互動,構架作業,它不參與作業真正的執行 -
BitSail的Source采用流批一體的設計思想,通過getSourceBoundedness方法設定作業的處理方式,通過configure方法定義readerConfiguration的配置,通過createTypeInfoConverter方法來進行資料型別轉換,可以通過FileMappingTypeInfoConverter得到用戶在 yaml 檔案中自定義的資料源型別和 BitSail 型別的轉換,實作自定義化的型別轉換, -
最后,定義資料源的資料分片格式
SourceSplit類和闖將管理Split的角色SourceSplitCoordinator類 -
最后完成
SourceReader實作從Split中進行資料的讀取,
-
每個
SourceReader都在獨立的執行緒中執行,并保證SourceSplitCoordinator分配給不同SourceReader的切片沒有交集 -
在
SourceReader的執行周期中,開發者只需要關注如何從構造好的切片中去讀取資料,之后完成資料型別對轉換,將外部資料型別轉換成BitSail的Row型別傳遞給下游即可
Reader 示例

Sink Connector
-
Sink:資料寫入組件的生命周期管理,主要負責和框架的互動,構架作業,它不參與作業真正的執行,
-
Writer:負責將接收到的資料寫到外部存盤,
-
WriterCommitter(可選):對資料進行提交操作,來完成兩階段提交的操作;實作 exactly-once 的語意,
開發者首先需要創建Sink類,實作Sink介面,主要負責資料寫入組件的生命周期管理,構架作業,通過configure方法定義writerConfiguration的配置,通過createTypeInfoConverter方法來進行資料型別轉換,將內部型別進行轉換寫到外部系統,同Source部分,之后我們再定義Writer類實作具體的資料寫入邏輯,在write方法呼叫時將BitSail Row型別把資料寫到快取佇列中,在flush方法呼叫時將快取佇列中的資料刷寫到目標資料源中,
Writer 示例

將連接器注冊到組態檔中
為你的連接器注冊組態檔,來使得框架可以在運行時動態發現它,組態檔的定義如下:
以 hive 為例,開發者需要在 resource 目錄下新增一個 json 檔案,名字示例為 bitsail-connector-hive.json,只要不和其他連接器重復即可
測驗模塊
在 Source 或者 Sink 連接器所在的模塊中,新增 ITCase 測驗用例,然后按照如下流程支持
-
通過 testcontainer 來啟動相應的組件
-
撰寫相應的組態檔

-
通過代碼 EmbeddedFlinkCluster.submit 來進行作業提交

提交 PR
當開發者實作自己的 Connector 后,就可以關聯自己的 issue,提交 PR 到 github 上了,提交之前,開發者記得 Connector 添加檔案,通過 review 之后,大家貢獻的 Connector 就成為 BitSail 的一部分了,我們按照貢獻程度會選取活躍的 Contributor 成為我們的 Committer,參與 BitSail 社區的重大決策,希望大家積極參與!
活動推薦
1.快來加入 BitSail 激勵計劃,成為 Contributor!????
不僅可以 Get 到新技術,提升生產效率,還能結識到一群志同道合的小伙伴一起探討和成長~????
更有藍牙耳機、鍵盤、音箱等激勵好禮,送給為 BitSail 做出積極貢獻的你!????
issue 認領鏈接:https://github.com/bytedance/bitsail/issues
2.【BitSail 種子用戶】招募 ing,福利多多!
????活動流程:
(1)填寫開源調研問卷,參與抽獎,
(2)我們將抽取部分填寫問卷用戶,參與一對一用戶訪談,
????問卷鏈接: https://www.wjx.cn/vm/w8po1FA.aspx#
填寫開源調研問卷,即可參與抽取位元組跳動資料平臺帆布包、充電寶等精美禮品,
立即跳轉 BitSail GitHub 代碼倉庫了解更多!
https://github.com/bytedance/bitsail
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/539338.html
標籤:其他
