%sql
SELECT
current_catalog() AS 今のカタログ, -- データの置き場所(大分類)
current_schema() AS 今のスキーマ, -- その中の仕切り
current_user() AS 自分, -- ログインしている自分
current_timestamp() AS 今;
%sql
-- 誰に何を許すか
GRANT SELECT ON TABLE learn.sales.orders TO `analyst@example.com`;
-- 今どんな権限が付いているか
SHOW GRANTS ON TABLE learn.sales.orders;
-- 取り消す
REVOKE SELECT ON TABLE learn.sales.orders FROM `analyst@example.com`;
%sql
SELECT
tpep_pickup_datetime AS 乗車時刻,
trip_distance AS 距離,
fare_amount AS 料金,
pickup_zip AS 乗車地の郵便番号
FROM samples.nyctaxi.trips
WHERE fare_amount >= 50 -- 料金が50以上の行だけ
ORDER BY fare_amount DESC -- 料金の大きい順(DESC=降順)
LIMIT 10;
NULLは「=」では比べられない 空っぽ(NULL)は「値がない」という状態なので、= NULL と書いても永久に一致しません。必ず IS NULL を使います。「なぜか0件になる」の原因の半分はこれです。
列に計算をさせる
SELECTには、既存の列だけでなく計算式も書けます。名前を付けるには AS を使います。
calc.sql
%sql
SELECT
trip_distance AS 距離,
fare_amount AS 料金,
round(fare_amount / trip_distance, 1) AS 単価, -- 1マイルあたり
CASE
WHEN trip_distance < 1 THEN '近距離'
WHEN trip_distance < 5 THEN '中距離'
ELSE '長距離'
END AS 距離区分
FROM samples.nyctaxi.trips
WHERE trip_distance > 0
LIMIT 20;
CASE WHEN は「もし〜なら」を書く仕組みで、Excelの IF 関数にあたります。上から順に判定し、最初に当てはまったところで止まります。実務でいちばん出番の多い構文なので、ここで慣れておいてください。
よく使う関数
やりたいこと
関数
例
四捨五入
round(値, 桁)
round(3.14159, 2) → 3.14
空欄を別の値に
coalesce(a, b)
coalesce(備考, 'なし')
文字をつなぐ
concat(a, b) / a || b
concat(姓, ' ', 名)
大文字・小文字
upper() / lower()
lower(メール)
前後の空白を消す
trim()
trim(会社名)
区切って取り出す
split_part(文字, 区切り, 何番目)
split_part(メール, '@', 2)
型を変える
cast(x AS INT) / try_cast()
try_cast(金額 AS DOUBLE)
try_ が付いた関数を覚えておくcast('あいう' AS INT) は失敗してエラーになりますが、try_cast は失敗するとNULLを返して処理を続けます。汚れたデータを扱う実務では、こちらのほうが役に立つ場面が多いです。
%sql
SELECT
tpep_pickup_datetime AS 乗車時刻,
trip_distance AS 距離,
fare_amount AS 料金,
round(fare_amount / trip_distance, 1) AS 単価
FROM samples.nyctaxi.trips
WHERE trip_distance >= 10
AND fare_amount IS NOT NULL
ORDER BY trip_distance DESC
LIMIT 5;
※ ORDER BY は必ず WHERE より後ろ、LIMIT は最後です。この順番は決まっていて、入れ替えると構文エラーになります。
集計とは、たくさんの行をある単位でまとめて、1行にすることです。「月ごとの売上」なら単位は月、「店舗ごとの客数」なら単位は店舗です。この単位を書くのが GROUP BY です。
group_by.sql
%sql
SELECT
pickup_zip AS 乗車地,
count(*) AS 件数,
round(avg(fare_amount), 1) AS 平均料金,
round(sum(fare_amount), 0) AS 合計料金
FROM samples.nyctaxi.trips
GROUP BY pickup_zip -- この単位でまとめる
HAVING count(*) >= 100 -- まとめた後の絞り込み
ORDER BY 合計料金 DESC
LIMIT 10;
%sql
SELECT
date_trunc('MONTH', tpep_pickup_datetime) AS 年月,
count(*) AS 件数,
round(sum(fare_amount), 0) AS 売上
FROM samples.nyctaxi.trips
GROUP BY 年月 -- ALL と書いてもよい(後述)
ORDER BY 年月;
Databricksの便利な書き方:GROUP BY ALL 集計関数ではない列を全部まとめの単位にする、という指定です。列を書き写す手間が消え、書き間違いも減ります。SELECT ... GROUP BY ALL と書くだけです。他のデータベースには無いことが多いので、他社製品に移すコードでは通常の書き方にしてください。
%sql
-- ① セッションのタイムゾーンを日本にする(そのノートブックの間だけ有効)
SET TIME ZONE 'Asia/Tokyo';
-- ② または、変換関数で明示する(こちらが確実)
SELECT
tpep_pickup_datetime AS utcの時刻,
from_utc_timestamp(tpep_pickup_datetime, 'Asia/Tokyo') AS 日本時間,
date_trunc('DAY', from_utc_timestamp(tpep_pickup_datetime, 'Asia/Tokyo')) AS 日本時間の日付
FROM samples.nyctaxi.trips
LIMIT 5;
%sql
SELECT * FROM (
SELECT
date_format(o_orderdate, 'yyyy-MM') AS 年月,
o_orderpriority AS 優先度,
o_totalprice AS 金額
FROM samples.tpch.orders
WHERE o_orderdate >= DATE'1998-01-01'
)
PIVOT (
round(sum(金額)/1000, 0)
FOR 優先度 IN ('1-URGENT', '2-HIGH', '3-MEDIUM')
)
ORDER BY 年月;
%sql
SELECT
c.c_name AS 顧客名,
o.o_orderdate AS 注文日,
o.o_totalprice AS 金額
FROM samples.tpch.orders AS o -- 左のテーブル
JOIN samples.tpch.customer AS c -- 右のテーブル
ON o.o_custkey = c.c_custkey -- つなぐ鍵(キー)
WHERE o.o_orderdate >= DATE'1998-01-01'
ORDER BY o.o_totalprice DESC
LIMIT 10;
AS o のように短い別名を付けておくと、以降 o.列名 と書けます。テーブルが3つ4つになるとこれが効いてきます。
JOINの種類
種類
結果に残る行
使う場面
INNER JOIN(既定)
両方にある行だけ
注文と、その顧客情報。確実に対応があるとき
LEFT JOIN
左は全部残る。右に無ければNULL
いちばん使う。「注文していない顧客も含めた一覧」
RIGHT JOIN
右が全部残る
LEFTで書き直せるので、ほぼ使わない
FULL OUTER JOIN
両方の全部
2つの名簿の突き合わせ
ANTI JOIN
右に無い左の行だけ
「まだ買っていない人」を探す
left_join.sql
%sql
-- 注文が1件も無い顧客も含めて、顧客ごとの注文件数を出す
SELECT
c.c_name AS 顧客名,
count(o.o_orderkey) AS 注文件数, -- ← count(*) にしない
coalesce(sum(o.o_totalprice), 0) AS 売上
FROM samples.tpch.customer AS c
LEFT JOIN samples.tpch.orders AS o
ON c.c_custkey = o.o_custkey
GROUP BY ALL
ORDER BY 注文件数 ASC
LIMIT 10;
LEFT JOINで count(*) を使わない 注文が無い顧客も、LEFT JOINでは1行(右側が全部NULL)として残ります。count(*) はそれを1件と数えてしまい、注文0件の人が「1件」になります。数えたい対象の列を指定して count(o.o_orderkey) と書けば、NULLは数えられず正しく0になります。
%sql
WITH 対象注文 AS (
-- ① まず期間で絞る
SELECT o_custkey, o_orderkey, o_totalprice
FROM samples.tpch.orders
WHERE o_orderdate >= DATE'1998-01-01'
),
顧客別 AS (
-- ② 顧客ごとにまとめる
SELECT
o_custkey,
count(*) AS 件数,
sum(o_totalprice) AS 売上
FROM 対象注文
GROUP BY ALL
)
-- ③ 最後に顧客名を付ける
SELECT
c.c_name AS 顧客名,
k.件数,
round(k.売上, 0) AS 売上
FROM 顧客別 AS k
JOIN samples.tpch.customer AS c ON k.o_custkey = c.c_custkey
ORDER BY k.売上 DESC
LIMIT 10;
%sql
WITH 明細数 AS (
SELECT
l_orderkey,
count(*) AS 明細数
FROM samples.tpch.lineitem
GROUP BY ALL
)
SELECT
round(avg(明細数), 2) AS 平均明細数,
max(明細数) AS 最大明細数,
count(*) AS 注文数
FROM 明細数;
%sql
-- 顧客ごとの売上サマリを作る
CREATE OR REPLACE TABLE customer_summary AS
SELECT
customer AS 顧客,
count(*) AS 件数,
sum(amount) AS 売上,
max(order_date) AS 最終注文日
FROM orders
GROUP BY ALL;
SELECT * FROM customer_summary;
CREATE OR REPLACE は安全 中身を入れ替えても、テーブルの履歴(Time Travel)は残ります。したがって、間違えて古いデータで作り直しても前の版に戻せます。DROP TABLE してから作り直すより、こちらを使ってください。
%sql
-- 今日届いた更新データ(ふだんは取り込んだテーブルを使う)
CREATE OR REPLACE TEMP VIEW 本日分 AS
SELECT * FROM VALUES
(2, '鈴木物産', DATE'2026-08-02', 9000, 'paid'), -- 金額が変わった
(5, '佐藤電機', DATE'2026-08-14', 22000, 'new') -- 新規
AS t(order_id, customer, order_date, amount, status);
MERGE INTO orders AS 元
USING 本日分 AS 新
ON 元.order_id = 新.order_id -- 何をもって「同じ行」とするか
WHEN MATCHED THEN UPDATE SET
元.amount = 新.amount,
元.status = 新.status,
元.updated_at = current_timestamp()
WHEN NOT MATCHED THEN INSERT
(order_id, customer, order_date, amount, status)
VALUES (新.order_id, 新.customer, 新.order_date, 新.amount, 新.status);
SELECT * FROM orders ORDER BY order_id;
実行結果
order_id customer order_date amount status
1 山田商店 2026-08-01 12000.00 paid
2 鈴木物産 2026-08-02 9000.00 paid ← 更新された
3 田中工業 2026-08-02 31500.00 new
5 佐藤電機 2026-08-14 22000.00 new ← 追加された
① learn.sales に customers テーブル(列:customer、area)を作り、3件ほどデータを入れてください。② 注文テーブルと結合して、地域ごとの売上を出してください。③ 同じ顧客のデータをもう一度MERGEしても件数が増えないことを確認してください。
解答を見る
answer.sql
%sql
CREATE OR REPLACE TABLE customers (customer STRING, area STRING);
INSERT INTO customers VALUES
('山田商店', '関東'), ('鈴木物産', '関西'), ('田中工業', '関東'), ('佐藤電機', '中部');
-- ② 地域ごとの売上
SELECT
c.area AS 地域,
count(*) AS 件数,
sum(o.amount) AS 売上
FROM orders AS o
JOIN customers AS c ON o.customer = c.customer
GROUP BY ALL
ORDER BY 売上 DESC;
-- ③ 同じデータをもう一度MERGEしても増えない(べき等)
MERGE INTO customers AS 元
USING (SELECT '山田商店' AS customer, '関東' AS area) AS 新
ON 元.customer = 新.customer
WHEN MATCHED THEN UPDATE SET 元.area = 新.area
WHEN NOT MATCHED THEN INSERT *;
SELECT count(*) FROM customers; -- 4件のまま
※ 「何度実行しても結果が同じ」ことをべき等(idempotent)といいます。処理が途中で失敗しても、そのまま流し直せば復旧できるため、自動化する処理はすべてべき等に作るのが原則です。INSERT INTO ではなく MERGE を使う理由がこれです。
%sql
CREATE OR REPLACE TABLE orders_silver AS
SELECT
cast(order_id AS BIGINT) AS order_id,
nullif(trim(customer), '') AS customer, -- 空文字はNULLに
try_to_date(order_date, 'yyyy-MM-dd') AS order_date,
try_cast(amount AS DECIMAL(12,2)) AS amount, -- 数値でなければNULL
lower(trim(status)) AS status,
_file_name,
_loaded_at
FROM orders_bronze
WHERE order_id IS NOT NULL;
-- 品質チェック:変換に失敗した行を数える
SELECT
count(*) AS 全行,
count(*) FILTER (WHERE customer IS NULL) AS 顧客名なし,
count(*) FILTER (WHERE amount IS NULL) AS 金額が変,
count(*) FILTER (WHERE order_date IS NULL) AS 日付が変
FROM orders_silver;
① 新しいCSVをもう1つボリュームに置き、同じ COPY INTO を実行して、増えた分だけ取り込まれることを確認してください。② orders_silver で、金額の変換に失敗した行の元データ(ファイル名つき)を表示してください。
解答を見る
answer.sql
%sql
-- ② 失敗した行を、元の文字列と一緒に確認する
SELECT
b.order_id,
b.amount AS 元の値,
b._file_name,
b._loaded_at
FROM orders_bronze AS b
LEFT JOIN orders_silver AS s ON cast(b.order_id AS BIGINT) = s.order_id
WHERE s.amount IS NULL
AND b.amount IS NOT NULL;
# ① PythonからSQLを実行して、結果をDataFrameで受け取る
df = spark.sql("""
SELECT pickup_zip, count(*) AS 件数
FROM samples.nyctaxi.trips
GROUP BY ALL
""")
display(df)
# ② SQLに値を渡す(文字列を直接つなげない。安全な渡し方)
zip_code = 10001
df2 = spark.sql(
"SELECT * FROM samples.nyctaxi.trips WHERE pickup_zip = :z LIMIT 5",
args={"z": zip_code},
)
display(df2)
# ③ DataFrameに名前を付けて、SQLセルから見えるようにする
df.createOrReplaceTempView("zip_summary")
SQLセルから使う
%sql
-- ③で作った一時ビューを、SQLセルからそのまま使える
SELECT * FROM zip_summary ORDER BY 件数 DESC LIMIT 5;
文字列をつなげてSQLを作らないf"... WHERE id = {入力}" のような書き方は、入力に妙な文字が混ざると意図しないSQLが実行されます(SQLインジェクション)。上の②のように args で渡すか、IDENTIFIER() を使ってください。社内データでも例外にしないこと。
# Sparkで集計してから、小さくなった結果だけpandasに渡す(正しい使い方)
small = spark.sql("SELECT pickup_zip, count(*) AS n FROM samples.nyctaxi.trips GROUP BY ALL")
pdf = small.toPandas()
pdf.plot(kind="bar", x="pickup_zip", y="n", figsize=(10, 4))
%sql
SELECT
event_id,
payload:user.id::BIGINT AS ユーザーid, -- : で潜り、:: で型を決める
payload:user.name::STRING AS 名前,
payload:action::STRING AS 行動,
payload:amount::DECIMAL(12,2) AS 金額,
payload:items AS 明細 -- 配列のまま
FROM events;
%sql
SELECT
e.event_id,
e.payload:user.name::STRING AS 名前,
item:sku::STRING AS 商品,
item:qty::INT AS 数量
FROM events AS e,
LATERAL variant_explode(e.payload:items) AS t(pos, item)
WHERE e.payload:action::STRING = 'purchase';
GROUP BY は行をまとめてしまうため、明細が消えます。明細を残したまま、順位や前月比を計算したいときに使うのがウィンドウ関数です。実務のSQLで、これが書けるかどうかは大きな差になります。
window.sql
%sql
WITH 月別 AS (
SELECT
date_format(o_orderdate, 'yyyy-MM') AS 年月,
sum(o_totalprice) AS 売上
FROM samples.tpch.orders
WHERE o_orderdate >= DATE'1998-01-01'
GROUP BY ALL
)
SELECT
年月,
round(売上) AS 売上,
round(lag(売上) OVER (ORDER BY 年月)) AS 前月,
round(売上 - lag(売上) OVER (ORDER BY 年月)) AS 増減,
round(sum(売上) OVER (ORDER BY 年月)) AS 累計,
round(avg(売上) OVER (ORDER BY 年月 ROWS BETWEEN 2 PRECEDING AND CURRENT ROW)) AS 三か月移動平均
FROM 月別
ORDER BY 年月;
関数
何が出るか
使いどころ
row_number()
1,2,3…(同点でも別番号)
最新1件の抽出、重複排除
rank() / dense_rank()
順位(同点は同順位)
売上ランキング
lag() / lead()
前の行 / 次の行の値
前月比、次回来店までの日数
sum() OVER (…)
累計
累計売上、在庫の推移
avg() OVER (… ROWS …)
移動平均
でこぼこをならして傾向を見る
PARTITION BY は「グループごとにやり直す」OVER (PARTITION BY 顧客 ORDER BY 日付) と書くと、顧客が変わるたびに番号が1に戻ります。「顧客ごとの初回購入日」「店舗ごとの売上順位」のように、実務で必要になるのはほぼこの形です。
%sql
-- 顧客ごとの最新注文だけを残す
SELECT
o_custkey AS 顧客,
o_orderkey AS 注文番号,
o_orderdate AS 注文日,
o_totalprice AS 金額
FROM samples.tpch.orders
QUALIFY row_number() OVER (PARTITION BY o_custkey ORDER BY o_orderdate DESC, o_orderkey DESC) = 1
ORDER BY 顧客
LIMIT 10;
これが「重複排除」の正体 STEP 5のMERGEで「取り込んだデータに同じキーが複数あると失敗する」と書きました。その解決策がこれです。QUALIFY row_number() OVER (PARTITION BY キー ORDER BY 更新時刻 DESC) = 1 で最新の1件に絞ってからMERGEすれば、安全に反映できます。この1行は、そのまま実務で使えます。覚えてください。
ORDER BYに「同点の決着」を入れる 更新時刻が同じ行が2つあると、どちらが残るかが実行のたびに変わり、結果が毎回違うという最悪の不具合になります。上の例のように、日付の後ろに o_orderkey DESC のような一意の列を足して、必ず順序が決まるようにしてください。
%sql
WITH 顧客別 AS (
SELECT
o_custkey AS 顧客,
sum(o_totalprice) AS 売上,
min(o_orderdate) AS 初回,
max(o_orderdate) AS 最終
FROM samples.tpch.orders
GROUP BY ALL
)
SELECT
rank() OVER (ORDER BY 売上 DESC) AS 順位,
顧客,
round(売上) AS 売上,
初回,
最終,
datediff(最終, 初回) AS 取引期間_日
FROM 顧客別
QUALIFY 順位 <= 10
ORDER BY 順位;
※ QUALIFY は、ウィンドウ関数の結果をそのまま条件にできる便利な句です。これが無い場合は、いったんCTEに入れてから WHERE 順位 <= 10 と書くことになります。
%sql
-- ① バージョン番号を指定して読む
SELECT * FROM learn.sales.orders VERSION AS OF 1;
-- ② 時刻を指定して読む
SELECT * FROM learn.sales.orders TIMESTAMP AS OF '2026-08-14 10:15:00';
-- ③ 今と過去の差分を見る(何が変わったのか)
SELECT * FROM learn.sales.orders VERSION AS OF 1
EXCEPT
SELECT * FROM learn.sales.orders;
-- ④ 本当に戻す
RESTORE TABLE learn.sales.orders TO VERSION AS OF 1;
%sql
-- ① 小さいファイルをまとめる
OPTIMIZE learn.sales.orders;
-- ② よく検索に使う列で、データを整列させておく(推奨のやり方)
ALTER TABLE learn.sales.orders CLUSTER BY (order_date, customer);
OPTIMIZE learn.sales.orders;
-- ③ 統計とファイル数の確認
DESCRIBE DETAIL learn.sales.orders;
手段
何をするか
いつ使うか
OPTIMIZE
小さなファイルを適切な大きさにまとめる
追記が続いたテーブル
リキッドクラスタリング CLUSTER BY
よく絞り込む列でデータを並べておく
いまはこれが基本。後から列を変えられる
ZORDER BY
同上(古い方式)
既存のテーブルで使われていれば維持
パーティション PARTITIONED BY
日付などでフォルダを物理的に分ける
1TB超の巨大テーブルだけ。小さいテーブルでやると逆に遅くなる
予測的最適化
Databricksが自動でOPTIMIZEとVACUUMを実行
有効にできるなら有効に。手間が消える
小さなテーブルをパーティションで分けない 「日付ごとに分ければ速そう」と考えて PARTITIONED BY (日付) にすると、1日あたり数十行しかない場合、極小ファイルが大量にできて劇的に遅くなります。よくある失敗です。迷ったらパーティションは使わず、CLUSTER BY にしてください。
① learn.sales.orders の全行の金額を、わざと UPDATE orders SET amount = 0; で壊してください。② 履歴を確認し、③ 壊す直前のバージョンとの差分を確認してから、④ 元に戻してください。
解答を見る
answer.sql
%sql
-- ① 事故を起こす(WHEREを忘れた想定)
UPDATE learn.sales.orders SET amount = 0;
-- ② 履歴を見る。いちばん上が事故のバージョン
DESCRIBE HISTORY learn.sales.orders;
-- ③ 直前(事故のversion - 1)と比べる
SELECT * FROM learn.sales.orders VERSION AS OF 5 -- ← 自分の番号に読み替える
EXCEPT
SELECT * FROM learn.sales.orders;
-- ④ 戻す
RESTORE TABLE learn.sales.orders TO VERSION AS OF 5;
SELECT * FROM learn.sales.orders ORDER BY order_id;
なぜ分けるのか ① 取り込みで失敗しても生データが残っているので、何度でもやり直せる ② 「どこで数字が変わったか」を層ごとに追える ③ 分析者は整ったsilver以降だけを見ればよく、生データの汚れに振り回されない。これは規模の大小に関係なく効きます。個人の練習でも、最初からこの3層で作ってください。
%sql
CREATE TABLE IF NOT EXISTS learn.sales.orders_silver (
order_id BIGINT,
customer STRING,
order_date DATE,
amount DECIMAL(12,2),
status STRING,
_file_name STRING,
_loaded_at TIMESTAMP
);
MERGE INTO learn.sales.orders_silver AS s
USING (
SELECT
cast(order_id AS BIGINT) AS order_id,
nullif(trim(customer), '') AS customer,
try_to_date(order_date, 'yyyy-MM-dd') AS order_date,
try_cast(amount AS DECIMAL(12,2)) AS amount,
lower(trim(status)) AS status,
_file_name,
_loaded_at
FROM learn.sales.orders_bronze
WHERE order_id IS NOT NULL
-- 同じ注文番号が複数あれば、いちばん新しく取り込んだ1件だけ残す
QUALIFY row_number() OVER (
PARTITION BY cast(order_id AS BIGINT)
ORDER BY _loaded_at DESC, _file_name DESC
) = 1
) AS b
ON s.order_id = b.order_id
WHEN MATCHED AND b._loaded_at > s._loaded_at THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;
SELECT count(*) AS silver件数 FROM learn.sales.orders_silver;
UPDATE SET * と INSERT * は、「列名が同じものを全部」という省略記法です。列が多いテーブルで重宝します。WHEN MATCHED AND b._loaded_at > s._loaded_at という条件を入れているのは、古いデータで新しいデータを上書きしてしまう事故を防ぐためです。
gold層を作る
silver_to_gold.sql
%sql
CREATE OR REPLACE TABLE learn.sales.monthly_sales_gold AS
SELECT
date_trunc('MONTH', order_date) AS 年月,
count(*) AS 件数,
count(DISTINCT customer) AS 顧客数,
sum(amount) AS 売上,
round(avg(amount), 0) AS 平均単価
FROM learn.sales.orders_silver
WHERE status <> 'canceled'
AND order_date IS NOT NULL
GROUP BY ALL
ORDER BY 年月;
SELECT * FROM learn.sales.monthly_sales_gold;
goldは「作り直す」でよい 集計表は元データから何度でも作り直せるので、CREATE OR REPLACE TABLE で丸ごと作り直すのがいちばん簡単で、間違いも起きません。データが大きくなって時間がかかるようになってから、差分更新を考えれば十分です。最初から難しく作らないこと。
%sql
-- ① 制約を付けて、違反データが入らないようにする
ALTER TABLE learn.sales.orders_silver
ADD CONSTRAINT amount_not_negative CHECK (amount IS NULL OR amount >= 0);
-- ② 毎日確認する数字を1本のSQLで出す
SELECT
current_date() AS 検査日,
count(*) AS 件数,
count(*) FILTER (WHERE customer IS NULL) AS 顧客名なし,
count(*) FILTER (WHERE amount IS NULL) AS 金額なし,
count(*) - count(DISTINCT order_id) AS 重複件数,
max(_loaded_at) AS 最終取り込み
FROM learn.sales.orders_silver;
もう一段進んだ自動化の仕組みがあります。Lakeflow宣言的パイプライン(以前の名称は Delta Live Tables / DLT)といい、「どういうテーブルを作りたいか」だけを書けば、実行の順番・差分処理・再試行・品質チェックをDatabricksが引き受けてくれます。
pipeline.sql(パイプライン用のノートブック)
-- ① 生データを取り込む(ファイルが増えた分だけ自動で処理される)
CREATE OR REFRESH STREAMING TABLE orders_bronze
AS SELECT *, _metadata.file_name AS _file_name, current_timestamp() AS _loaded_at
FROM STREAM read_files(
'/Volumes/learn/sales/files/',
format => 'csv',
header => true
);
-- ② 整える。品質ルールを付けられる
CREATE OR REFRESH STREAMING TABLE orders_silver (
CONSTRAINT 注文番号あり EXPECT (order_id IS NOT NULL) ON VIOLATION DROP ROW,
CONSTRAINT 金額が正 EXPECT (amount >= 0)
)
AS SELECT
cast(order_id AS BIGINT) AS order_id,
nullif(trim(customer), '') AS customer,
try_to_date(order_date) AS order_date,
try_cast(amount AS DECIMAL(12,2)) AS amount,
lower(trim(status)) AS status,
_file_name, _loaded_at
FROM STREAM(orders_bronze);
-- ③ 集計する
CREATE OR REFRESH MATERIALIZED VIEW monthly_sales_gold
AS SELECT date_trunc('MONTH', order_date) AS 年月,
count(*) AS 件数, sum(amount) AS 売上
FROM orders_silver
WHERE status <> 'canceled'
GROUP BY ALL;
%sql
-- 3段階すべてに「使ってよい」が必要(ここが最初の関門)
GRANT USE CATALOG ON CATALOG learn TO `data_analyst`;
GRANT USE SCHEMA ON SCHEMA learn.sales TO `data_analyst`;
GRANT SELECT ON SCHEMA learn.sales TO `data_analyst`; -- 中の全テーブル
-- 特定のテーブルだけ渡す場合
GRANT SELECT ON TABLE learn.sales.monthly_sales_gold TO `sales_team`;
-- 書き込みも許す(データ担当だけ)
GRANT MODIFY, SELECT ON SCHEMA learn.sales TO `data_engineer`;
-- 確認する
SHOW GRANTS ON SCHEMA learn.sales;
SHOW GRANTS `data_analyst` ON CATALOG learn;
-- 取り消す
REVOKE SELECT ON TABLE learn.sales.orders_bronze FROM `data_analyst`;
「権限がありません」の9割はこれ テーブルに SELECT を付けただけでは読めません。その上のカタログとスキーマに USE が要ります。3階層すべてを通す、と覚えてください。ビューを使う場合は、ビュー自体に権限があれば、元テーブルの権限は不要です(これがビューを使う理由のひとつです)。
%sql
-- ① 列マスク:特定のグループ以外にはぼかして見せる
CREATE OR REPLACE FUNCTION learn.sales.mask_email(email STRING)
RETURN CASE
WHEN is_account_group_member('pii_reader') THEN email
ELSE '***@***'
END;
ALTER TABLE learn.sales.customers
ALTER COLUMN email SET MASK learn.sales.mask_email;
-- ② 行フィルタ:見てよい行だけに絞る
CREATE OR REPLACE FUNCTION learn.sales.area_filter(area STRING)
RETURN is_account_group_member('all_area_reader') OR area = current_user_area();
ALTER TABLE learn.sales.customers
SET ROW FILTER learn.sales.area_filter ON (area);
簡易な方法:ビューで隠す マスク機能が使えない環境では、必要な列だけのビューを作り、ビューにだけ権限を渡す方法があります。CREATE VIEW customers_safe AS SELECT id, name, area FROM customers; のようにして、元テーブルの権限は誰にも渡しません。素朴ですが確実です。
タグとリネージ ── どこから来て、どこへ行くのか
tag.sql
%sql
-- 個人情報を含む列に印を付ける(後で一括で探せる)
ALTER TABLE learn.sales.customers
ALTER COLUMN email SET TAGS ('pii' = 'true');
ALTER SCHEMA learn.sales SET TAGS ('owner_team' = 'データ基盤チーム');
-- 印の付いた列を全社から探す
SELECT * FROM system.information_schema.column_tags
WHERE tag_name = 'pii';
%sql
-- 監査ログ(有効化されている環境で参照できる)
SELECT
event_time,
user_identity.email AS 実行者,
action_name,
request_params
FROM system.access.audit
WHERE event_date >= current_date() - INTERVAL 7 DAYS
AND action_name IN ('getTable', 'generateTemporaryTableCredential')
ORDER BY event_time DESC
LIMIT 50;
あなたのチームに、次の3種類の人がいます。それぞれに必要な権限を GRANT 文で書いてください。①データ基盤チーム(全部)②分析担当(goldだけ読める)③営業部(月次売上テーブルだけ読める)。
解答を見る
answer.sql
%sql
-- 前提:スキーマを役割で分けておく(learn.bronze / learn.silver / learn.gold)
-- ① データ基盤チーム
GRANT ALL PRIVILEGES ON CATALOG learn TO `data_engineer`;
-- ② 分析担当:goldだけ
GRANT USE CATALOG ON CATALOG learn TO `data_analyst`;
GRANT USE SCHEMA ON SCHEMA learn.gold TO `data_analyst`;
GRANT SELECT ON SCHEMA learn.gold TO `data_analyst`;
-- ③ 営業部:1テーブルだけ
GRANT USE CATALOG ON CATALOG learn TO `sales_team`;
GRANT USE SCHEMA ON SCHEMA learn.gold TO `sales_team`;
GRANT SELECT ON TABLE learn.gold.monthly_sales TO `sales_team`;
-- 確認
SHOW GRANTS `sales_team` ON CATALOG learn;
%sql
-- 直近30日の、種類別・日別のDBU消費
SELECT
usage_date AS 日付,
sku_name AS 種類,
round(sum(usage_quantity), 2) AS DBU
FROM system.billing.usage
WHERE usage_date >= current_date() - INTERVAL 30 DAYS
GROUP BY ALL
ORDER BY 日付 DESC, DBU DESC;
cost_by_job.sql
%sql
-- 何にいくらかかっているか(概算の金額つき)
SELECT
u.usage_metadata.job_id AS ジョブid,
u.sku_name AS 種類,
round(sum(u.usage_quantity), 1) AS DBU,
round(sum(u.usage_quantity * p.pricing.effective_list.default), 2) AS 概算ドル
FROM system.billing.usage AS u
JOIN system.billing.list_prices AS p
ON u.sku_name = p.sku_name
AND u.usage_end_time >= p.price_start_time
AND (p.price_end_time IS NULL OR u.usage_end_time < p.price_end_time)
WHERE u.usage_date >= current_date() - INTERVAL 30 DAYS
GROUP BY ALL
ORDER BY 概算ドル DESC
LIMIT 20;
%sql
SELECT
executed_by AS 実行者,
round(total_duration_ms / 1000, 1) AS 秒,
round(read_bytes / 1024 / 1024 / 1024, 2) AS 読んだGB,
read_rows AS 読んだ行数,
left(statement_text, 120) AS クエリ
FROM system.query.history
WHERE start_time >= current_date() - INTERVAL 7 DAYS
AND total_duration_ms > 60000 -- 1分以上かかったもの
ORDER BY total_duration_ms DESC
LIMIT 20;
① 直近7日で、いちばんDBUを使っている sku_name は何かを調べてください。② 自分が実行したクエリのうち、いちばん時間がかかったものを見つけ、Query Profileを開いて「読んだデータ量」を確認してください。
解答を見る
answer.sql
%sql
-- ①
SELECT sku_name, round(sum(usage_quantity), 2) AS DBU
FROM system.billing.usage
WHERE usage_date >= current_date() - INTERVAL 7 DAYS
GROUP BY ALL
ORDER BY DBU DESC;
-- ② 自分のクエリだけを対象にする
SELECT statement_id, round(total_duration_ms/1000,1) AS 秒,
round(read_bytes/1024/1024,1) AS 読んだMB, left(statement_text, 100) AS クエリ
FROM system.query.history
WHERE executed_by = current_user()
ORDER BY total_duration_ms DESC
LIMIT 5;
from databricks import sql
import os
with sql.connect(
server_hostname = os.environ["DATABRICKS_HOST"],
http_path = os.environ["DATABRICKS_HTTP_PATH"], # SQLウェアハウスの接続先
access_token = os.environ["DATABRICKS_TOKEN"], # コードに直接書かない
) as conn:
with conn.cursor() as cur:
cur.execute("SELECT * FROM learn.sales.monthly_sales_gold ORDER BY 年月")
for row in cur.fetchall():
print(row)
請求で失敗しないための4か条 ① すべてのコンピュートに短い自動停止を設定する ② 定期処理はジョブコンピュートで動かす ③ コンピュートポリシーと予算アラートを設定する ④ 月に1回 system.billing.usage を見る。この4つだけで、Databricksの「思わぬ高額請求」のほぼすべてを防げます。
3か月の学習プラン例
1日1時間、週5日で進めた場合の目安
時期
やること
その週の到達目標
1週目
STEP 0〜1
Free Editionを作り、自分のカタログとスキーマを持つ
2〜3週目
STEP 2〜3
SELECTと集計で、月別・分類別の表を作れる
4〜5週目
STEP 4〜5
JOINとCTEが書ける。MERGEでテーブルを更新できる
6〜7週目
STEP 6
手元のCSVを取り込んで集計できる
8〜9週目
STEP 7
PySparkで同じ処理が書ける。SQLとの使い分けが分かる
10週目
STEP 8
JSONを扱い、ランキングと前月比を出せる
11週目
STEP 9+復習
Time Travelで戻せる。テーブルを速くできる
12〜13週目
STEP 10
毎朝自動で更新される集計表を作る
14週目
STEP 11〜12
権限を設計し、コストと性能を説明できる
15週目〜
STEP 13
Gitとダッシュボードをつなぎ、卒業課題を完成させる
続けるためのコツ ① 自分の仕事のデータを使う。サンプルデータだけで進めると、途中で必ず飽きます ② 完璧に理解してから次へ進もうとしない。8割わかったら先に進み、必要になったときに戻るほうが結局は速い ③ うまくいったコードは必ずファイルに保存して、Gitに置く。3か月後の自分が必ず助かります。
まず落ち着いて、それ以上の操作をやめてください。Delta LakeにはTime Travelがあり、既定で30日分の履歴が残っています(STEP 9)。DESCRIBE HISTORY テーブル で事故のバージョンを特定し、VERSION AS OF で直前の状態を確認してから、RESTORE TABLE ... TO VERSION AS OF n で戻します。慌てて上書きを重ねるのが最悪の対応です。なお、VACUUM を実行した後の古いバージョンには戻れません。
資格(Databricks認定)は取るべきですか
実務に入る前の目標としては良い教材です。Data Engineer Associateが入門で、このページのSTEP 0〜11の範囲がおおむね対応します。分析寄りなら Data Analyst Associate です。ただし、資格があるから任せてもらえるわけではありません。「自分で取り込んで、自動で回して、権限とコストを説明できる」ものを1つ作った経験のほうが、採用でも社内でも評価されます。資格は、学習の抜け漏れを埋める道具として使うのがおすすめです。