未載入 Skill5485 tokens · 114100 ms
## BigQuery 匯入與合併設計
以下是規劃稿,未假設已登入 Google Cloud,也未執行任何雲端操作。由於 `orders.csv` 欄位未提供,以下以 `order_id` 為訂單唯一鍵、`updated_at` 為來源更新時間;若沒有這兩個欄位,必須先定義可靠的業務鍵與版本規則,否則無法可靠去重或判斷新舊資料。
### 1. 資料分層
| 層級 | 資源 | 用途 |
|---|---|---|
| 原始檔 | Cloud Storage `orders-raw` bucket | 原封不動保存收到的 CSV |
| 批次追蹤 | `ops.ingestion_batches` | 記錄每批狀態、錯誤、job ID、筆數 |
| 暫存 | `staging.orders_<batch_id>` | 以 `STRING` 欄位載入單一批次 |
| 原始列 | `raw.orders` | 永久保留可解析的原始列及批次資訊 |
| 拒絕列 | `raw.order_rejects` | 保存驗證失敗的列及原因 |
| 正式表 | `core.orders` | 每個 `order_id` 最多一筆有效訂單 |
Cloud Storage 原始檔是完整稽核來源;即使 CSV 結構損壞到 BigQuery 無法解析,檔案仍會保留。
### 2. 接收與批次識別
收到檔案後:
1. 對檔案原始位元組計算 SHA-256。
2. 使用雜湊作為 `batch_id`。
3. 上傳到不可變路徑,例如:
```text
gs://orders-raw/received_date=2026-09-03/<sha256>/orders.csv
```
4. 上傳時使用「僅在物件不存在時建立」的 generation precondition。
5. 若相同內容重送,SHA-256 與路徑相同,直接視為同一批,不再次處理。
6. 若檔名相同但內容不同,雜湊不同,因此會保存成另一批,不會覆蓋原檔。
可另外設定 bucket retention policy;是否鎖定需審慎決定,因為鎖定後可能不可逆。若允許同一路徑被覆寫,則至少啟用版本保留並把 `object_generation` 記入批次表。[Cloud Storage 物件版本說明](https://docs.cloud.google.com/storage/docs/object-versioning)
### 3. 批次狀態表
`ops.ingestion_batches` 建議至少包含:
```text
batch_id STRING
source_uri STRING
object_generation INT64
source_sha256 STRING
source_size INT64
received_at TIMESTAMP
status STRING
load_job_id STRING
merge_job_id STRING
source_row_count INT64
accepted_row_count INT64
rejected_row_count INT64
error_reason STRING
error_details JSON
started_at TIMESTAMP
finished_at TIMESTAMP
```
建議狀態:
```text
RECEIVED
LOADING
LOAD_FAILED
VALIDATING
VALIDATION_FAILED
MERGING
MERGE_FAILED
SUCCEEDED
DUPLICATE
```
狀態表以 `batch_id` 做邏輯唯一鍵;BigQuery 不應被當成會強制唯一鍵的 OLTP 資料庫,因此工作程序仍須使用確定性的 job ID 和單批次序列化處理。
### 4. 載入流程
1. 以確定性 job ID,例如 `orders_load_<sha256>`,將單一原始物件載入其批次暫存表。
2. 暫存欄位原則上先載成 `STRING`,後續使用 `SAFE_CAST` 驗證,避免一個日期或數字格式錯誤使整批無法保留。
3. CSV 欄數不符、引號未關閉等結構性錯誤仍可能使 load job 失敗;將 job 的 `errorResult`、`errors`、job ID 寫入批次表並標記 `LOAD_FAILED`。
4. 不要設定會靜默略過錯誤列的高 `max_bad_records`,否則可能無法完整說明哪些資料被遺漏。
5. 遇到逾時或網路結果不明時,不可直接換新 job ID 重跑。先以原 job ID 查詢狀態;BigQuery 對指定 job ID 的 `WRITE_APPEND` 重試可提供 at-most-once 保護。[BigQuery 重試可靠性](https://docs.cloud.google.com/bigquery/docs/reliability-intro)、[查詢 job 狀態](https://docs.cloud.google.com/bigquery/docs/running-jobs)
### 5. 驗證、保存與合併
先將所有成功解析的列寫入 `raw.orders`,並加上:
```text
batch_id
source_uri
object_generation
ingested_at
row_fingerprint
原始 CSV 各欄位(STRING)
```
`raw.orders` 建議依 `DATE(ingested_at)` 分割,並以 `batch_id, order_id` clustering。寫入時以 `batch_id` 防止整批重複;原始檔本身仍保留在 Cloud Storage。
驗證至少包含:
- `order_id` 不可為空。
- 日期、數量、金額可用 `SAFE_CAST` 成功轉型。
- 必填欄位存在。
- 金額、數量符合業務範圍。
- 同一批內相同 `order_id` 的處理規則明確。
不合格列寫入 `raw.order_rejects`,包含 `batch_id`、原始欄位及 `reject_reason`。建議預設採嚴格模式:只要有拒絕列,就將批次標為 `VALIDATION_FAILED`,不更新正式表;若業務允許部分成功,則另設 `SUCCEEDED_WITH_REJECTS`,不可把它標成完全成功。
合併前必須先讓來源中每個 `order_id` 最多一列,否則 BigQuery `MERGE` 在一個目標列對應多個來源列時可能失敗。[BigQuery MERGE 語法](https://docs.cloud.google.com/bigquery/docs/reference/standard-sql/dml-syntax)
示意 SQL:
```sql
CREATE TEMP TABLE deduplicated_source AS
SELECT * EXCEPT(rn)
FROM (
SELECT
SAFE_CAST(order_id AS STRING) AS order_id,
SAFE_CAST(updated_at AS TIMESTAMP) AS source_updated_at,
-- 其他 SAFE_CAST 後的正式欄位
batch_id,
ingested_at,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY
SAFE_CAST(updated_at AS TIMESTAMP) DESC,
ingested_at DESC,
batch_id DESC
) AS rn
FROM `PROJECT.raw.orders`
WHERE batch_id = @batch_id
)
WHERE rn = 1;
BEGIN TRANSACTION;
MERGE `PROJECT.core.orders` AS T
USING deduplicated_source AS S
ON T.order_id = S.order_id
WHEN MATCHED
AND S.source_updated_at > T.source_updated_at
THEN UPDATE SET
source_updated_at = S.source_updated_at,
last_batch_id = S.batch_id
-- 其他欄位
WHEN NOT MATCHED THEN
INSERT (
order_id,
source_updated_at,
last_batch_id
-- 其他欄位
)
VALUES (
S.order_id,
S.source_updated_at,
S.batch_id
-- 其他欄位
);
UPDATE `PROJECT.ops.ingestion_batches`
SET
status = 'SUCCEEDED',
finished_at = CURRENT_TIMESTAMP()
WHERE batch_id = @batch_id;
COMMIT TRANSACTION;
```
`MERGE` 和成功狀態更新放在同一個 transaction,使兩者一起提交或回滾;BigQuery 支援交易內的 DML 與原子提交。[BigQuery 多陳述式交易](https://docs.cloud.google.com/bigquery/docs/transactions)
如果訂單沒有可信的 `updated_at`:
- 若訂單不可修改:正式表只使用 `WHEN NOT MATCHED THEN INSERT`。
- 若訂單可修改:必須由來源提供版本號、事件時間或其他明確排序依據;不能以檔案到達時間推測業務資料的新舊。
### 6. 失敗與重跑規則
- `LOAD_FAILED`:修正檔案或解析設定後,以新的 attempt/job ID 重跑,但沿用相同 `batch_id`。
- `VALIDATION_FAILED`:原始資料不動;修正轉換規則後重跑驗證。
- `MERGE_FAILED`:重跑相同批次;transaction 可避免正式表更新成功但批次狀態未更新。
- job 回應不明:先用既有 job ID 查詢,不立即建立新 job。
- `SUCCEEDED` 或相同 SHA-256 重送:標記為 `DUPLICATE` 或直接回報既有成功結果,不再次合併。
- 錯誤通知只需包含 `batch_id` 和控制表連結,避免在通知中散布訂單敏感資料。
工作程序查詢自己建立的 job 時,建立 job 後會自動取得該 job 的 `get` 權限;因此不必為了失敗追蹤授予可查看全專案 job 的 `roles/bigquery.resourceViewer`。[BigQuery job 權限](https://docs.cloud.google.com/bigquery/docs/running-jobs)
### 7. 最小必要權限
應使用專用 runtime service account,並將「一次性建置權限」和「每日執行權限」分開。
#### 檔案上傳者
在指定 raw bucket:
```text
roles/storage.objectCreator
```
它只能建立物件,不能讀取、刪除或覆寫既有物件,適合不可變落地。[Cloud Storage IAM roles](https://docs.cloud.google.com/storage/docs/access-control/iam-roles)
#### 每日匯入/合併 service account
專案層級:
```text
roles/bigquery.jobUser
```
提供 `bigquery.jobs.create`。
指定 Cloud Storage bucket:
```text
roles/storage.objectViewer
```
若使用精確物件 URI,自訂角色可縮到:
```text
storage.buckets.get
storage.objects.get
```
只有使用 URI wildcard 或列舉檔案時才需要:
```text
storage.objects.list
```
指定 BigQuery datasets:
```text
roles/bigquery.dataEditor
```
應只授予 `staging`、`raw`、`core`、`ops` 這些必要 dataset,不授予整個專案。若組織要求更嚴格,可建立自訂角色,依實際 SQL 和預先建立的資料表包含:
```text
bigquery.tables.create
bigquery.tables.get
bigquery.tables.getData
bigquery.tables.update
bigquery.tables.updateData
```
BigQuery 官方列出的 load job 權限包括 `bigquery.tables.create`、`bigquery.tables.updateData`、`bigquery.tables.update` 與 `bigquery.jobs.create`。[BigQuery 批次載入權限](https://docs.cloud.google.com/bigquery/docs/batch-loading-data)
#### 一次性建置者
由管理者建立 bucket、retention policy、datasets、tables、IAM 與排程。建置完成後移除較高權限;每日 runtime 不需要:
```text
roles/owner
roles/editor
roles/bigquery.admin
roles/storage.admin
```
若使用 Cloud Scheduler、Cloud Run、Workflows 或 Eventarc,還需加入該執行架構的最小呼叫權限,但這取決於實際選用的服務,不能在尚未決定架構時宣稱為必要權限。
已載入 Skill +7366 tokens · 148027 ms
# BigQuery 每日 `orders.csv` 匯入與合併規劃
> 以下為設計方案;尚未登入 Google Cloud,也未建立資源、上傳檔案或執行任何 BigQuery 工作。
## 1. 核心假設
- `order_id` 是訂單的穩定且唯一業務鍵。若來源沒有可靠的 `order_id`,必須先與來源系統定義複合鍵,否則無法保證不重複。
- 若同一訂單可能更新,來源應提供 `updated_at` 或版本號。
- 如果沒有版本欄位,建議採「首次寫入為準」,將內容不同的同鍵訂單列為衝突,不要任意覆寫。
## 2. 資料層次
| 層次 | 建議資源 | 用途 |
|---|---|---|
| 原始檔 | `gs://BUCKET/orders/.../<sha256>.csv` | 保存收到的原始位元組,不修改、不覆寫 |
| 批次控制 | `ops.ingestion_batches` | 追蹤批次狀態、工作 ID、筆數與錯誤 |
| 到件紀錄 | `ops.delivery_events` | 每次收到檔案都留紀錄,包括被判定為重送者 |
| 暫存 | `staging.orders_<batch_id>` | 以明確 schema 載入當批 CSV |
| 原始列 | `raw.orders` | 保存成功解析的原始欄位及批次資訊 |
| 正式表 | `curated.orders` | 每個 `order_id` 最多一筆目前版本 |
| 異常列 | `ops.order_rejections` | 保存型別或業務規則不合格的列 |
Cloud Storage bucket 和 BigQuery dataset 應放在相容的區域;從 Cloud Storage 載入時,兩者需位於相同 regional 或 multi-regional location。[Google Cloud:Loading CSV data](https://docs.cloud.google.com/bigquery/docs/loading-data-cloud-storage-csv)
## 3. 批次識別與防重送
1. 收到檔案後先計算完整檔案的 SHA-256:
```text
batch_id = lowercase(hex(sha256(file_bytes)))
```
2. 每次到件都新增一筆 `ops.delivery_events`,記錄:
- `delivery_id`
- `batch_id`
- 原始檔名
- `received_at`
- Cloud Storage URI 與 object generation
- `action`:`PROCESS`、`SKIP_DUPLICATE` 或 `RETRY`
3. 以 `batch_id` 查詢 `ops.ingestion_batches`:
- 已為 `SUCCEEDED`:不再載入或合併,記為 `SKIP_DUPLICATE`。
- 為 `FAILED`:沿用同一 `batch_id` 重試。
- 不存在:建立 `RECEIVED` 批次。
- 相同檔名但內容不同:SHA-256 不同,視為新批次。
4. 所有 BigQuery load/query job 都使用由 `batch_id` 推導的固定 job ID,例如:
```text
orders_load_<batch_id>
orders_merge_<batch_id>
```
API 呼叫逾時後重試時,先查詢相同 job ID 的既有結果,不重新提交另一份工作。這可避免「工作其實成功,但呼叫端沒收到回應」造成重複處理。
5. 流程同一時間只允許一個執行者處理特定 `batch_id`。如果可能有多個 worker,應在具備條件寫入能力的協調層取得鎖;不能依賴 BigQuery 的 `PRIMARY KEY ... NOT ENFORCED` 作為唯一性約束。
## 4. 匯入與合併流程
### A. 保存原始檔
將檔案寫入不可覆寫的路徑,例如:
```text
gs://BUCKET/orders/received_date=2026-09-03/<batch_id>/orders.csv
```
建議:
- 上傳時使用 object generation precondition,避免意外覆寫。
- 啟用 Object Versioning 或 retention policy。
- 處理失敗也不刪除原始檔。
- 若法規要求保存每一次實際到件,而不只是每種內容,可為每個 `delivery_id` 保存獨立物件。
### B. 建立批次控制紀錄
`ops.ingestion_batches` 至少包含:
```text
batch_id, source_uri, object_generation, file_sha256,
status, failed_step, load_job_id, merge_job_id,
received_at, started_at, completed_at,
source_row_count, valid_row_count, rejected_row_count,
inserted_count, updated_count, duplicate_count,
retry_count, error_reason, error_message
```
狀態流程:
```text
RECEIVED → LOADING → LOADED → VALIDATING → MERGING → SUCCEEDED
└──────── 任一步驟失敗 ────────→ FAILED
```
失敗狀態必須由協調程式在捕捉到 load/query job 錯誤後,以另一個獨立查詢寫回;不能只在已失敗的交易內記錄錯誤。另保存 Google Cloud 的 job ID,方便從 BigQuery Job History 或 Audit Logs 追查。
### C. 嚴格載入暫存表
- 使用明確 schema,不使用 autodetect。
- 首次解析時,容易出錯的來源欄位先載為 `STRING`,後續用 `SAFE_CAST` 驗證。
- `max_bad_records=0`。
- 不允許多餘欄位或缺欄。
- 暫存表名稱由 `batch_id` 決定,並使用 `WRITE_EMPTY`,避免重試時追加第二份資料。
- 暫存表可設定 7~30 天到期,但原始 Cloud Storage 物件與批次紀錄不得隨之刪除。
結構性 CSV 錯誤會使整個 load job 失敗;原始檔仍保存在 Cloud Storage,批次記為 `FAILED/LOAD`。
### D. 保存解析後原始列
載入成功後,將該批所有列原樣寫入 `raw.orders`,並補上:
```text
batch_id, source_uri, file_sha256, ingested_at
```
此表保留來源字串值,不以正式表型別覆蓋。寫入工作使用固定 job ID,且只有批次尚未標記 `raw_loaded` 時才執行,避免重試造成第二份原始列。
### E. 驗證及批內去重
用 `SAFE_CAST` 建立驗證結果,例如:
```sql
SELECT
order_id,
SAFE_CAST(order_date AS DATE) AS order_date,
SAFE_CAST(amount AS NUMERIC) AS amount,
SAFE_CAST(updated_at AS TIMESTAMP) AS source_updated_at
FROM `PROJECT.staging.orders_BATCH_ID`;
```
以下任一情況寫入 `ops.order_rejections`:
- `order_id` 為空。
- 日期、金額或時間無法轉型。
- 金額、狀態等違反業務規則。
- 同一 `order_id`、同一版本卻有不同內容。
建議採全批次原子策略:只要存在 rejected row,就不修改正式表,批次標為 `FAILED/VALIDATION`。若業務接受部分成功,才改為隔離錯誤列並合併有效列,而且批次狀態應使用 `SUCCEEDED_WITH_REJECTIONS`。
同一檔案內重複的 `order_id` 必須先縮成一列,否則 `MERGE` 可能因一筆目標列對到多筆來源列而失敗:
```sql
CREATE TEMP TABLE batch_dedup AS
SELECT * EXCEPT(rn)
FROM (
SELECT
v.*,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY source_updated_at DESC, row_fingerprint DESC
) AS rn
FROM validated_rows AS v
)
WHERE rn = 1;
```
`row_fingerprint` 應由正規化後的完整列計算,讓同一版本的選擇結果具決定性。
### F. 原子合併正式表
正式表建議包含:
```text
order_id, order_date, amount, status,
source_updated_at, source_batch_id,
created_at, last_merged_at
```
在同一個 BigQuery transaction 中執行 `MERGE` 和批次成功狀態更新:
```sql
BEGIN TRANSACTION;
MERGE `PROJECT.curated.orders` AS T
USING batch_dedup AS S
ON T.order_id = S.order_id
WHEN MATCHED
AND S.source_updated_at > T.source_updated_at
THEN UPDATE SET
order_date = S.order_date,
amount = S.amount,
status = S.status,
source_updated_at = S.source_updated_at,
source_batch_id = @batch_id,
last_merged_at = CURRENT_TIMESTAMP()
WHEN NOT MATCHED THEN
INSERT (
order_id, order_date, amount, status,
source_updated_at, source_batch_id,
created_at, last_merged_at
)
VALUES (
S.order_id, S.order_date, S.amount, S.status,
S.source_updated_at, @batch_id,
CURRENT_TIMESTAMP(), CURRENT_TIMESTAMP()
);
UPDATE `PROJECT.ops.ingestion_batches`
SET
status = 'SUCCEEDED',
completed_at = CURRENT_TIMESTAMP()
WHERE batch_id = @batch_id;
COMMIT TRANSACTION;
```
BigQuery multi-statement transaction 可將 DML 變更一起提交或回滾;若查詢在提交前失敗,交易會回滾。[Google Cloud:Multi-statement transactions](https://docs.cloud.google.com/bigquery/docs/transactions)
如果來源沒有 `updated_at`,將 `WHEN MATCHED` 改成不更新;內容不同者寫入衝突表供人工處理。
另外應每日檢查:
```sql
SELECT order_id, COUNT(*) AS n
FROM `PROJECT.curated.orders`
GROUP BY order_id
HAVING n > 1;
```
結果必須為零;BigQuery 宣告的非強制主鍵不能取代這項檢查。
## 5. 失敗與重試規則
| 失敗點 | 正式表是否變更 | 重試方式 |
|---|---:|---|
| 上傳失敗 | 否 | 使用相同檔案雜湊重新上傳 |
| CSV load 失敗 | 否 | 修正設定或來源後,沿用批次重新提交 |
| 驗證失敗 | 否 | 原始檔及 raw 資料保留;修正資料後以新內容建立新批次 |
| `MERGE` 失敗 | 否 | 交易回滾;查明原因後沿用相同 job/batch 規則重試 |
| 狀態回寫失敗 | 正式表狀態不確定 | 先查 merge job ID 與正式表的 `source_batch_id`,不可直接重跑 |
`error_reason` 應保存結構化 API reason;`error_message` 只供人閱讀,不要用解析錯誤文字的方式決定程式流程。
## 6. 最小必要權限
建議分離「部署身分」與「每日執行服務帳號」。每日執行帳號不應擁有 Project Editor、Owner 或 BigQuery Admin。
### 每日執行服務帳號
BigQuery:
- 在執行 job 的 project 授予 `roles/bigquery.jobUser`,提供 `bigquery.jobs.create`。
- 在 `staging`、`raw`、`curated`、`ops` datasets 授予 `roles/bigquery.dataEditor`,這是較容易維護的預先定義角色方案。
- 更嚴格時建立 custom role,僅包含:
- `bigquery.tables.create`:建立每批暫存表。
- `bigquery.tables.update`:load job 更新暫存表。
- `bigquery.tables.updateData`:載入及寫入資料。
- `bigquery.tables.getData`:讀取暫存、raw、正式及控制表。
- `bigquery.tables.get`:讀取必要的表中繼資料。
- 若所有暫存資源均由部署流程預先建立,可移除執行帳號的 `bigquery.tables.create`。
- 不授予 dataset 建立、刪除、IAM 修改或專案管理權限。
CSV load 官方列出的 BigQuery 必要權限為 `bigquery.tables.create`、`bigquery.tables.updateData`、`bigquery.tables.update` 與 `bigquery.jobs.create`。[Google Cloud:CSV load required permissions](https://docs.cloud.google.com/bigquery/docs/loading-data-cloud-storage-csv)
Cloud Storage:
- 若服務帳號只讀取已上傳的檔案,授予 bucket 層級 `roles/storage.objectViewer`;不授予整個 project。
- 最嚴格的 custom role:
- `storage.buckets.get`
- `storage.objects.get`
- 只有使用 wildcard URI 時才加 `storage.objects.list`
- 若同一服務帳號負責保存新原始檔,再於該 bucket 授予 `roles/storage.objectCreator`;不要授予物件刪除或覆寫權限。
BigQuery 從 Cloud Storage 載入所需的精確 Storage 權限為 `storage.buckets.get`、`storage.objects.get`,使用 wildcard 時另需 `storage.objects.list`。[Google Cloud:Cloud Storage load permissions](https://docs.cloud.google.com/bigquery/docs/loading-data-cloud-storage-csv)
### 部署/管理身分
僅在部署時使用,負責:
- 建立 bucket、datasets、tables、service account。
- 設定 IAM、retention policy、排程與監控。
- 啟用必要 API。
這些管理權限不應授予每日執行服務帳號。若人員需要查看失敗批次,可只對 `ops` dataset 授予 `roles/bigquery.dataViewer`;需要查看 Cloud Logging 時另授予 `roles/logging.viewer`。