Writing to GCP(google cloud) lakehouse runtime catalog iceberg tables with dbt

dbtでbigquery native tableからlakehouse runtime catalogのiceberg tableへ書き込む
gcp
bigquery
dbt
iceberg
ai_generated
Author

Masaya Kameyama

Published

August 24, 2026

Introduction

この記事は検証結果をもとに生成AIで作成した. dbtでlakehouse runtime catalogへ書き込みができるところまでは検証したが残念ポイントその3は未検証である.

前回, GCPのlakehouse runtime catalogとsnowflakeの統合について書いた. その中で残念ポイントその2として, bigquery側のデータをlakehouse runtime catalogのテーブルへ書き込めないため, 連携したいデータをカタログ側のテーブルへコピーする追加のワークフローが必要になると述べた. その後この部分はプレビュー機能として解消しつつあり, bigqueryのDDL/DMLでlakehouse runtime catalogのiceberg tableを直接操作できるようになった.

そこで, dbtで中間層まではbigquery native tableとして変換し, 最終層(gold層)だけをlakehouse runtime catalogのiceberg tableとして書き出すパイプラインを実際に作って検証した. 結論としては実現できたが, dbtの標準機能では届かず, カスタムマテリアライゼーションを書く必要があった. この記事ではその方法と, 実際に動いた範囲, そして新たに見つかった残念ポイントをまとめる. 例によって変更の激しい領域なので, 2026年8月現在の情報であることに注意してほしい.

bigqueryからlakehouse runtime catalogのテーブルを操作する

lakehouse runtime catalogのテーブルをbigqueryから操作する場合, 4パート識別子が要求される. 前回紹介したmanaged iceberg tableの構文と比較すると, かなり簡潔である:

CREATE TABLE `PROJECT_ID.CATALOG_ID.NAMESPACE.TABLE_NAME` (id int, data string);
INSERT INTO `PROJECT_ID.CATALOG_ID.NAMESPACE.TABLE_NAME` VALUES (1, "foo");

managed iceberg tableで必要だったWITH CONNECTIONOPTIONS(file_format, table_format, storage_uri)も不要である. データの保存先はカタログ側が決めるためで, CATALOG_IDはカタログのIDを指す. catalog_typeCATALOG_TYPE_GCS_BUCKETを使った場合はバケット名がそのままカタログIDになるので, ハイフンを含む名前もバッククォートの内側にそのまま書けばよい. UPDATE / DELETE / MERGEも同じ識別子で通る.

読み書きの可否はテーブルプロパティgcp.biglake.bigquery-dml.enabledで決まる. bigqueryのDDLで作成したテーブルはこれが既定で有効になるため, 作成直後からbigqueryのDMLが使える. 一方でsparkなどのOSSエンジンが作成したテーブルは既定で無効なので, bigqueryからはSELECTしかできず, DMLを投げると

DML statements are only supported over tables that have data stored in BigQuery.

というエラーになる. この場合はTBLPROPERTIESまたはALTER TABLEで明示的にオプトインすればbigqueryからのDMLが通る. つまり前回書いた「iceberg tableは作成した側のサービスからしか書き込めない」という制約は, カタログ側のテーブルに関してはプロパティで解除できるようになったということである.

なお, この機能に関するドキュメントは記述が分散しており, 古いset upページには今でも「bigqueryのDDL/DMLでは作成も変更もできない」という記述が残っている. 実挙動は上記のとおりなので, 迷ったら実機で確認するのが早い.

terraform側の準備

カタログとnamespaceはterraformで作れる. bigqueryのCREATE TABLEはnamespaceが存在している前提なので, テーブルを作る前にnamespaceを用意しておく必要がある:

resource "google_storage_bucket" "lakehouse" {
  name                        = "my-lakehouse-bucket"
  location                    = "us-central1"
  force_destroy               = true
  uniform_bucket_level_access = true
}

resource "google_biglake_iceberg_catalog" "main" {
  name             = google_storage_bucket.lakehouse.name
  catalog_type     = "CATALOG_TYPE_GCS_BUCKET"
  credential_mode  = "CREDENTIAL_MODE_VENDED_CREDENTIALS"
  primary_location = "us-central1"
}

resource "google_biglake_iceberg_namespace" "gold" {
  catalog      = google_biglake_iceberg_catalog.main.name
  namespace_id = "gold"
}

# vended credentialsではカタログのサービスエージェントが
# 利用者に代わってGCS上のmetadata/dataを読み書きするため, この権限が必要
resource "google_storage_bucket_iam_member" "catalog_agent" {
  bucket = google_storage_bucket.lakehouse.name
  role   = "roles/storage.objectAdmin"
  member = "serviceAccount:${google_biglake_iceberg_catalog.main.biglake_service_account}"
}

最後のIAMを忘れると, テーブル作成時にstorage.objects.create deniedの403で落ちる. vended credentialsモードではクエリエンジンが自分の権限でGCSに書くのではなく, カタログのサービスエージェントが発行する短命のトークンを使ってGCSにアクセスするためである.

dbtの標準機能では届かない

dbtはバージョン1.10前後からcatalogs.ymlによる外部カタログのサポートを入れており, dbt-bigqueryにもcatalog_type: biglake_metastoreが用意されている. しかしこれで作られるのはbigquery managed iceberg tableである. アダプタの実装を読むと, iceberg指定時に付与されるのはtable_format / file_format / storage_uriのOPTIONSであり, 識別子は通常の3パートproject.dataset.tableのままだった. catalog integrationの種類もbiglake_metastoreの1種類しかなく, 既定のカタログ名はmanaged_icebergになっている.

より根本的な問題として, dbt coreのRelationはdatabase.schema.identifierの3パートまでしか表現できない. lakehouse runtime catalogのテーブルはproject, catalog, namespace, tableの4つを単一のバッククォートで囲む必要があるため, 標準のマテリアライゼーションでは構文上どうしても届かない.

カスタムマテリアライゼーション

そこでRelationを経由せず, configから4パート識別子を組み立ててDDL/DMLを直接発行するマテリアライゼーションを書く. 実装は次のようになる:

{% macro lakehouse_iceberg_merge_sql(table_ref, source_sql, column_names, unique_keys) %}
  {%- set update_columns = column_names | reject('in', unique_keys) | list -%}
  merge into {{ table_ref }} as dbt_dest
  using (
    {{ source_sql }}
  ) as dbt_src
  on
    {%- for key in unique_keys %}
    {{ 'and' if not loop.first }} dbt_dest.`{ key }` = dbt_src.`{ key }`
    {%- endfor %}
  {% if update_columns | length > 0 -%}
  when matched then update set
    {%- for col in update_columns %}
    `{ col }` = dbt_src.`{ col }`{{ ',' if not loop.last }}
    {%- endfor %}
  {% endif -%}
  when not matched then insert
    ({% for col in column_names %}`{ col }`{{ ', ' if not loop.last }}{% endfor %})
  values
    ({% for col in column_names %}dbt_src.`{ col }`{{ ', ' if not loop.last }}{% endfor %})
{% endmacro %}


{% materialization lakehouse_iceberg, adapter='bigquery' %}

  {%- set identifier = model['alias'] -%}
  {%- set catalog = var('iceberg_catalog') -%}
  {%- set namespace = var('iceberg_namespace') -%}
  {%- set table_ref = '`' ~ target.project ~ '.' ~ catalog ~ '.' ~ namespace ~ '.' ~ identifier ~ '`' -%}

  {%- set unique_key = config.get('unique_key') -%}
  {%- set unique_keys = [unique_key] if unique_key is string else (unique_key or []) -%}
  {%- set full_refresh_mode = should_full_refresh() -%}

  {{ run_hooks(pre_hooks) }}

  {#- 4パート識別子はRelationではないためget_columns_in_relationが使えない.
      compiled SQLを空サブクエリ化して列名と型を取得する -#}
  {%- set columns = adapter.get_column_schema_from_query(get_empty_subquery_sql(sql)) -%}
  {%- set column_names = columns | map(attribute='name') | list -%}
  {%- set column_ddl = [] -%}
  {%- for c in columns -%}
    {%- do column_ddl.append('`' ~ c.name ~ '` ' ~ c.data_type) -%}
  {%- endfor -%}

  {% if full_refresh_mode %}
    {% do run_query('drop table if exists ' ~ table_ref) %}
  {% endif %}

  {% do run_query('create table if not exists ' ~ table_ref ~ ' (' ~ column_ddl | join(', ') ~ ')') %}

  {% call statement('main') %}
    {% if unique_keys | length > 0 and not full_refresh_mode %}
      {{ lakehouse_iceberg_merge_sql(table_ref, sql, column_names, unique_keys) }}
    {% else %}
      insert into {{ table_ref }}
        ({% for col in column_names %}`{ col }`{{ ', ' if not loop.last }}{% endfor %})
      select {% for col in column_names %}`{ col }`{{ ', ' if not loop.last }}{% endfor %}
      from (
        {{ sql }}
      )
    {% endif %}
  {% endcall %}

  {{ run_hooks(post_hooks) }}

  {#- 4パートのiceberg tableはrelation cacheで表現できないため空を返す -#}
  {{ return({'relations': []}) }}

{% endmaterialization %}

ポイントは3つある.

ひとつめは列定義の取得方法である. 4パート識別子はRelationではないのでadapter.get_columns_in_relationが使えない. 代わりにget_empty_subquery_sqlでcompiled SQLをwhere false limit 0のサブクエリに包み, adapter.get_column_schema_from_queryに渡すことで列名とbigqueryの型を得ている. これはdbtがcontract検証で使っているのと同じ手である.

ふたつめは書き込みの分岐である. 通常のrunではCREATE TABLE IF NOT EXISTSでテーブルを用意してからMERGEを投げる. 初回はテーブルが空なので結果的に全件INSERTになり, 2回目以降は差分更新になる. --full-refreshのときはDROPCREATEINSERT ... SELECTに切り替える.

みっつめはreturn({'relations': []})である. 作ったテーブルをdbtのrelation cacheに載せないことで, 存在しない3パートのrelationをdbtが参照しようとして壊れるのを避けている. 代償として後述の制約が付く.

モデル側はマテリアライゼーションを指定するだけでよい:

{{ config(
    materialized='lakehouse_iceberg',
    schema='gold',
    unique_key=['event_date', 'tenant_id']
) }}

select
    event_date,
    tenant_id,
    impressions,
    clicks,
    current_timestamp() as loaded_at
from {{ ref('int_events_daily') }}

実際に動いたこと

公式ドキュメントに載っているのは列指定のCREATE TABLEINSERT ... VALUESだけなので, dbtが必要とする操作を個別に切り分けて試した. 結果は次のとおりで, いずれも未文書ながら動作した:

操作 結果
CREATE TABLE(列指定) ⭕ ドキュメント記載どおり
INSERT ... VALUES ⭕ ドキュメント記載どおり
INSERT ... SELECT(native table → iceberg table) ⭕ 未文書
MERGE(native tableをUSINGに取る) ⭕ 未文書
CREATE TABLE AS SELECT ⭕ 未文書

native tableとiceberg tableをまたぐINSERT ... SELECTMERGEが通ることが, この方式が成立するかどうかの分岐点だった. ここが駄目ならINSERT ... VALUESしか書き込み手段がなくなり, dbtから使うのは現実的ではなくなる.

dbt側の挙動も期待どおりだった. loaded_atcurrent_timestamp()を入れておくと差分更新の確認がしやすい:

  • 1回目のrun: CREATE TABLE + MERGEでテーブルが作られる
  • 2回目のrun: 行数が変わらずloaded_atだけが進む. すなわち追記ではなくMERGEとして既存行が更新されている
  • --full-refresh: dbtのログがMERGEからINSERTに変わり, GCS上のテーブルディレクトリが別のIDに切り替わる

書き込み先が本当にlakehouse runtime catalog側のテーブルなのかは, 独立した3つの観測で確認した.

  1. gold層に対応するbigquery datasetのテーブル一覧がである. managed iceberg tableならdataset側に見えるはずである
  2. iceberg REST APIのテーブル一覧に出てくる
GET https://biglake.googleapis.com/iceberg/v1/restcatalog/v1/projects/PROJECT_ID/catalogs/CATALOG_ID/namespaces/NAMESPACE/tables

{ "identifiers": [ { "namespace": ["gold"], "name": "gold_events_daily" } ] }
  1. icebergのクライアントからスナップショットと実データが読める

pyicebergで読む場合は以下のようになる. sparkやtrinoを立てるよりも手軽に相互運用性を確認できる:

import google.auth
import google.auth.transport.requests
from pyiceberg.catalog.rest import RestCatalog

creds, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"])
creds.refresh(google.auth.transport.requests.Request())

catalog = RestCatalog("main", **{
    "uri": "https://biglake.googleapis.com/iceberg/v1/restcatalog",
    # バケット名だけを渡すと400になる
    "warehouse": "gs://my-lakehouse-bucket",
    "token": creds.token,
    "header.x-goog-user-project": "PROJECT_ID",
    "header.X-Iceberg-Access-Delegation": "vended-credentials",
    # RestCatalog既定のOAuthフローを回避しADCのトークンをそのまま使う
    "rest.auth.type": "noop",
})

t = catalog.load_table("gold.gold_events_daily")
print(t.current_snapshot())
print(t.scan().to_arrow())

pyiceberg[gcsfs]だけではFileIOが初期化できずload_tableが失敗するので, pyiceberg[pyarrow,gcsfs]を入れる必要がある1.

読み取ったloaded_atはdbtのrun時刻と一致した. EXPORT TABLE METADATAを挟まずに, dbtがbigqueryのエンジン経由で書いた結果が外部エンジンからそのまま見えているということである. 前回の記事で問題にしていたexportとテーブル作り直しのループが, 最終層をカタログ側のテーブルにするだけで消えたことになる.

残念ポイントその3

とはいえ制約は残る. まずマテリアライゼーション自体に内在するものとして,

  • 終端専用になる. relation cacheに載せないので, このモデルを下流のref()で参照できない. したがってiceberg table同士を連結したリネージをdbtで組むことはできず, あくまでgold層の出口としてしか使えない
  • dbt testdbt docsが使えない. これらも3パートのrelationを前提にしているためである. データ品質のチェックはsilver層までで済ませ, 出力されたiceberg tableは4パート識別子への直接クエリで確認することになる
  • 列の追加が反映されない. CREATE TABLE IF NOT EXISTSで作っているので, スキーマを変えたら--full-refreshが必要になる

次に, プラットフォーム側の挙動として,

  • DROP TABLEしてもGCSのファイルは残る. カタログの登録は消えるが, metadataとparquetはバケットに残ったままだった. --full-refreshを繰り返すとゴミが蓄積していくので, ライフサイクルルールか定期的な掃除を前提にした方がよい
  • native tableとiceberg tableを同一クエリで扱うにはロケーションを揃える必要がある. bigqueryはUSマルチリージョンとus-central1を別ロケーションとして扱い, これをまたぐクエリを拒否する. silver層のdatasetとカタログのバケットのロケーションが食い違っていると, そもそもINSERT ... SELECTが組めない
  • プレビュー機能への依存. 上で動作したもののうち3つは公式ドキュメントに記載がない. GAまでの間に挙動が変わる可能性は当然ある

そして最大の残念ポイントは, これがdbtの標準機能ではないことそのものである. dbt-bigqueryのRelationが4パートを扱えるようになるか, catalogs.ymlにlakehouse runtime catalogを指すcatalog typeが追加されれば, このマテリアライゼーションは不要になる. それまでは自前で保守する必要がある.

Summary

  • 前回残念ポイントその2として挙げた「bigqueryからlakehouse runtime catalogのテーブルへ書き込めない」問題は, プレビュー機能として解消しつつある. `PROJECT_ID.CATALOG_ID.NAMESPACE.TABLE_NAME`の4パート識別子を使えば, WITH CONNECTIONOPTIONSもなしにDDL/DMLが通る.
  • 書き込みの可否はテーブルプロパティgcp.biglake.bigquery-dml.enabledで決まり, bigqueryが作成したテーブルは既定で有効, OSSエンジンが作成したテーブルはオプトインで有効化できる.
  • dbt-bigqueryのcatalogs.yml(catalog_type: biglake_metastore)ではmanaged iceberg tableしか作れず, dbt coreのRelationも3パートまでのため, 標準機能ではlakehouse runtime catalogのテーブルを作れない.
  • 4パート識別子を直接組み立てるカスタムマテリアライゼーションを書けば, 中間層はbigquery native tableで変換し, 最終層だけをlakehouse runtime catalogのiceberg tableとして書き出すパイプラインが成立する. 列定義はget_column_schema_from_queryget_empty_subquery_sqlで取得し, 通常runはMERGE, --full-refreshDROPCREATEINSERT ... SELECTとした.
  • 公式ドキュメント未記載のINSERT ... SELECT, MERGE(native sourceをUSING), CREATE TABLE AS SELECTはいずれも動作した. 出力されたテーブルはbigquery dataset側には現れず, iceberg REST APIとpyicebergから参照できる. EXPORT TABLE METADATAは不要である.
  • ただしiceberg tableは終端専用(下流ref()不可)になり, dbt test / dbt docsも使えない. DROP TABLEでGCSのファイルが残る点と, native tableとのロケーション一致が必要な点にも注意が必要である.
Back to top

Footnotes

  1. 余談だが, ユーザーのADCにはquota_project_idが埋め込まれており, これが全APIリクエストの課金先プロジェクトになる. ここが無効なプロジェクトを指していると, どのバケットを触ってもUserProjectAccountProblemの403が返る. gcloud config set projectでは変わらないのでgcloud auth application-default set-quota-projectで更新する.↩︎