-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathsetup_databricks.sql
More file actions
197 lines (180 loc) · 7.16 KB
/
Copy pathsetup_databricks.sql
File metadata and controls
197 lines (180 loc) · 7.16 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
-- =============================================================================
-- Medallion architecture for the OpenMetadata + Kafka demo
-- Run this in your Databricks Free Edition SQL Editor after creating the
-- Unity Catalog storage credential + external location for the S3 staging bucket.
-- =============================================================================
CREATE CATALOG IF NOT EXISTS demo;
USE CATALOG demo;
CREATE SCHEMA IF NOT EXISTS bronze COMMENT 'Raw CDC events from Kafka (append-only)';
CREATE SCHEMA IF NOT EXISTS silver COMMENT 'Cleaned current-state tables';
CREATE SCHEMA IF NOT EXISTS gold COMMENT 'Business aggregates for analytics';
-- ---------------------------------------------------------------------------
-- Bronze: Landing tables for the Confluent Delta Lake Sink.
-- Connector settings delta.lake.catalog / delta.lake.database target demo.bronze
-- with short table names (customers, products, ...). auto.create is off —
-- create these tables before starting the sink.
--
-- Schema MUST match Debezium ExtractNewRecordState + staging Avro:
-- - INT ids (not BIGINT for PKs)
-- - TIMESTAMPTZ → STRING (io.debezium.time.ZonedTimestamp ISO-8601)
-- - __op / __table / __source_ts_ms / __deleted (unwrap rewrite adds __deleted)
-- - no partition column (staging files omit it)
-- - order_items has no line_total (Postgres GENERATED; compute in silver)
-- If COPY INTO fails with schema mismatch, run databricks/recreate_bronze.sql
-- ---------------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS demo.bronze.customers (
id INT,
email STRING,
full_name STRING,
phone STRING,
created_at STRING,
__op STRING,
__table STRING,
__source_ts_ms BIGINT,
__deleted STRING
) USING DELTA
COMMENT 'Raw customer CDC events from Kafka topic dbz.ecommerce.public.customers';
CREATE TABLE IF NOT EXISTS demo.bronze.products (
id INT,
sku STRING,
name STRING,
category STRING,
unit_price STRING,
updated_at STRING,
__op STRING,
__table STRING,
__source_ts_ms BIGINT,
__deleted STRING
) USING DELTA
COMMENT 'Raw product CDC events from Kafka topic dbz.ecommerce.public.products';
CREATE TABLE IF NOT EXISTS demo.bronze.orders (
id INT,
customer_id INT,
status STRING,
currency STRING,
total_amount STRING,
placed_at STRING,
updated_at STRING,
__op STRING,
__table STRING,
__source_ts_ms BIGINT,
__deleted STRING
) USING DELTA
COMMENT 'Raw order CDC events from Kafka topic dbz.ecommerce.public.orders';
CREATE TABLE IF NOT EXISTS demo.bronze.order_items (
id INT,
order_id INT,
product_id INT,
quantity INT,
unit_price STRING,
__op STRING,
__table STRING,
__source_ts_ms BIGINT,
__deleted STRING
) USING DELTA
COMMENT 'Raw order_item CDC events from Kafka topic dbz.ecommerce.public.order_items';
-- ---------------------------------------------------------------------------
-- (Optional legacy path) If an older sink run landed in workspace.default.*,
-- promote into demo.bronze — normally unnecessary once catalog/database are set.
-- ---------------------------------------------------------------------------
-- INSERT INTO demo.bronze.customers SELECT * FROM workspace.default.customers;
-- INSERT INTO demo.bronze.products SELECT * FROM workspace.default.products;
-- INSERT INTO demo.bronze.orders SELECT * FROM workspace.default.orders;
-- INSERT INTO demo.bronze.order_items SELECT * FROM workspace.default.order_items;
-- ---------------------------------------------------------------------------
-- Silver: current-state tables (dedupe CDC append log → latest row per PK).
-- Temporal columns: ZonedTimestamp ISO-8601 strings → TIMESTAMP via to_timestamp().
-- Re-run after the sink has flushed data.
-- ---------------------------------------------------------------------------
CREATE OR REPLACE TABLE demo.silver.customers AS
SELECT id, email, full_name, phone,
to_timestamp(created_at) AS created_at,
__source_ts_ms AS last_cdc_ts
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY id ORDER BY __source_ts_ms DESC) AS rn
FROM demo.bronze.customers
WHERE COALESCE(__op, 'c') <> 'd'
AND COALESCE(__deleted, 'false') <> 'true'
)
WHERE rn = 1;
CREATE OR REPLACE TABLE demo.silver.products AS
SELECT id, sku, name, category,
CAST(unit_price AS DECIMAL(12, 2)) AS unit_price,
to_timestamp(updated_at) AS updated_at,
__source_ts_ms AS last_cdc_ts
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY id ORDER BY __source_ts_ms DESC) AS rn
FROM demo.bronze.products
WHERE COALESCE(__op, 'c') <> 'd'
AND COALESCE(__deleted, 'false') <> 'true'
)
WHERE rn = 1;
CREATE OR REPLACE TABLE demo.silver.orders AS
SELECT id, customer_id, status, currency,
CAST(total_amount AS DECIMAL(12, 2)) AS total_amount,
to_timestamp(placed_at) AS placed_at,
to_timestamp(updated_at) AS updated_at,
__source_ts_ms AS last_cdc_ts
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY id ORDER BY __source_ts_ms DESC) AS rn
FROM demo.bronze.orders
WHERE COALESCE(__op, 'c') <> 'd'
AND COALESCE(__deleted, 'false') <> 'true'
)
WHERE rn = 1;
CREATE OR REPLACE TABLE demo.silver.order_items AS
SELECT id, order_id, product_id, quantity,
CAST(unit_price AS DECIMAL(12, 2)) AS unit_price,
quantity * CAST(unit_price AS DECIMAL(12, 2)) AS line_total,
__source_ts_ms AS last_cdc_ts
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY id ORDER BY __source_ts_ms DESC) AS rn
FROM demo.bronze.order_items
WHERE COALESCE(__op, 'c') <> 'd'
AND COALESCE(__deleted, 'false') <> 'true'
)
WHERE rn = 1;
-- ---------------------------------------------------------------------------
-- Gold: analytics aggregates (Unity Catalog records lineage from silver → gold)
-- ---------------------------------------------------------------------------
CREATE OR REPLACE TABLE demo.gold.customer_order_summary AS
SELECT
c.id AS customer_id,
c.email,
c.full_name,
COUNT(DISTINCT o.id) AS order_count,
COALESCE(SUM(o.total_amount), 0) AS lifetime_value,
MAX(o.placed_at) AS last_order_at
FROM demo.silver.customers c
LEFT JOIN demo.silver.orders o
ON c.id = o.customer_id
AND o.status <> 'cancelled'
GROUP BY c.id, c.email, c.full_name;
CREATE OR REPLACE TABLE demo.gold.daily_revenue AS
SELECT
CAST(o.placed_at AS DATE) AS order_date,
o.currency,
COUNT(DISTINCT o.id) AS orders,
SUM(oi.line_total) AS revenue
FROM demo.silver.orders o
JOIN demo.silver.order_items oi ON o.id = oi.order_id
WHERE o.status <> 'cancelled'
GROUP BY CAST(o.placed_at AS DATE), o.currency;
CREATE OR REPLACE TABLE demo.gold.product_performance AS
SELECT
p.id AS product_id,
p.sku,
p.name,
p.category,
SUM(oi.quantity) AS units_sold,
SUM(oi.line_total) AS revenue
FROM demo.silver.products p
LEFT JOIN demo.silver.order_items oi ON p.id = oi.product_id
LEFT JOIN demo.silver.orders o
ON oi.order_id = o.id
AND o.status <> 'cancelled'
GROUP BY p.id, p.sku, p.name, p.category;