【東東老師 X 思想科技】資料匯入到 BigQuery 的方法

(二) 使用 Data Transfer Service 立即或排程匯入

如果你要正式傳輸大量資料到 Google Cloud BigQuery,推薦使用 Data Transfer Service( (簡稱 DTS))。它除了可以從地端的資料庫上傳之外、也支援從 Google 相關的行銷平台如 Google Ads、Ad Manager 和 Google Play,甚至其他雲端如 AWS S3 Storage、AWS Redshift、Azure Blog Storage、Salesforce 等,更多資料來源可以參考這份文件

針對不同來源,操作方法都差不多,我們就以 AWS S3 Storage 來源為例子,說明步驟如下。

我們先從 BigQuery 主選單進去,找到「資料移轉」再建立「移轉作業」: 

37 進入 Data Transfer Service
37 進入 Data Transfer Service
資料來源:擷圖自 GCP 主控台

接下來從「來源類型」點擊下拉式選單,你會看到它有非常多的來源可以選擇,我們就直接選 「Amazon S3」: 

38  從來源類型選擇 Amazon S3
38  從來源類型選擇 Amazon S3
資料來源:擷圖自 GCP 主控台

我們就來設定轉移的作業名稱,以及排程的執行頻率,例如每天執行一次就設定為 24 小時。 

39 設定 DTS 轉移名稱和頻率
39 設定 DTS 轉移名稱和頻率
資料來源:擷圖自 GCP 主控台

在目的地設定,我們輸入轉入  BigQuery 的表格名稱。而在 Amazon S3 的資料來源設定當中,有兩個重要的欄位叫做 Access Key ID 和 Secret Access Key: 

40 填寫 S3 資料來源詳細資料
40 填寫 S3 資料來源詳細資料
資料來源:擷圖自 GCP 主控台

我們就去 AWS 的主控台然後點擊右上角自己的 ID,進入安全憑證頁面:

41 AWS 主控台進入 Security Credentials
41 AWS 主控台進入 Security Credentials
資料來源:擷圖自 AWS 主控台

接著點擊「建立存取金鑰」:

42 建立 Access Key
42 建立 Access Key
資料來源:擷圖自 AWS 主控台

接著就進入到存取金鑰的頁面,它有提供 Access Key ID 和 Secret Access Key,如果你只是短暫測試,就直接將兩個值複製起來,不要下載 CSV 檔案,以免金鑰被駭客拿到。 

43 下載或複製 Access Key
43 下載或複製 Access Key
資料來源:擷圖自 AWS 主控台

Access Key ID 和 Secret Access Key 複製後貼入 DTS 的內容欄位。

另外 Write Disposition 指的是,你匯入資料後,要不要清除原始資料 (WRITE_TRUNCATE),或是一直附加到表格的最後面 (WRITE_APPEND)。

44 AWS S3 Access Key 和寫入方式設定
44 AWS S3 Access Key 和寫入方式設定
資料來源:擷圖自 GCP 主控台

而在傳輸選項的部分:

Number of Error Allowed

跟之前建立 BigQuery 一樣,如果資料量夠大,建議設定一些允許的錯誤,以免整個傳輸作業有任何一點錯誤,導致全部重來。

Decimal Target Types 

這是非常重要的欄位,指的是關於如何處理從 S3 中的 decimal (十進位數) 格式資料轉換到 BigQuery 時的資料型態轉換規則。

你可以先指定比較偏好的型態,再指定退而求其次的型態,最多三個。在傳輸過程中,先用第一個型態來接收資料,如果接收到不符合這個型態的,整個欄位自動調整成第二個型態來接收看看,不行再改成第三個型態。

例如你有一個 decimal(38,9) 的欄位 (總共38位數,小數點後9位),你先設定「NUMERIC, BIGNUMERIC, STRING」:

  • 系統會先檢查 NUMERIC 是否能處理這個精確度 (這麼多位的數字)。
  • 結果後來發現有一個數字太大,NUMERIC 不足以處理 (因為 NUMERIC 最多 38 位),就會嘗試把欄位整個改成 BIGNUMERIC 繼續接收資料。
  • 結果數字又更大,或是非數字的資料出現,BIGNUMERIC 也不行,最後會轉成 STRING,幾乎什麼格式的內容都可以接受。

這樣做的好處是:

  • 避免資料傳輸中斷 – 不會因為突然出現意外的資料格式就失敗,不要全部重來太辛苦了。
  • 自動調適 – 不需要事先精確知道所有資料的格式和範圍,如果你懶得事先指定格式,就給一個順序讓 DTS 自動判斷就好。
  • 資料完整性 – 確保所有資料都能被正確保存,即使需要轉換成比較寬鬆的格式,寧可先收進來,格式後續再慢慢處理就好了。

Ignore Unknown Values

很容易理解,無法判斷的內容就直接略過,也可以避免因錯誤而中斷。

Header Rows to Skip

跟之前匯入 CSV 檔的操作一樣,如果第一列是標題就讓它略過。

45 傳輸選項設定
45 傳輸選項設定
資料來源:擷圖自 GCP 主控台

最後服務帳戶的部分,要指定一個讓它能有權限代替你執行整個傳輸工作,你不需要事先研究要設定給這個 Service Account 哪些權限,你直接指定一個帳戶給它,它會自動授予必要的權限角色給這個帳戶,當然你可以先產生一個不具備任何權限的服務帳戶,後面碰到權限問題再調整就好。

然後勾選「電子郵件通知」,這樣就不用時時進來確認進度。沒問題就按下「儲存」,它就會自動開始傳輸了。

46 Service Account 和通知設定
46 Service Account 和通知設定
資料來源:擷圖自 GCP 主控台

整個設定如上,除了你可以在 Console 設定,如果你有安裝指令套件,也可以下 bq 指令;或是寫程式呼叫 API,或操作 Java 的 Library 都可以執行傳輸工作。

更多設定的細節可以參考這份文件。如果是從 Azure Blob Storage 來傳輸,也可以參考這份文件

(三) 匯入串流資料

前兩種都是批次,代表一次性把資料匯入完成。而串流資料是持續不斷產生的 (例如社群留言、股票報價、使用者點擊遊戲或工廠 IoT 設備回傳),可能要即時處理和反應 (例如即時推薦商品或掉出寶物),就必須使用串流資料的傳輸方法,依照複雜程度分述如下:

1. BigQuery Storage Write API

這個方法針對小規模、簡單的串流需求,直接從應用程式寫入資料,並且不需要複雜的資料轉換,最適合這種方法。

BigQuery 早期主要使用 tabledata.insertAll API。這個 API 的設計相對簡單,主要特性包括:

(1) 使用 REST over HTTP 通訊協定。

(2) 資料格式採用 JSON。

(3) 提供基本的串流寫入功能。

(4) 操作直觀容易上手。

但它同時也存在一些明顯的問題:

(1) 效能限制

  • 使用 REST over HTTP 通訊協定,傳輸效率較低。
  • JSON 格式佔用較大的網路頻寬。
  • 屬於短期 Request,每個連接都有建立和關閉的動作,耗費系統資源且傳輸量受限。

(2) 資料一致性問題

  • 沒有完整的交易支援,例如有一部分失敗,不能全部倒回處理。
  • 無法保證資料的精確一次處理,可能出現資料重複的情況。

(3) 成本考量:

.資料傳輸成本較高。

.沒有免費用量額度。

Storage Write API 相對的優勢如下:

(1) 採用 gRPC 串流取代 REST over HTTP

雖然每個資料表最多 100 個連線,但這 100 個連線是持久的,可以持續使用,效能提昇不少。

(2) 資料一致性強化

它使用一個編號系統「串流偏移」(Stream Offset) 來追蹤每筆資料,當發現資料重複就不會傳送,提供「完全一次性」處理保證的選項。

(3) 成本效益

Storage Write API 提供每月 2 TiB 的免費擷取額度,而且 gRPC 串流能夠使用更少的資源來處理資料,提升資源的使用效率。

兩者比較總結如下表:

47 BigQuery tabledata.insertAll 與Storage Write API 比較
47 BigQuery tabledata.insertAll 與Storage Write API 比較
資料來源:東東老師自行整理

因此,如果需要高頻率的短期連線,可以使用 tabledata.insertAll API,如果需要穩定的長期資料串流以及資料一致性保證或交易支援,就選擇 Storage Write API。

2. Datastream

如果不寫程式,最好的方法就是使用 Datastream,它可以支援資料庫的即時同步,包含 MySQL、PostgreSQL、Oracle 等,如果來來源資料庫有變更,它也能同步到 BigQuery,也就是所謂的 CDC (Change Data Capture)。

設定步驟很簡單,首先進入 Datastream 來建立連線設定檔:

48 進入 Datastream 建立連線設定檔
48 進入 Datastream 建立連線設定檔
資料來源:擷圖自 GCP 主控台

我們真是 PostgreSQL 資料來為範例的資料來源,點擊 PostgreSQL:

49 選擇 PostgreSQL 做為 Datastream 來源
49 選擇 PostgreSQL 做為 Datastream 來源
資料來源:擷圖自 GCP 主控台

我這裡沒有現成的資料庫,以下使用國外網友的示範圖片,輸入連線設定檔名稱和 ID,並且指定 Region:

50 輸入 Datastream 連線設定檔名稱並指定 Region
50 輸入 Datastream 連線設定檔名稱並指定 Region
資料來源:擷圖自 YouTube 影片

接著輸入資料庫的 IP、Port、連線帳號、密碼和資料庫名稱: 

51 提供資料庫連的連線資訊
51 提供資料庫連的連線資訊
資料來源:擷圖自 YouTube 影片

它有提供三種連線方式,我們選擇最簡單的 IP 白名單,這裡指的是允許地端的資料庫可以從Datastream 的 IP 連線過去存取資料。 

52 選擇資料傳輸方式
52 選擇資料傳輸方式
資料來源:擷圖自 YouTube 影片

看到它跳出 IP 位址,直接複製起來貼到地端的防火牆去設白名單。

53 設定來源網路允許 Datastream 的 IP 去存取資料庫
53 設定來源網路允許 Datastream 的 IP 去存取資料庫
資料來源:擷圖自 YouTube 影片

接著要測試連線,這裡一定要測試成功,才可以建立設定檔,所以上方的欄位不能隨便亂打喔!

54 執行連線測試或直接建立
54 執行連線測試或直接建立
資料來源:擷圖自 YouTube 影片

若測試通過,就可以建立設定檔了。

55 測試連線成功後建立設定檔
55 測試連線成功後建立設定檔
資料來源:擷圖自 YouTube 影片

接著畫面會秀出,我們剛剛建好的連線設定檔內容,我們回到上一層: 

56 確認連線資訊並回到上一層
56 確認連線資訊並回到上一層
資料來源:擷圖自 YouTube 影片

我們要針對傳輸的目的地 BigQuery 也建立連線設定檔,點擊「Create Profile」: 

57 建立 BigQuery 的連線設定檔
57 建立 BigQuery 的連線設定檔
資料來源:擷圖自 YouTube 影片

選擇 BigQuery:

58 選擇 BigQuery 做為傳輸目的地
58 選擇 BigQuery 做為傳輸目的地
資料來源:擷圖自 YouTube 影片

設定名稱和 Region 並按下建立:

59 設定名稱和 Region 並按下建立
59 設定名稱和 Region 並按下建立
資料來源:擷圖自 YouTube 影片

接下來就看到它設定完成了:

60 確認 BigQuery 連線資訊
60 確認 BigQuery 連線資訊
資料來源:擷圖自 YouTube 影片

我們現在可以正式建立串流工作,點擊「Create Stream」: 

61 進入 Stream 選單建立串流
61 進入 Stream 選單建立串流
資料來源:擷圖自 YouTube 影片

輸入串流的名稱、ID、Region、來源跟目的地: 

62 設定 Stream 名稱、來源和目的地
62 設定 Stream 名稱、來源和目的地
資料來源:擷圖自 YouTube 影片

然後按「Continue」進行下一步:

63 確認無誤點擊繼續
63 確認無誤點擊繼續
資料來源:擷圖自 YouTube 影片

在這裡我們選擇剛剛建立的 PostgreSQL 連線設定檔: 

64 選擇剛建立的 PostgreSQL 連線設定檔
64 選擇剛建立的 PostgreSQL 連線設定檔
資料來源:擷圖自 YouTube 影片

在這裡我們一樣要測試一下資料庫的連線,成功後再按「Continue」:

65 連線測試成功並繼續
65 連線測試成功並繼續
資料來源:擷圖自 YouTube 影片

在 Configure Source 中,我們要設定 Publication Slot Name 和 Publication Name。

這個 Slot 不是 BigQuery 本身的 Slot 喔!在 PostgreSQL 的資料複寫 (Replication) 機制中,Replication Slot 是一個重要的概念。

它能夠追蹤資料變更,即使複寫目標 (Datastream) 暫時離線或延遲,Slot 也會確保需要的 WAL (Write-Ahead Logging) 檔案,直到 Datastream 收到這些變更,也能避免資料遺失。

你需要在 PostgreSQL 資料庫中預先建立這個 Replication Slot,然後在設定 Datastream 時提供這個 Slot 名稱,這樣 Datastream 就能透過這個 Slot 來追蹤和接收資料庫的變更。

而 Publication 是 PostgreSQL 10 版之後推出的邏輯複寫 (Logical Replication) 機制中的一個重要元素。你可以指定特定的表格要被複寫,以及要複寫哪些類型的操作。

66 設定 Publication Slot 和 Publication Name
66 設定 Publication Slot 和 Publication Name
資料來源:擷圖自 YouTube 影片

再往下選擇要轉移資料的表格跟欄位,選好再下一步: 

67 指定要同步資料的表格和欄位
67 指定要同步資料的表格和欄位
資料來源:擷圖自 YouTube 影片

在目的地的部分,我們就選擇剛剛建立的 BigQuery 連線設定檔: 

68 選擇 BigQuery 的設定檔並繼續
68 選擇 BigQuery 的設定檔並繼續
資料來源:擷圖自 YouTube 影片

這裡是問你要把轉移出來的資料要放在單一 Region 還是多個 Region 的 Dataset。如果你後續的的應用主要在特定區域運作,選擇 Region 即可,如果你需要更高的配額限制和更高的可用性,選擇 Multi-region。

69 指定資料要同步到 BigQuery 的 Region
資料來源:擷圖自 YouTube 影片

Stream Write Mode 指的是寫入模式,包含 Merge (合併模式,會更新舊資料)、Append (追加模式,不會更新或刪除現有資料) 和Replace (替換模式,會刪除所有現有資料)。

而 Staleness Limit (資料陳舊限制),指的是資料是否需要立即處理,設定為 0 秒表示資料會立即被處理,用來保持最即時的狀態。但可能會導致較高的查詢成本,因為每次變更都需要立即處理。

如果您想降低成本,並且可以接受一點資料延遲的話,可以增加一點 Staleness Limit。

70 指定寫入模式和資料陳舊限制
70 指定寫入模式和資料陳舊限制
資料來源:擷圖自 YouTube 影片

最後再按一下「RUN VALIDATION」來確認所有的設定正確無誤,如果都沒問題就可直接按下「CREATE & START」直接開始傳輸:

71 驗證所有設定並建立工作
71 驗證所有設定並建立工作
資料來源:擷圖自 YouTube 影片

它跳出一個確認視窗,再次按下「CREATE & START」:

72 確認建立 Datastream 工作
72 確認建立 Datastream 工作
資料來源:擷圖自 YouTube 影片

接著我們就看到 Datastream 開始同步資料了: 

73 確認 Datastream 開始同步資料
73 確認 Datastream 開始同步資料
資料來源:擷圖自 YouTube 影片

3. Pub/Sub 搭配 Cloud Functions

Pub/Sub 是一個完全代管的訊息佇列 (Message Queue) 服務,基於發布/訂閱 (Publish/Subscribe) 模式,主要用於事件驅動架構和串流資料處理。主要有三個特性:

(1) 解耦 (Decoupling)

透過發布/訂閱模式,發送方和接收方之間不需要直接互動,讓系統的各個部分能夠獨立發展和擴展,大幅提升了系統的靈活性和可維護性。

(2) 可靠性 (Reliability)

系統會將所有訊息永久儲存,即使在過程中發生異常,訊息也不會遺失。

(3) 擴展性 (Scalability)

系統能夠因應流量變化自動擴充,無需手動干預。

而 Cloud Functions 本身是輕量級的應用程式平台,內建訂閱 Pub/Sub Topic 的功能 很適合處理中小規模的串流資料,尤其是針對簡單的資料轉換。當新訊息到達時,能夠自動觸發 Cloud Function 執行,最適合用於事件驅動的架構,特別是那些資料流時有時無的場景。

優勢:

因為是無伺服器架構,不用管理任何基礎設施。系統能夠根據實際負載自動擴展,只要為實際使用的資源付費。開發過程也相對簡單,工程師可以專注於業務邏輯的實現,不必擔心底層基礎建設的管理。這種按需付費的模式通常能提供較好的成本效益。

限制:

Cloud Functions 的執行時間有上限,最長只能運行 540 秒,不適合需要長時間處理的任務。例如資料轉換太複雜的時候,並且冷啟動也會帶來延遲。另外對於併發處理能力也有一定的限制。

注意事項:

首先要設定合理的 Time-Out 時間,確保能夠完整處理資料。同時,必須實作穩健可靠的錯誤處理和重試機制,以應對可能的失敗情況。你也要建立完善的監控機制來追踪 Cloud Functions  的執行狀態。

4. Pub/Sub + Cloud Run

Cloud Run 結合 Pub/Sub 提供了一個更加靈活的串流資料處理方案。特別適合那些需要自訂處理邏輯,或需要特定運作環境的場景。與 Cloud Functions 相比,它支援更長的運作時間,最長可達到 60 分鐘,並且能夠處理中大型的串流資料。

優勢:

由於容器化的部署方式,這提供很大的靈活性。您可以使用任何程式語言,安裝任何所需的相依性套件。Cloud Run 能夠自動擴充 (Autoscale),同時能運作較長時間,讓它能夠適用於複雜的應用場景。

限制:

首先是容器映像檔的維護,這需要一個完整的容器管理或 CI/CD 流程。相比 Cloud Functions,Cloud Run 的設定也相對複雜一些,需要更多的維運工作。從成本角度來看,如果沒設定好 Autoscale 的參數,可能會導致較高的費用。

注意事項:

要注意容器映像檔的維護和更新策略。對資源的配置需要仔細優化,包括記憶體、CPU 等參數的設定。建立完善的監控和警報機制也很重要,需要及時發現和解決潛在問題。另外,為了控制成本也要小心設定 Autoscale 的策略。

5.  Pub/Sub + Dataflow

Dataflow 結合 Pub/Sub 是一個企業級的串流處理方案,適合處理大規模資料和複雜的轉換。Dataflow 具有高度的擴展性,特別是那些複雜的業務邏輯,或需要進行大規模數據處理的企業級應用場景。

優勢:

Dataflow 提供了豐富的資料轉換和處理功能,能夠處理複雜的業務邏輯。它也能 Autoscale,能夠自動處理負載量的變化。還內建完善的容錯機制,能夠確保數據處理的可靠性。此外,它還支援複雜的視窗操作 (Window Operations),將連續不斷的資料流切分成有限時間區段來處理,以及自訂的轉換邏輯。

限制:

Dataflow 也是最複雜的,它需要較多的專業知識,從設定到後續維護都很複雜。從成本角度來看,它是會自動建立虛擬機器來運作,也有較高的運營成本,通常需要專門的團隊和較長的時間投入。

注意事項:

首先是仔細規劃資料管道 (Data Pipeline) 的設計,確保能夠高效地處理資料流。成本優化也是一個重要考量,需要合理設置資源使用策略,從 CPU、記憶體、硬碟,到 Autoscale、處理窗格和管道設計等等。建立完善的監控機制是很重要的。

6.  串流方法比較表

以下整理五種串流資料的方法,你可以根據場景來決定處理方式:

74 各種串流資料處理方法比較表
74 各種串流資料處理方法比較表
資料來源:東東老師自行整理

文章轉載自《東東GCP 教學》網站

由專業團隊提供的全面技術學習與支援

歡迎您與我們聯絡
我們會協助您取得最佳解決方案!

歡迎您與我們聯絡
我們會協助您取得最佳解決方案!

思想科技 Master Concept
微信公众号:Master_Concept