-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy path02_observability.sql
More file actions
72 lines (66 loc) · 3.24 KB
/
Copy path02_observability.sql
File metadata and controls
72 lines (66 loc) · 3.24 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
-- =====================================================================
-- 02_observability.sql — SLO & monitoring du pipeline crypto temps réel
-- Rôle conseillé : CRYPTO_PIPELINE_ROLE (sauf §6 coût : ACCOUNTADMIN possible).
-- =====================================================================
USE ROLE CRYPTO_PIPELINE_ROLE;
USE WAREHOUSE WH_CRYPTO_XS;
-- 1) LATENCE d'ingestion (event Binance -> réception consumer) ---------
-- SLO cible : p95 < 15 s end-to-end. Ici on mesure event -> réception ;
-- ajouter ~5-10 s de commit Snowpipe Streaming pour l'end-to-end.
SELECT
symbol,
COUNT(*) AS trades_5min,
ROUND(AVG(DATEDIFF('millisecond', traded_at, ingest_time))/1000.0, 3) AS avg_latency_s,
ROUND(PERCENTILE_CONT(0.95) WITHIN GROUP (
ORDER BY DATEDIFF('millisecond', traded_at, ingest_time))/1000.0, 3) AS p95_latency_s
FROM ANALYTICS.PUBLIC_STAGING.STG_TRADES
WHERE ingest_time >= DATEADD('minute', -5, CURRENT_TIMESTAMP())
GROUP BY symbol
ORDER BY symbol;
-- 2) FRAÎCHEUR (dernier event / ingest vs maintenant) -----------------
-- SLO : warn > 30 s, error > 120 s.
SELECT
'trades' AS source,
MAX(traded_at) AS last_event_time,
MAX(ingest_time) AS last_ingest_time,
DATEDIFF('second', MAX(ingest_time), CURRENT_TIMESTAMP()) AS freshness_seconds
FROM ANALYTICS.PUBLIC_STAGING.STG_TRADES
UNION ALL
SELECT
'depth',
NULL,
MAX(ingest_time),
DATEDIFF('second', MAX(ingest_time), CURRENT_TIMESTAMP())
FROM ANALYTICS.PUBLIC_STAGING.STG_DEPTH_LEVELS;
-- 3) DÉBIT (lignes/min, 15 dernières min) -----------------------------
SELECT
DATE_TRUNC('minute', ingest_time) AS minute,
COUNT(*) AS trades
FROM ANALYTICS.PUBLIC_STAGING.STG_TRADES
WHERE ingest_time >= DATEADD('minute', -15, CURRENT_TIMESTAMP())
GROUP BY 1
ORDER BY 1 DESC;
-- 4) DYNAMIC TABLES — succès & lag des refreshs ------------------------
-- Historique des rafraîchissements (état SUCCEEDED attendu).
SELECT name, state, refresh_start_time, refresh_end_time, data_timestamp
FROM TABLE(ANALYTICS.INFORMATION_SCHEMA.DYNAMIC_TABLE_REFRESH_HISTORY())
ORDER BY refresh_end_time DESC
LIMIT 20;
-- Lag cible vs réel (target_lag / mean_lag / max_lag) :
SHOW DYNAMIC TABLES IN SCHEMA ANALYTICS.PUBLIC_MARTS;
-- 5) QUALITÉ — taux de déduplication (raw vs silver) ------------------
SELECT
(SELECT COUNT(*) FROM RAW.CRYPTO.RAW_TRADES) AS raw_trades,
(SELECT COUNT(*) FROM ANALYTICS.PUBLIC_STAGING.STG_TRADES) AS dedup_trades,
ROUND(100 * (1 - (SELECT COUNT(*) FROM ANALYTICS.PUBLIC_STAGING.STG_TRADES)
/ NULLIF((SELECT COUNT(*) FROM RAW.CRYPTO.RAW_TRADES), 0)), 1)
AS duplicate_pct;
-- 6) COÛT (crédits du warehouse, 24 dernières h) — FinOps --------------
SELECT
DATE_TRUNC('hour', start_time) AS hour,
ROUND(SUM(credits_used), 4) AS credits
FROM TABLE(ANALYTICS.INFORMATION_SCHEMA.WAREHOUSE_METERING_HISTORY(
DATE_RANGE_START => DATEADD('day', -1, CURRENT_DATE()),
WAREHOUSE_NAME => 'WH_CRYPTO_XS'))
GROUP BY 1
ORDER BY 1 DESC;