-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmaster_data.sql
More file actions
355 lines (311 loc) · 12.3 KB
/
Copy pathmaster_data.sql
File metadata and controls
355 lines (311 loc) · 12.3 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
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
-- QUERY FOR USER PRIVILEGE
-- Uses SNOWFLAKE.ACCOUNT_USAGE views available in all Snowflake accounts
CREATE OR REPLACE TEMPORARY TABLE CTE_USER_PRIVILEGE AS
SELECT
TABLE_NAME,
LISTAGG((P ||' -> '|| USER ), ' | ') WITHIN GROUP (ORDER BY P) AS USER_PRIVILEGE
FROM (
SELECT
CONCAT(R.TABLE_CATALOG,'.',R.TABLE_SCHEMA,'.',R.NAME) AS TABLE_NAME,
R.PRIVILEGE AS P,
LISTAGG (DISTINCT (U.GRANTEE_NAME), ', ') AS USER
FROM
SNOWFLAKE.ACCOUNT_USAGE.GRANTS_TO_ROLES R
LEFT JOIN
SNOWFLAKE.ACCOUNT_USAGE.GRANTS_TO_USERS U ON R.GRANTEE_NAME = U.ROLE
WHERE
U.DELETED_ON IS NULL
AND U.GRANTED_TO = 'USER'
AND R.GRANTED_ON IN ('TABLE', 'VIEW', 'MATERIALIZED VIEW')
AND R.GRANTED_TO = 'ROLE'
GROUP BY
TABLE_NAME,
R.PRIVILEGE
) AS Subquery
GROUP BY
TABLE_NAME ;
CREATE OR REPLACE TEMPORARY TABLE DATA_LINEAGE_MASTER_DATA_TMP AS
WITH COPY_HIST AS (
-- admin.ctrl_lambda_integration_details removed (Piramal-specific Lambda control table)
SELECT TABLE_SCHEMA_NAME, TABLE_CATALOG_NAME, TABLE_NAME, FILE_NAME, STAGE_LOCATION, LAST_LOAD_TIME, '' as load_type
FROM SNOWFLAKE.ACCOUNT_USAGE.COPY_HISTORY
WHERE last_load_time >= DATE_TRUNC('MONTH', CURRENT_DATE())
AND last_load_time < DATEADD('MONTH', 1, DATE_TRUNC('MONTH', CURRENT_DATE()))
AND stage_location ilike '%S3%'
QUALIFY ROW_NUMBER() OVER (PARTITION BY (TABLE_CATALOG_NAME||'.'||TABLE_SCHEMA_NAME||'.'||TABLE_NAME) ORDER BY LAST_LOAD_TIME DESC) = 1
),
copy_hist_complete_data as (
SELECT DISTINCT
cop.file_name AS SOURCE_OBJECT_NAME,
'' AS SOURCE_DATABASE ,
'' SOURCE_SCHEMA,
cop.file_name AS SOURCE_TABLE ,
'' AS S_TABLE_TYPE,
OBJECT_CONSTRUCT(
'S3_PATH', COP.STAGE_LOCATION,
'LOAD_TYPE', COP.LOAD_TYPE
) AS SOURCE_METADATA,
'' AS S_USER_PRIVILEGE,
'' AS S_COLUMN_LIST ,
UPPER(CONCAT('ECOMMERCE', '.', COP.TABLE_SCHEMA_NAME, '.', COP.TABLE_NAME)) AS TARGET_OBJECT_NAME,
'ECOMMERCE' AS TARGET_DATABASE,
UPPER(COP.TABLE_SCHEMA_NAME) AS TARGET_SCHEMA,
UPPER(COP.TABLE_NAME)AS TARGET_TABLE,
'PERMANENT' AS T_TABLE_TYPE,
OBJECT_CONSTRUCT(
'TABLE_OWNER', T_META.TABLE_OWNER,
'TABLE_TYPE', T_META.TABLE_TYPE,
'ROW_COUNT', T_META.ROW_COUNT,
'BYTES', T_META.BYTES,
'CREATED', TO_CHAR(T_META.CREATED::DATE, 'DD-MM-YYYY'),
'LAST_ALTERED', TO_CHAR(T_META.LAST_ALTERED::DATE, 'DD-MM-YYYY'),
'LAST_DDL', TO_CHAR(T_META.LAST_DDL::DATE, 'DD-MM-YYYY'),
'LAST_DDL_BY', T_META.LAST_DDL_BY,
'IS_ICEBERG', T_META.IS_ICEBERG,
'IS_TRANSIENT',T_META.IS_TRANSIENT,
'USER',''
) AS TARGET_METADATA,
T_USER_PRIVILEGE.USER_PRIVILEGE AS T_USER_PRIVILEGE,
'' AS T_COLUMN_LIST,
'' AS S_T_COLUMN_MAPPING,
'' AS QUERY_TAG,
'' AS DATA_PRODUCT_FLAG ,
'S3' AS DATA_PIPELINE_FLAG,
object_construct() AS PIPELINE_METADATA,
TO_CHAR(LAST_LOAD_TIME::DATE, 'DD-MM-YYYY') AS LAST_SYNC_DATE,
'' AS QUERY_TEXT
FROM
COPY_HIST AS COP
LEFT JOIN
INFORMATION_SCHEMA.TABLES AS T_META ON UPPER(CONCAT('ECOMMERCE', '.', COP.TABLE_SCHEMA_NAME, '.', COP.TABLE_NAME)) = UPPER(CONCAT_WS('.', T_META.TABLE_CATALOG, T_META.TABLE_SCHEMA, T_META.TABLE_NAME))
LEFT JOIN
CTE_USER_PRIVILEGE AS T_USER_PRIVILEGE ON UPPER(CONCAT('ECOMMERCE', '.',COP.TABLE_SCHEMA_NAME, '.',COP.TABLE_NAME)) = T_USER_PRIVILEGE.TABLE_NAME
GROUP BY ALL
),
-- CTE FOR QUERY HISTORY
-- Captures tables, views, and materialized views
TempLatestQueryHistory AS (
SELECT
QUERY_ID,
QUERY_TEXT,
QUERY_TAG,
ROLE_NAME,
ROWS_INSERTED,
ROWS_DELETED,
ROWS_UPDATED,
SCHEMA_NAME,
START_TIME AS LAST_SYNC_DATE,
ROW_NUMBER() OVER (
PARTITION BY QUERY_TEXT, SCHEMA_NAME, ROLE_NAME
ORDER BY START_TIME DESC
) AS RN
FROM
SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORY
WHERE
database_name IN ('ECOMMERCE')
--AND SCHEMA_NAME IN ('DATA_PRODUCTS')
AND START_TIME >= DATE_TRUNC('MONTH', CURRENT_DATE())
AND START_TIME < DATEADD('MONTH', 1, DATE_TRUNC('MONTH', CURRENT_DATE()))
AND EXECUTION_STATUS = 'SUCCESS'
AND QUERY_TYPE IN ('INSERT','CREATE_TABLE_AS_SELECT','CREATE_VIEW','CREATE_MATERIALIZED_VIEW')
)
,
-- CTE FOR ACCESS HISTORY
CTE_DATA_LINEAGE_ACCESS_HISTORY AS (
SELECT DISTINCT
t.user_name,
QUERY_ID,
baseSources.value: "objectName" as source_object_name,
baseSources.value: "columnName" as source_column_name,
om.value: "objectName" as target_object_name,
columns_modified.value: "columnName" as target_column_name
FROM
(SELECT * FROM snowflake.account_usage.access_history
WHERE QUERY_ID IN (SELECT QUERY_ID FROM TempLatestQueryHistory)) t,
LATERAL FLATTEN(input => t.OBJECTS_MODIFIED) om,
LATERAL FLATTEN(input => om.value: "columns", outer => true) columns_modified,
LATERAL FLATTEN(input => columns_modified.value: "baseSources", outer => true) baseSources
WHERE
source_object_name IS NOT NULL AND target_object_name IS NOT NULL
)
-- FINAL SELECT
SELECT DISTINCT
-- QH.SCHEMA_NAME,
REPLACE(AH.SOURCE_OBJECT_NAME,'"','' )AS SOURCE_OBJECT_NAME,
SPLIT_PART(REPLACE(AH.SOURCE_OBJECT_NAME, '"', ''), '.', 1) AS SOURCE_DATABASE,
SPLIT_PART(REPLACE(AH.SOURCE_OBJECT_NAME, '"', ''), '.', 2) AS SOURCE_SCHEMA,
SPLIT_PART(REPLACE(AH.SOURCE_OBJECT_NAME, '"', ''), '.', 3) AS SOURCE_TABLE,
CASE
WHEN S_META.TABLE_TYPE = 'VIEW' THEN 'VIEW'
WHEN S_META.TABLE_TYPE = 'MATERIALIZED VIEW' THEN 'MATERIALIZED VIEW'
WHEN S_META.IS_TRANSIENT = 'YES' THEN 'TRANSIENT'
WHEN S_META.BYTES IS NULL THEN 'TEMPORARY'
ELSE 'PERMANENT'
END AS S_TABLE_TYPE,
OBJECT_CONSTRUCT(
'TABLE_OWNER', S_META.TABLE_OWNER,
'TABLE_TYPE', S_META.TABLE_TYPE,
'ROW_COUNT', S_META.ROW_COUNT,
'BYTES', S_META.BYTES,
'CREATED', TO_CHAR(S_META.CREATED::DATE, 'DD-MM-YYYY'),
'LAST_ALTERED', TO_CHAR(S_META.LAST_ALTERED::DATE, 'DD-MM-YYYY'),
'LAST_DDL', TO_CHAR(S_META.LAST_DDL::DATE, 'DD-MM-YYYY'),
'LAST_DDL_BY', S_META.LAST_DDL_BY,
'IS_ICEBERG', S_META.IS_ICEBERG,
'IS_TRANSIENT',S_META.IS_TRANSIENT,
'USER',AH.user_name
) AS SOURCE_METADATA,
S_USER_PRIVILEGE.USER_PRIVILEGE AS S_USER_PRIVILEGE,
LISTAGG(AH.source_column_name, ',') AS S_COLUMN_LIST,
REPLACE(AH.TARGET_OBJECT_NAME,'"','' )AS TARGET_OBJECT_NAME,
SPLIT_PART(REPLACE(AH.TARGET_OBJECT_NAME, '"', ''), '.', 1) AS TARGET_DATABASE,
SPLIT_PART(REPLACE(AH.TARGET_OBJECT_NAME, '"', ''), '.', 2) AS TARGET_SCHEMA,
SPLIT_PART(REPLACE(AH.TARGET_OBJECT_NAME, '"', ''), '.', 3) AS TARGET_TABLE,
CASE
WHEN T_META.TABLE_TYPE = 'VIEW' THEN 'VIEW'
WHEN T_META.TABLE_TYPE = 'MATERIALIZED VIEW' THEN 'MATERIALIZED VIEW'
WHEN T_META.IS_TRANSIENT = 'YES' THEN 'TRANSIENT'
WHEN T_META.BYTES IS NULL THEN 'TEMPORARY'
ELSE 'PERMANENT'
END AS T_TABLE_TYPE,
OBJECT_CONSTRUCT(
'TABLE_OWNER', T_META.TABLE_OWNER,
'TABLE_TYPE', T_META.TABLE_TYPE,
'ROW_COUNT', T_META.ROW_COUNT,
'BYTES', T_META.BYTES,
'CREATED', TO_CHAR(T_META.CREATED::DATE, 'DD-MM-YYYY'),
'LAST_ALTERED', TO_CHAR(T_META.LAST_ALTERED::DATE, 'DD-MM-YYYY'),
'LAST_DDL', TO_CHAR(T_META.LAST_DDL::DATE, 'DD-MM-YYYY'),
'LAST_DDL_BY', T_META.LAST_DDL_BY,
'IS_ICEBERG', T_META.IS_ICEBERG,
'IS_TRANSIENT',T_META.IS_TRANSIENT,
'USER',AH.user_name
) AS TARGET_METADATA,
T_USER_PRIVILEGE.USER_PRIVILEGE AS T_USER_PRIVILEGE,
LISTAGG(AH.target_column_name, ',') AS T_COLUMN_LIST,
LISTAGG(CONCAT(AH.source_column_name, ' -> ', AH.target_column_name), ' , ') WITHIN GROUP (ORDER BY AH.source_column_name) AS S_T_COLUMN_MAPPING,
QH.QUERY_TAG,
CASE
WHEN QH.QUERY_TAG ILIKE 'DP-%' THEN 'DATA PRODUCT'
ELSE ''
END AS DATA_PRODUCT_FLAG,
'' AS DATA_PIPELINE_FLAG,
OBJECT_CONSTRUCT(
) AS PIPELINE_METADATA,
TO_CHAR(LAST_SYNC_DATE::DATE, 'DD-MM-YYYY') AS LAST_SYNC_DATE,
QH.QUERY_TEXT AS QUERY_TEXT
FROM
CTE_DATA_LINEAGE_ACCESS_HISTORY AH
LEFT JOIN
TempLatestQueryHistory QH ON AH.QUERY_ID = QH.QUERY_ID
LEFT JOIN
INFORMATION_SCHEMA.TABLES S_META ON UPPER(AH.SOURCE_OBJECT_NAME) = UPPER(CONCAT_WS('.', S_META.TABLE_CATALOG, S_META.TABLE_SCHEMA, S_META.TABLE_NAME))
LEFT JOIN
INFORMATION_SCHEMA.TABLES T_META ON UPPER(AH.TARGET_OBJECT_NAME) = UPPER(CONCAT_WS('.', T_META.TABLE_CATALOG, T_META.TABLE_SCHEMA, T_META.TABLE_NAME))
LEFT JOIN
CTE_USER_PRIVILEGE AS S_USER_PRIVILEGE ON AH.SOURCE_OBJECT_NAME = S_USER_PRIVILEGE.TABLE_NAME
LEFT JOIN
CTE_USER_PRIVILEGE AS T_USER_PRIVILEGE ON AH.TARGET_OBJECT_NAME = T_USER_PRIVILEGE.TABLE_NAME
WHERE QH.RN = 1
GROUP BY ALL
-- CTRL_TABLE_DETAILS, sftp_data, and api_data were Piramal-specific tables not available in a generic Snowflake account.
-- Replace the UNION ALL blocks below with your own custom data sources if needed.
-- UNION ALL SELECT * FROM CTRL_TABLE_DETAILS
-- UNION ALL SELECT * FROM sftp_data
-- UNION ALL SELECT * FROM api_data
UNION ALL
SELECT * FROM copy_hist_complete_data;
--CREATE IS DELETE TEMP TABLE BASED ON ALL FILTER CONDITIONS
CREATE OR REPLACE TABLE DATA_LINEAGE_MASTER_DATA AS
--GETTING DISTINCT DATA FOR FILTERS
WITH FilteredMasterData AS (
SELECT DISTINCT
TARGET_DATABASE,
TARGET_SCHEMA
FROM
DATA_LINEAGE_MASTER_DATA_TMP
GROUP BY ALL
UNION
SELECT DISTINCT
SOURCE_DATABASE,
SOURCE_SCHEMA
FROM
DATA_LINEAGE_MASTER_DATA_TMP
WHERE SOURCE_DATABASE = 'ECOMMERCE'
GROUP BY ALL
),
MIN_SYNC_CTE AS (
SELECT
TO_CHAR(MIN(TO_DATE(LAST_SYNC_DATE, 'DD-MM-YYYY')), 'YYYY-MM-DD') AS MIN_LAST_SYNC_DATE
FROM
DATA_LINEAGE_MASTER_DATA_TMP
WHERE
COALESCE(DATA_PIPELINE_FLAG,'') = ''
),
-- FILTERED QUERY HISTORY
-- Captures drops of tables, views, and materialized views
FilteredHistory AS (
SELECT
QH.QUERY_TEXT,
QH.QUERY_ID,
QH.SCHEMA_NAME,
QH.DATABASE_NAME,
QH.START_TIME,
QH.USER_NAME,
ROW_NUMBER() OVER (PARTITION BY QH.SCHEMA_NAME, QH.DATABASE_NAME, QH.QUERY_TEXT ORDER BY QH.START_TIME) AS rn
FROM
SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORY AS QH
JOIN
FilteredMasterData AS FMD ON QH.DATABASE_NAME = FMD.TARGET_DATABASE AND QH.SCHEMA_NAME = FMD.TARGET_SCHEMA
INNER JOIN
MIN_SYNC_CTE
WHERE
QH.QUERY_TYPE IN ('DROP', 'DROP_VIEW', 'DROP_MATERIALIZED_VIEW')
AND QH.execution_status = 'SUCCESS'
AND QH.START_TIME >= DATE_TRUNC('MONTH', CURRENT_DATE())
AND QH.START_TIME < DATEADD('MONTH', 1, DATE_TRUNC('MONTH', CURRENT_DATE()))
)
--MAIN SELECT FOR IS DELETE
SELECT
master.*,
OBJECT_CONSTRUCT(
'SOURCE_DELETE_FLAG' ,CASE
WHEN SOURCE_history.QUERY_ID IS NOT NULL THEN 'DROPPED'
ELSE NULL
END ,
'QUERY_TEXT' , SOURCE_history.QUERY_TEXT,
'USER' ,SOURCE_history.USER_NAME,
'DELETE_DATE',TO_CHAR(SOURCE_history.START_TIME::DATE, 'DD-MM-YYYY')
) AS S_DELETE_METADATA ,
OBJECT_CONSTRUCT(
'TARGET_DELETE_FLAG' , CASE
WHEN TARGET_history.QUERY_ID IS NOT NULL THEN 'DROPPED'
ELSE NULL
END ,
'QUERY_TEXT' ,TARGET_history.QUERY_TEXT,
'USER' ,TARGET_history.USER_NAME,
'DELETE_DATE',TO_CHAR(TARGET_history.START_TIME::DATE, 'DD-MM-YYYY' )
) AS T_DELETE_METADATA
FROM
DATA_LINEAGE_MASTER_DATA_TMP AS master
LEFT JOIN
(
SELECT *
FROM
FilteredHistory
WHERE
rn = 1
) AS SOURCE_history ON SOURCE_history.DATABASE_NAME = master.SOURCE_DATABASE
AND SOURCE_history.SCHEMA_NAME = master.SOURCE_SCHEMA
AND REGEXP_SUBSTR(SOURCE_history.QUERY_TEXT, 'DROP\\s+.*\\b' || master.SOURCE_TABLE || '\\b', 1, 1, 'i') IS NOT NULL
LEFT JOIN
(
SELECT *
FROM
FilteredHistory
WHERE
rn = 1
) AS
TARGET_history ON TARGET_history.DATABASE_NAME = master.TARGET_DATABASE
AND TARGET_history.SCHEMA_NAME = master.TARGET_SCHEMA
AND REGEXP_SUBSTR(TARGET_history.QUERY_TEXT, 'DROP\\s+.*\\b' || master.TARGET_TABLE || '\\b', 1, 1, 'i') IS NOT NULL ;