Compare commits
2 Commits
test_chara
...
exp360-cus
| Author | SHA1 | Date | |
|---|---|---|---|
| 403577275c | |||
| b9cd8ab1e8 |
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Document","metadata":{"name":"enrich-360.json","description":"","runtime":"spark","kind":"storage#object","id":"ocular-unstructured-data-cod-uat/content/workspaces/test_charan/content/enrich-360.json/v1/1784873517595730","selfLink":"https://www.googleapis.com/storage/v1/b/ocular-unstructured-data-cod-uat/o/content%2Fworkspaces%2Ftest_charan%2Fcontent%2Fenrich-360.json%2Fv1","mediaLink":"https://storage.googleapis.com/download/storage/v1/b/ocular-unstructured-data-cod-uat/o/content%2Fworkspaces%2Ftest_charan%2Fcontent%2Fenrich-360.json%2Fv1?generation=1784873517595730&alt=media","bucket":"ocular-unstructured-data-cod-uat","generation":"1784873517595730","metageneration":"1","contentType":"application/json","storageClass":"STANDARD","size":47856,"md5Hash":"UYTX41oC9i1GR4kcOncgtQ==","crc32c":"KI2P7g==","etag":"CNL48P/T6pUDEAE=","timeCreated":"2026-07-24T06:11:57.637Z","updated":"2026-07-24T06:11:57.637Z","timeStorageClassUpdated":"2026-07-24T06:11:57.637Z","timeFinalized":"2026-07-24T06:11:57.637Z","type":"file","mtime":"2026-07-24T06:11:57.637000Z","ctime":"2026-07-24T06:11:57.637000Z","version":1,"absolute_path":"gs://ocular-unstructured-data-cod-uat/content/workspaces/test_charan/content/enrich-360.json/v1","content_type":"application/json","content_size":47856,"storage_type":"gs"},"spec":{"ui":{},"blocks":[]}}
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Document","metadata":{"name":"main.py","description":"","runtime":"spark","kind":"storage#object","id":"ocular-unstructured-data-cod-uat/content/workspaces/test_charan/content/main.py/v1/1784872964382918","selfLink":"https://www.googleapis.com/storage/v1/b/ocular-unstructured-data-cod-uat/o/content%2Fworkspaces%2Ftest_charan%2Fcontent%2Fmain.py%2Fv1","mediaLink":"https://storage.googleapis.com/download/storage/v1/b/ocular-unstructured-data-cod-uat/o/content%2Fworkspaces%2Ftest_charan%2Fcontent%2Fmain.py%2Fv1?generation=1784872964382918&alt=media","bucket":"ocular-unstructured-data-cod-uat","generation":"1784872964382918","metageneration":"1","contentType":"text/x-python","storageClass":"STANDARD","size":3487,"md5Hash":"Ut4LddZ5CJG6TuQ5HxuLWA==","crc32c":"j1sgnQ==","etag":"CMbBi/jR6pUDEAE=","timeCreated":"2026-07-24T06:02:44.431Z","updated":"2026-07-24T06:02:44.431Z","timeStorageClassUpdated":"2026-07-24T06:02:44.431Z","timeFinalized":"2026-07-24T06:02:44.431Z","type":"file","mtime":"2026-07-24T06:02:44.431000Z","ctime":"2026-07-24T06:02:44.431000Z","version":1,"absolute_path":"gs://ocular-unstructured-data-cod-uat/content/workspaces/test_charan/content/main.py/v1","content_type":"text/x-python","content_size":3487,"storage_type":"gs"},"spec":{"ui":{},"blocks":[]}}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing",
|
||||
"type": "PHYSICAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-24T09:56:34.021820",
|
||||
"created_by": null,
|
||||
"schema_hash": "411e8293819c",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-24T09:56:34.021820",
|
||||
"updated_at": "2026-07-24T09:56:34.021820",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing_physical_dataset1",
|
||||
"type": "PHYSICAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T05:57:12.382527",
|
||||
"created_by": null,
|
||||
"schema_hash": "6ec727aad414",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing_physical_dataset1"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T05:57:12.382527",
|
||||
"updated_at": "2026-07-22T05:57:12.382527",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1,26 +0,0 @@
|
||||
{
|
||||
"name": "del",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-24T08:28:52.025865",
|
||||
"created_by": null,
|
||||
"schema_hash": "86a8f4d9d691",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "del",
|
||||
"updated_at": "2026-07-24T09:35:52.585382"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-24T08:28:52.025865",
|
||||
"updated_at": "2026-07-24T09:35:52.585382",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "del", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-cod-customer\".bills", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "del"], "version": "v1", "created_date": "2026-07-24 09:35:52.585382", "updated_date": "2026-07-24 09:35:52.585382", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": "20260724_093552", "schedule": null, "lineage": null, "status": "active", "schema_hash": "86a8f4d9d691"}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test charan123", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit\n", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test charan123"], "version": "v1", "created_date": "2026-07-22 06:55:20.664226", "updated_date": "2026-07-22 06:55:20.664226", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6ec727aad414"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_c",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-24T07:24:23.076999",
|
||||
"created_by": null,
|
||||
"schema_hash": "79a17f89b99c",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_c"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-24T07:24:23.076999",
|
||||
"updated_at": "2026-07-24T07:24:23.076999",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_c", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-cod-customer\".actionsaudit", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_c"], "version": "v1", "created_date": "2026-07-24 07:24:23.076999", "updated_date": "2026-07-24 07:24:23.076999", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "79a17f89b99c"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_charan",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T05:44:00.404637",
|
||||
"created_by": null,
|
||||
"schema_hash": "6ec727aad414",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_charan"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T05:44:00.404637",
|
||||
"updated_at": "2026-07-22T05:44:00.404637",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_charan", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit\n", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_charan"], "version": "v1", "created_date": "2026-07-22 05:44:00.404637", "updated_date": "2026-07-22 05:44:00.404637", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6ec727aad414"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_ddel",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-24T08:26:51.469548",
|
||||
"created_by": null,
|
||||
"schema_hash": "6244490a0e95",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_ddel"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-24T08:26:51.469548",
|
||||
"updated_at": "2026-07-24T08:26:51.469548",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_ddel", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".customermapping", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_ddel"], "version": "v1", "created_date": "2026-07-24 08:26:51.469548", "updated_date": "2026-07-24 08:26:51.469548", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6244490a0e95"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_v_23072026_1",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "testing",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-23T12:21:22.814571",
|
||||
"created_by": null,
|
||||
"schema_hash": "16d7910d16bc",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_v_23072026_1"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-23T12:21:22.814571",
|
||||
"updated_at": "2026-07-23T12:21:22.814571",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_v_23072026_1", "description": "testing", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit LIMIT 2;", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_v_23072026_1"], "version": "v1", "created_date": "2026-07-23 12:21:22.814571", "updated_date": "2026-07-23 12:21:22.814571", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "16d7910d16bc"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_v_2326",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "testing",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-23T15:39:19.830617",
|
||||
"created_by": null,
|
||||
"schema_hash": "44c5c04d90a4",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_v_2326"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-23T15:39:19.830617",
|
||||
"updated_at": "2026-07-23T15:39:19.830617",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
eyJlbnRpdHlUeXBlIjogImRhdGFzZXQiLCAibmFtZSI6ICJ0ZXN0X3ZfMjMyNiIsICJkZXNjcmlwdGlvbiI6ICJ0ZXN0aW5nIiwgInR5cGUiOiAiVklSVFVBTF9EQVRBU0VUIiwgInN1Yl90eXBlIjogIlN0YW5kYXJkIiwgInNxbCI6ICJzZWxlY3QgKiBmcm9tICBcImV4cDM2MC1jb2QtY3VzdG9tZXJcIi5hY3Rpb25zYXVkaXQgbGltaXQgMjsiLCAic3FsQ29udGV4dCI6IFtdLCAicmVmZXJlbmNlcyI6IHt9LCAicGF0aCI6IFsiT2N1bGFyIiwgInRlc3RfY2hhcmFuIiwgInRlc3Rfdl8yMzI2Il0sICJ2ZXJzaW9uIjogInYxIiwgImNyZWF0ZWRfZGF0ZSI6ICIyMDI2LTA3LTIzIDE1OjM5OjE5LjgzMDYxNyIsICJ1cGRhdGVkX2RhdGUiOiAiMjAyNi0wNy0yMyAxNTozOToxOS44MzA2MTciLCAiY3JlYXRlZF9ieSI6IG51bGwsICJhdXRvX3ZlcnNpb24iOiB0cnVlLCAic25hcHNob3RfbW9kZSI6IHRydWUsICJvdmVyd3JpdGVfbW9kZSI6IGZhbHNlLCAic25hcHNob3RfaWQiOiBudWxsLCAic2NoZWR1bGUiOiBudWxsLCAibGluZWFnZSI6IG51bGwsICJzdGF0dXMiOiAiYWN0aXZlIiwgInNjaGVtYV9oYXNoIjogIjQ0YzVjMDRkOTBhNCJ9
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_v_dataset_230726",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-23T09:25:35.169674",
|
||||
"created_by": null,
|
||||
"schema_hash": "16d7910d16bc",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_v_dataset_230726"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-23T09:25:35.169674",
|
||||
"updated_at": "2026-07-23T09:25:35.169674",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_v_dataset_230726", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit LIMIT 2;", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_v_dataset_230726"], "version": "v1", "created_date": "2026-07-23 09:25:35.169674", "updated_date": "2026-07-23 09:25:35.169674", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "16d7910d16bc"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_v_ds_1",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-23T13:01:10.745645",
|
||||
"created_by": null,
|
||||
"schema_hash": "474fb9bc0a29",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_v_ds_1"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-23T13:01:10.745645",
|
||||
"updated_at": "2026-07-23T13:01:10.745645",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_v_ds_1", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".customermapping limit 2;", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_v_ds_1"], "version": "v1", "created_date": "2026-07-23 13:01:10.745645", "updated_date": "2026-07-23 13:01:10.745645", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "474fb9bc0a29"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "test_vaishnavi_3",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "testing",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-23T12:33:00.216838",
|
||||
"created_by": null,
|
||||
"schema_hash": "474fb9bc0a29",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "test_vaishnavi_3"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-23T12:33:00.216838",
|
||||
"updated_at": "2026-07-23T12:33:00.216838",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "test_vaishnavi_3", "description": "testing", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".customermapping limit 2;", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "test_vaishnavi_3"], "version": "v1", "created_date": "2026-07-23 12:33:00.216838", "updated_date": "2026-07-23 12:33:00.216838", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "474fb9bc0a29"}
|
||||
@@ -1,26 +0,0 @@
|
||||
{
|
||||
"name": "testds241",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-24T05:38:01.411790",
|
||||
"created_by": null,
|
||||
"schema_hash": "79a17f89b99c",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testds241",
|
||||
"updated_at": "2026-07-24T05:39:15.660190"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-24T05:38:01.411790",
|
||||
"updated_at": "2026-07-24T05:39:15.660190",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testds241", "description": "testds241", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-cod-customer\".actionsaudit", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testds241"], "version": "v1", "created_date": "2026-07-24 05:39:15.660190", "updated_date": "2026-07-24 05:39:15.660190", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": "20260724_053915", "schedule": null, "lineage": null, "status": "active", "schema_hash": "79a17f89b99c"}
|
||||
@@ -1,26 +0,0 @@
|
||||
{
|
||||
"name": "testing",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T05:16:24.156402",
|
||||
"created_by": null,
|
||||
"schema_hash": "6ec727aad414",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing",
|
||||
"updated_at": "2026-07-22T05:37:53.696331"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T05:16:24.156402",
|
||||
"updated_at": "2026-07-22T05:37:53.696331",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testing", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit\n", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testing"], "version": "v1", "created_date": "2026-07-22 05:39:27.537040", "updated_date": "2026-07-22 05:39:27.537040", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": "20260722_053927", "schedule": null, "lineage": null, "status": "active", "schema_hash": "6ec727aad414"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing12",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-27T06:54:52.618875",
|
||||
"created_by": null,
|
||||
"schema_hash": "6244490a0e95",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing12"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-27T06:54:52.618875",
|
||||
"updated_at": "2026-07-27T06:54:52.618875",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testing12", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".customermapping", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testing12"], "version": "v1", "created_date": "2026-07-27 06:54:52.618875", "updated_date": "2026-07-27 06:54:52.618875", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6244490a0e95"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing_1",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T05:36:41.264740",
|
||||
"created_by": null,
|
||||
"schema_hash": "6ec727aad414",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing_1"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T05:36:41.264740",
|
||||
"updated_at": "2026-07-22T05:36:41.264740",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testing_1", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit\n", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testing_1"], "version": "v1", "created_date": "2026-07-22 05:36:41.264740", "updated_date": "2026-07-22 05:36:41.264740", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6ec727aad414"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing_c",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-27T06:55:38.766187",
|
||||
"created_by": null,
|
||||
"schema_hash": "6244490a0e95",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing_c"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-27T06:55:38.766187",
|
||||
"updated_at": "2026-07-27T06:55:38.766187",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testing_c", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".customermapping", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testing_c"], "version": "v1", "created_date": "2026-07-27 06:55:38.766187", "updated_date": "2026-07-27 06:55:38.766187", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6244490a0e95"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "testing_v",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "testing",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T05:59:02.120520",
|
||||
"created_by": null,
|
||||
"schema_hash": "cbbc63fa5c32",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "testing_v"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T05:59:02.120520",
|
||||
"updated_at": "2026-07-22T05:59:02.120520",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "testing_v", "description": "testing", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-cod-customer\".actionsaudit;", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "testing_v"], "version": "v1", "created_date": "2026-07-22 05:59:02.120520", "updated_date": "2026-07-22 05:59:02.120520", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "cbbc63fa5c32"}
|
||||
@@ -1,25 +0,0 @@
|
||||
{
|
||||
"name": "virtual1",
|
||||
"type": "VIRTUAL_DATASET",
|
||||
"sub_type": "Standard",
|
||||
"description": "",
|
||||
"latest_version": "v1",
|
||||
"versions": [
|
||||
{
|
||||
"version": "v1",
|
||||
"created_at": "2026-07-22T07:21:13.118407",
|
||||
"created_by": null,
|
||||
"schema_hash": "6ec727aad414",
|
||||
"status": "active",
|
||||
"is_current": true,
|
||||
"breaking_change": false,
|
||||
"change_notes": null,
|
||||
"dremio_view_name": "virtual1"
|
||||
}
|
||||
],
|
||||
"schedule": null,
|
||||
"created_at": "2026-07-22T07:21:13.118407",
|
||||
"updated_at": "2026-07-22T07:21:13.118407",
|
||||
"created_by": null,
|
||||
"status": "active"
|
||||
}
|
||||
@@ -1 +0,0 @@
|
||||
{"entityType": "dataset", "name": "virtual1", "description": "", "type": "VIRTUAL_DATASET", "sub_type": "Standard", "sql": "SELECT * FROM \"exp360-customer\".actionsaudit\n", "sqlContext": [], "references": {}, "path": ["Ocular", "test_charan", "virtual1"], "version": "v1", "created_date": "2026-07-22 07:21:13.118407", "updated_date": "2026-07-22 07:21:13.118407", "created_by": null, "auto_version": true, "snapshot_mode": true, "overwrite_mode": false, "snapshot_id": null, "schedule": null, "lineage": null, "status": "active", "schema_hash": "6ec727aad414"}
|
||||
147
test/main.py
147
test/main.py
@@ -1,147 +0,0 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
@@ -1,167 +0,0 @@
|
||||
import marimo
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
app = marimo.App()
|
||||
|
||||
|
||||
@app.cell
|
||||
def init():
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
return LOGGER, collect_metrics, log_info, materialization, os, spark, time
|
||||
|
||||
|
||||
@app.cell
|
||||
def finalize(
|
||||
LOGGER,
|
||||
collect_metrics,
|
||||
log_info,
|
||||
materialization,
|
||||
os,
|
||||
spark,
|
||||
time,
|
||||
):
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
return
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app.run()
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Notebook","metadata":{"name":"test","description":"test","runtime":"spark"},"spec":{"ui":{},"blocks":[]}}
|
||||
147
test123/main.py
147
test123/main.py
@@ -1,147 +0,0 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test123'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
@@ -1,167 +0,0 @@
|
||||
import marimo
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
app = marimo.App()
|
||||
|
||||
|
||||
@app.cell
|
||||
def init():
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test123'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
return LOGGER, collect_metrics, log_info, materialization, os, spark, time
|
||||
|
||||
|
||||
@app.cell
|
||||
def finalize(
|
||||
LOGGER,
|
||||
collect_metrics,
|
||||
log_info,
|
||||
materialization,
|
||||
os,
|
||||
spark,
|
||||
time,
|
||||
):
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
return
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app.run()
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Notebook","metadata":{"name":"test123","description":"test123","runtime":"spark"},"spec":{"ui":{},"blocks":[]}}
|
||||
158
test14/main.py
158
test14/main.py
@@ -1,158 +0,0 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test14'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
code_transform__0_start_time=time.time()
|
||||
|
||||
print('hellos')
|
||||
|
||||
code_transform__0_end_time=time.time()
|
||||
|
||||
code_transform__0_dependency_key="code_transform__0"
|
||||
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
@@ -1,181 +0,0 @@
|
||||
import marimo
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
app = marimo.App()
|
||||
|
||||
|
||||
@app.cell
|
||||
def init():
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'test14'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
return LOGGER, collect_metrics, log_info, materialization, os, spark, time
|
||||
|
||||
|
||||
@app.cell
|
||||
def code_transform__0(time):
|
||||
|
||||
code_transform__0_start_time=time.time()
|
||||
|
||||
print('hellos')
|
||||
|
||||
code_transform__0_end_time=time.time()
|
||||
|
||||
code_transform__0_dependency_key="code_transform__0"
|
||||
|
||||
return
|
||||
|
||||
|
||||
@app.cell
|
||||
def finalize(
|
||||
LOGGER,
|
||||
collect_metrics,
|
||||
log_info,
|
||||
materialization,
|
||||
os,
|
||||
spark,
|
||||
time,
|
||||
):
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
return
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app.run()
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"VisualBuilder","metadata":{"name":"test14","description":" ","runtime":"spark"},"spec":{"ui":{"edges":[],"nodes":[{"id":"code-transform__0","type":"workflowNode","position":{"x":1448.8333333333335,"y":-438},"data":{"nodeType":"code-transform","id":"code-transform__0"},"measured":{"width":240,"height":112},"selected":true}],"nodesData":{"code-transform__0":{"name":"code_transform__0","type":"CodeTransform","language":"python","datasource":"","code":"print('hellos')","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":[]}}},"blocks":[{"name":"code_transform__0","type":"CodeTransform","options":{},"language":"python","datasource":"","code":"print('hellos')","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":[]}]}}
|
||||
115
testws11/main.py
115
testws11/main.py
@@ -1,115 +0,0 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit
|
||||
from pyspark.sql.types import StringType, IntegerType
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
from dremio.flight.endpoint import DremioFlightEndpoint
|
||||
from dremio.flight.query import DremioFlightEndpointQuery
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'testws11'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options={'key': secrets.get('S3_ACCESS_KEY'), 'secret': secrets.get('S3_SECRET_KEY'), 'region': secrets.get('S3_REGION')})
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
|
||||
"spark.jars.packages": "com.amazonaws:aws-java-sdk-bundle:1.12.262,com.github.ben-manes.caffeine:caffeine:3.2.0,org.apache.iceberg:iceberg-aws-bundle:1.8.1,org.apache.iceberg:iceberg-common:1.8.1,org.apache.iceberg:iceberg-core:1.8.1,org.apache.iceberg:iceberg-spark:1.8.1,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.11.901,org.apache.hadoop:hadoop-common:3.3.4,org.apache.hadoop:hadoop-cloud-storage:3.3.4,org.apache.hadoop:hadoop-client-runtime:3.3.4,org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.8.1,org.projectnessie.nessie-integrations:nessie-spark-extensions-3.5_2.12:0.103.2,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.2,za.co.absa.cobrix:spark-cobol_2.12:2.8.0,ch.cern.sparkmeasure:spark-measure_2.12:0.26"
|
||||
}
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False'}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
spark.stop()
|
||||
@@ -1,127 +0,0 @@
|
||||
import marimo
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
app = marimo.App()
|
||||
|
||||
|
||||
@app.cell
|
||||
def init():
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit
|
||||
from pyspark.sql.types import StringType, IntegerType
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
from dremio.flight.endpoint import DremioFlightEndpoint
|
||||
from dremio.flight.query import DremioFlightEndpointQuery
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'testws11'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options={'key': secrets.get('S3_ACCESS_KEY'), 'secret': secrets.get('S3_SECRET_KEY'), 'region': secrets.get('S3_REGION')})
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
|
||||
"spark.jars.packages": "com.amazonaws:aws-java-sdk-bundle:1.12.262,com.github.ben-manes.caffeine:caffeine:3.2.0,org.apache.iceberg:iceberg-aws-bundle:1.8.1,org.apache.iceberg:iceberg-common:1.8.1,org.apache.iceberg:iceberg-core:1.8.1,org.apache.iceberg:iceberg-spark:1.8.1,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.11.901,org.apache.hadoop:hadoop-common:3.3.4,org.apache.hadoop:hadoop-cloud-storage:3.3.4,org.apache.hadoop:hadoop-client-runtime:3.3.4,org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.8.1,org.projectnessie.nessie-integrations:nessie-spark-extensions-3.5_2.12:0.103.2,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.2,za.co.absa.cobrix:spark-cobol_2.12:2.8.0,ch.cern.sparkmeasure:spark-measure_2.12:0.26"
|
||||
}
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
return LOGGER, collect_metrics, log_info, materialization, spark, time
|
||||
|
||||
|
||||
@app.cell
|
||||
def finalize(LOGGER, collect_metrics, log_info, materialization, spark, time):
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False'}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
spark.stop()
|
||||
return
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app.run()
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Notebook","metadata":{"name":"testws11","description":"testws11","runtime":"spark"},"spec":{"ui":{},"blocks":[]}}
|
||||
@@ -1,147 +0,0 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'testws241'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
@@ -1,167 +0,0 @@
|
||||
import marimo
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
app = marimo.App()
|
||||
|
||||
|
||||
@app.cell
|
||||
def init():
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit, input_file_name
|
||||
from pyspark.sql.types import StringType, IntegerType, MapType, StructType,StructField
|
||||
from postal.parser import parse_address
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim, lpad, format_number
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
import orjson
|
||||
|
||||
from ocular_ai_sdk import OcularClient
|
||||
from ocular_ai_sdk.exceptions import (
|
||||
OcularSDKException,
|
||||
AuthenticationError,
|
||||
ResourceNotFoundError
|
||||
)
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
import ssl
|
||||
from urllib.request import Request, urlopen
|
||||
from urllib.parse import urlencode
|
||||
from urllib.error import HTTPError
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'testws241'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
client = OcularClient(
|
||||
pat_token=secrets.get('OCULAR_AI_PAT_TOKEN')
|
||||
)
|
||||
|
||||
if 'AZURE_SERVICE_PRINCIPAL' in secrets:
|
||||
_storage_options=orjson.loads(secrets['AZURE_SERVICE_PRINCIPAL'])
|
||||
else:
|
||||
_storage_options = {
|
||||
'key': secrets.get('S3_ACCESS_KEY'),
|
||||
'secret': secrets.get('S3_SECRET_KEY'),
|
||||
'region': secrets.get('S3_REGION')
|
||||
}
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options=_storage_options)
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.jars.ivy": "/opt/spark/.ivy2/",
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
|
||||
}
|
||||
|
||||
if filesystemManager.storage_type == SupportedFilesystemType.AZUREBLOB:
|
||||
_params[f"fs.azure.account.auth.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "OAuth"
|
||||
_params[f"fs.azure.account.oauth.provider.type.{_storage_options['account_name']}.dfs.core.windows.net"] = "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"
|
||||
_params[f"fs.azure.account.oauth2.client.id.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_id']
|
||||
_params[f"fs.azure.account.oauth2.client.secret.{_storage_options['account_name']}.dfs.core.windows.net"] = _storage_options['client_secret']
|
||||
_params[f"fs.azure.account.oauth2.client.endpoint.{_storage_options['account_name']}.dfs.core.windows.net"] = f"https://login.microsoftonline.com/{_storage_options['tenant_id']}/oauth2/v2.0/token"
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
return LOGGER, collect_metrics, log_info, materialization, os, spark, time
|
||||
|
||||
|
||||
@app.cell
|
||||
def finalize(
|
||||
LOGGER,
|
||||
collect_metrics,
|
||||
log_info,
|
||||
materialization,
|
||||
os,
|
||||
spark,
|
||||
time,
|
||||
):
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False', 'execution_order': os.environ.get('EXECUTION_ORDER')}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
if os.getenv('EXECUTION_ENVIRONMENT'):
|
||||
spark.stop()
|
||||
return
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
app.run()
|
||||
@@ -1 +0,0 @@
|
||||
{"version":"v1alpha","kind":"Notebook","metadata":{"name":"testws241","description":"testws241","runtime":"spark"},"spec":{"ui":{},"blocks":[]}}
|
||||
Reference in New Issue
Block a user