4 Commits

Author SHA1 Message Date
admin
a4a9c45ffc Workflow saved 2026-07-29 12:14:09 +00:00
admin
448aac4163 Workflow saved 2026-07-29 08:25:10 +00:00
admin
08b78e6aa8 v1 2026-07-24 14:15:45 +00:00
admin
5b6f1ad2ff v1 2026-07-24 14:15:09 +00:00
7 changed files with 7582 additions and 0 deletions

2879
forecasting_workflow/main.py Normal file

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

File diff suppressed because one or more lines are too long

View File

@@ -0,0 +1,394 @@
__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 'exp360-cus-uat'
workflow = 'service_request_metrics'
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('S3_ACCESS_KEY'),
"spark.hadoop.fs.s3a.secret.key": secrets.get('S3_SECRET_KEY'),
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "us-west-1",
"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.driver.extraJavaOptions": "--add-opens=java.base/java.nio=ALL-UNNAMED",
"spark.executor.extraJavaOptions": "--add-opens=java.base/java.nio=ALL-UNNAMED"
}
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()
# %%
ActionsAuditData_start_time=time.time()
ActionsAuditData_fail_on_error=""
try:
_ActionsAuditData_options = {
'jdbc':{
'dbtable': """actionsaudit""",
'url':secrets.get(''),
'driver':''
},
'kafka' : {
'kafka.bootstrap.servers' : secrets.get('OCULAR_KAFKA_BOOTSTRAP_SERVERS'),
'subscribe' : '',
'startingOffsets' : 'earliest'
},
'cobol' : {
'copybook' : '',
'encoding' : '',
'is_text': False,
'schema_retention_policy' : 'collapse_root'
}
}
_ActionsAuditData_reader = spark.read.format('iceberg')
ActionsAuditData_df = _ActionsAuditData_reader.load('dremio.actionsaudit')
ActionsAuditData_df, ActionsAuditData_observer = observe_metrics("ActionsAuditData_df", ActionsAuditData_df)
ActionsAuditData_df.createOrReplaceTempView('ActionsAuditData_df')
ActionsAuditData_dependency_key="ActionsAuditData"
ActionsAuditData_execute_status="SUCCESS"
except Exception as e:
ActionsAuditData_error = e
log_error(LOGGER, f"Component ActionsAuditData Failed", e)
ActionsAuditData_execute_status="ERROR"
raise e
finally:
ActionsAuditData_end_time=time.time()
# %%
data_mapper__1_start_time=time.time()
data_mapper__1_fail_on_error=""
try:
_data_mapper__1_select_clause=ActionsAuditData_df.columns if False else []
_data_mapper__1_select_clause.append('''DATE(action_date) AS action_date''')
_data_mapper__1_select_clause.append('''sub_category AS service_type''')
_data_mapper__1_select_clause.append('''action_count AS action_count''')
try:
data_mapper__1_df=spark.sql(("SELECT " + ', '.join(_data_mapper__1_select_clause) + " FROM ActionsAuditData_df").replace("{job_id}",f"'{job_id}'"))
except Exception as e:
data_mapper__1_df = ActionsAuditData_df.limit(0)
log_info(LOGGER, f"error while mapping the data :{e} " )
data_mapper__1_df, data_mapper__1_observer = observe_metrics("data_mapper__1_df", data_mapper__1_df)
data_mapper__1_df.createOrReplaceTempView("data_mapper__1_df")
data_mapper__1_dependency_key="data_mapper__1"
print(ActionsAuditData_dependency_key)
data_mapper__1_execute_status="SUCCESS"
except Exception as e:
data_mapper__1_error = e
log_error(LOGGER, f"Component data_mapper__1 Failed", e)
data_mapper__1_execute_status="ERROR"
raise e
finally:
data_mapper__1_end_time=time.time()
# %%
LatestServiceRequests_start_time=time.time()
print(data_mapper__1_df.columns)
LatestServiceRequests_fail_on_error=""
try:
try:
LatestServiceRequests_df = spark.sql("select * from data_mapper__1_df where action_date >= COALESCE((SELECT MAX(DATE(action_date)) FROM dremio.servicemetrics), (SELECT MIN(action_date) FROM data_mapper__1_df))")
except AnalysisException as e:
log_info(LOGGER, f"error while filtering data : {e}")
LatestServiceRequests_df = data_mapper__1_df.limit(0)
except Exception as e:
log_info(LOGGER, f"Unexpected error: {e}")
LatestServiceRequests_df = data_mapper__1_df.limit(0)
LatestServiceRequests_df, LatestServiceRequests_observer = observe_metrics("LatestServiceRequests_df", LatestServiceRequests_df)
LatestServiceRequests_df.createOrReplaceTempView('LatestServiceRequests_df')
LatestServiceRequests_dependency_key="LatestServiceRequests"
LatestServiceRequests_execute_status="SUCCESS"
except Exception as e:
LatestServiceRequests_error = e
log_error(LOGGER, f"Component LatestServiceRequests Failed", e)
LatestServiceRequests_execute_status="ERROR"
raise e
finally:
LatestServiceRequests_end_time=time.time()
# %%
aggregate__3_start_time=time.time()
aggregate__3_fail_on_error="True"
try:
aggregate__3_df = LatestServiceRequests_df.groupBy(
"action_date",
"service_type"
).agg(
sum('action_count').alias("service_count")
)
aggregate__3_df, aggregate__3_observer = observe_metrics("aggregate__3_df", aggregate__3_df)
aggregate__3_df.createOrReplaceTempView('aggregate__3_df')
aggregate__3_dependency_key="aggregate__3"
print(LatestServiceRequests_dependency_key)
aggregate__3_execute_status="SUCCESS"
except Exception as e:
aggregate__3_error = e
log_error(LOGGER, f"Component aggregate__3 Failed", e)
aggregate__3_execute_status="ERROR"
raise e
finally:
aggregate__3_end_time=time.time()
# %%
CheckpointOutput_start_time=time.time()
CheckpointOutput_df = aggregate__3_df.localCheckpoint()
# aggregate__3_df.persist()
CheckpointOutput_df.createOrReplaceTempView("CheckpointOutput_df")
CheckpointOutput_end_time=time.time()
CheckpointOutput_dependency_key="CheckpointOutput"
print(aggregate__3_dependency_key)
# %%
ServiceRequestMetricsWriter_start_time=time.time()
ServiceRequestMetricsWriter_fail_on_error=""
try:
_ServiceRequestMetricsWriter_options = {
'jdbc':{
'dbtable': 'servicemetrics',
'url':secrets.get(''),
'driver':'',
'stringtype': 'unspecified'
},
'kafka' : {
'kafka.bootstrap.servers' : secrets.get('OCULAR_KAFKA_BOOTSTRAP_SERVERS'),
'topic' : ''
}
}
_CheckpointOutput_df = CheckpointOutput_df
_ServiceRequestMetricsWriter_writer = _CheckpointOutput_df.write.format('iceberg').mode('append')
_ServiceRequestMetricsWriter_writer.save('dremio.servicemetrics')
ServiceRequestMetricsWriter_dependency_key="ServiceRequestMetricsWriter"
print(CheckpointOutput_dependency_key)
ServiceRequestMetricsWriter_execute_status="SUCCESS"
except Exception as e:
ServiceRequestMetricsWriter_error = e
log_error(LOGGER, f"Component ServiceRequestMetricsWriter Failed", e)
ServiceRequestMetricsWriter_execute_status="ERROR"
raise e
finally:
ServiceRequestMetricsWriter_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()

View File

@@ -0,0 +1,480 @@
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 'exp360-cus-uat'
workflow = 'service_request_metrics'
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('S3_ACCESS_KEY'),
"spark.hadoop.fs.s3a.secret.key": secrets.get('S3_SECRET_KEY'),
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "us-west-1",
"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.driver.extraJavaOptions": "--add-opens=java.base/java.nio=ALL-UNNAMED",
"spark.executor.extraJavaOptions": "--add-opens=java.base/java.nio=ALL-UNNAMED"
}
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 (
AnalysisException,
LOGGER,
collect_metrics,
job_id,
log_error,
log_info,
materialization,
observe_metrics,
os,
secrets,
spark,
sum,
time,
)
@app.cell
def ActionsAuditData(LOGGER, log_error, observe_metrics, secrets, spark, time):
ActionsAuditData_start_time=time.time()
ActionsAuditData_fail_on_error=""
try:
_ActionsAuditData_options = {
'jdbc':{
'dbtable': """actionsaudit""",
'url':secrets.get(''),
'driver':''
},
'kafka' : {
'kafka.bootstrap.servers' : secrets.get('OCULAR_KAFKA_BOOTSTRAP_SERVERS'),
'subscribe' : '',
'startingOffsets' : 'earliest'
},
'cobol' : {
'copybook' : '',
'encoding' : '',
'is_text': False,
'schema_retention_policy' : 'collapse_root'
}
}
_ActionsAuditData_reader = spark.read.format('iceberg')
ActionsAuditData_df = _ActionsAuditData_reader.load('dremio.actionsaudit')
ActionsAuditData_df, ActionsAuditData_observer = observe_metrics("ActionsAuditData_df", ActionsAuditData_df)
ActionsAuditData_df.createOrReplaceTempView('ActionsAuditData_df')
ActionsAuditData_dependency_key="ActionsAuditData"
ActionsAuditData_execute_status="SUCCESS"
except Exception as e:
ActionsAuditData_error = e
log_error(LOGGER, f"Component ActionsAuditData Failed", e)
ActionsAuditData_execute_status="ERROR"
raise e
finally:
ActionsAuditData_end_time=time.time()
return ActionsAuditData_dependency_key, ActionsAuditData_df
@app.cell
def data_mapper__1(
ActionsAuditData_dependency_key,
ActionsAuditData_df,
LOGGER,
job_id,
log_error,
log_info,
observe_metrics,
spark,
time,
):
data_mapper__1_start_time=time.time()
data_mapper__1_fail_on_error=""
try:
_data_mapper__1_select_clause=ActionsAuditData_df.columns if False else []
_data_mapper__1_select_clause.append('''DATE(action_date) AS action_date''')
_data_mapper__1_select_clause.append('''sub_category AS service_type''')
_data_mapper__1_select_clause.append('''action_count AS action_count''')
try:
data_mapper__1_df=spark.sql(("SELECT " + ', '.join(_data_mapper__1_select_clause) + " FROM ActionsAuditData_df").replace("{job_id}",f"'{job_id}'"))
except Exception as e:
data_mapper__1_df = ActionsAuditData_df.limit(0)
log_info(LOGGER, f"error while mapping the data :{e} " )
data_mapper__1_df, data_mapper__1_observer = observe_metrics("data_mapper__1_df", data_mapper__1_df)
data_mapper__1_df.createOrReplaceTempView("data_mapper__1_df")
data_mapper__1_dependency_key="data_mapper__1"
print(ActionsAuditData_dependency_key)
data_mapper__1_execute_status="SUCCESS"
except Exception as e:
data_mapper__1_error = e
log_error(LOGGER, f"Component data_mapper__1 Failed", e)
data_mapper__1_execute_status="ERROR"
raise e
finally:
data_mapper__1_end_time=time.time()
return (data_mapper__1_df,)
@app.cell
def LatestServiceRequests(
AnalysisException,
LOGGER,
data_mapper__1_df,
log_error,
log_info,
observe_metrics,
spark,
time,
):
LatestServiceRequests_start_time=time.time()
print(data_mapper__1_df.columns)
LatestServiceRequests_fail_on_error=""
try:
try:
LatestServiceRequests_df = spark.sql("select * from data_mapper__1_df where action_date >= COALESCE((SELECT MAX(DATE(action_date)) FROM dremio.servicemetrics), (SELECT MIN(action_date) FROM data_mapper__1_df))")
except AnalysisException as e:
log_info(LOGGER, f"error while filtering data : {e}")
LatestServiceRequests_df = data_mapper__1_df.limit(0)
except Exception as e:
log_info(LOGGER, f"Unexpected error: {e}")
LatestServiceRequests_df = data_mapper__1_df.limit(0)
LatestServiceRequests_df, LatestServiceRequests_observer = observe_metrics("LatestServiceRequests_df", LatestServiceRequests_df)
LatestServiceRequests_df.createOrReplaceTempView('LatestServiceRequests_df')
LatestServiceRequests_dependency_key="LatestServiceRequests"
LatestServiceRequests_execute_status="SUCCESS"
except Exception as e:
LatestServiceRequests_error = e
log_error(LOGGER, f"Component LatestServiceRequests Failed", e)
LatestServiceRequests_execute_status="ERROR"
raise e
finally:
LatestServiceRequests_end_time=time.time()
return LatestServiceRequests_dependency_key, LatestServiceRequests_df
@app.cell
def aggregate__3(
LOGGER,
LatestServiceRequests_dependency_key,
LatestServiceRequests_df,
log_error,
observe_metrics,
sum,
time,
):
aggregate__3_start_time=time.time()
aggregate__3_fail_on_error="True"
try:
aggregate__3_df = LatestServiceRequests_df.groupBy(
"action_date",
"service_type"
).agg(
sum('action_count').alias("service_count")
)
aggregate__3_df, aggregate__3_observer = observe_metrics("aggregate__3_df", aggregate__3_df)
aggregate__3_df.createOrReplaceTempView('aggregate__3_df')
aggregate__3_dependency_key="aggregate__3"
print(LatestServiceRequests_dependency_key)
aggregate__3_execute_status="SUCCESS"
except Exception as e:
aggregate__3_error = e
log_error(LOGGER, f"Component aggregate__3 Failed", e)
aggregate__3_execute_status="ERROR"
raise e
finally:
aggregate__3_end_time=time.time()
return aggregate__3_dependency_key, aggregate__3_df
@app.cell
def ServiceRequestMetricsWriter(
CheckpointOutput_dependency_key,
CheckpointOutput_df,
LOGGER,
log_error,
secrets,
time,
):
ServiceRequestMetricsWriter_start_time=time.time()
ServiceRequestMetricsWriter_fail_on_error=""
try:
_ServiceRequestMetricsWriter_options = {
'jdbc':{
'dbtable': 'servicemetrics',
'url':secrets.get(''),
'driver':'',
'stringtype': 'unspecified'
},
'kafka' : {
'kafka.bootstrap.servers' : secrets.get('OCULAR_KAFKA_BOOTSTRAP_SERVERS'),
'topic' : ''
}
}
_CheckpointOutput_df = CheckpointOutput_df
_ServiceRequestMetricsWriter_writer = _CheckpointOutput_df.write.format('iceberg').mode('append')
_ServiceRequestMetricsWriter_writer.save('dremio.servicemetrics')
ServiceRequestMetricsWriter_dependency_key="ServiceRequestMetricsWriter"
print(CheckpointOutput_dependency_key)
ServiceRequestMetricsWriter_execute_status="SUCCESS"
except Exception as e:
ServiceRequestMetricsWriter_error = e
log_error(LOGGER, f"Component ServiceRequestMetricsWriter Failed", e)
ServiceRequestMetricsWriter_execute_status="ERROR"
raise e
finally:
ServiceRequestMetricsWriter_end_time=time.time()
return
@app.cell
def CheckpointOutput(aggregate__3_dependency_key, aggregate__3_df, time):
CheckpointOutput_start_time=time.time()
CheckpointOutput_df = aggregate__3_df.localCheckpoint()
# aggregate__3_df.persist()
CheckpointOutput_df.createOrReplaceTempView("CheckpointOutput_df")
CheckpointOutput_end_time=time.time()
CheckpointOutput_dependency_key="CheckpointOutput"
print(aggregate__3_dependency_key)
return CheckpointOutput_dependency_key, CheckpointOutput_df
@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()

View File

@@ -0,0 +1,497 @@
__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,
run_component, apply_data_quality, compute_dq_stats, enforce_error_threshold,
build_dq_error_log, build_api_error_log, ERROR_LOG_SCHEMA, RetryConfig, with_retry,
app_scoped_error_code,
registry_error_code,
rewrite_response_body_json_access,
rewrite_response_body_json_access_if_json,
)
from exception_utils import (
ErrorMessage,
Severity,
ConnectionException,
AuthenticationException,
SSLException,
RateLimitException,
ServiceUnavailableException,
TimeoutException,
ValidationException,
SchemaMappingException,
ExpressionException,
MergeException,
ConfigurationException,
RetryExhaustedException,
format_exception,
mask_pii,
mask_pii_dict,
)
from py4j.protocol import Py4JJavaError
from component_error_handler import handle_analysis_error, handle_java_error, classify_java_error
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer, set_correlation_id, set_workflow_context
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
init_start_time=time.time()
LOGGER = get_logger()
alias_str='abcdefghijklmnopqrstuvwxyz'
workspace = os.getenv('WORKSPACE') or 'exp360uat'
workflow = 'service_request_metrics'
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 ''
correlation_id = job_id
set_correlation_id(correlation_id)
set_workflow_context(workspace=workspace, workflow=workflow, job_id=job_id, retry_job_id=retry_job_id, execution_environment=execution_environment)
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}', Correlation Id: '{correlation_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)
import dremio_operations
dremio_operations.configure(secrets)
import kb_query
kb_query.configure(secrets)
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('S3_ACCESS_KEY'),
"spark.hadoop.fs.s3a.secret.key": secrets.get('S3_SECRET_KEY'),
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "us-east-1",
"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()
# %%
ActionsAuditData_start_time=time.time()
ActionsAuditData_fail_on_error=""
try:
_ActionsAuditData_options = {
'jdbc':{
'dbtable': """actionsaudit""",
'url':secrets.get(''),
'driver':''
},
'kafka' : {
'kafka.bootstrap.servers' : secrets.get('OCULAR_KAFKA_BOOTSTRAP_SERVERS'),
'subscribe' : '',
'startingOffsets' : 'earliest'
},
'cobol' : {
'copybook' : '',
'encoding' : '',
'is_text': False,
'schema_retention_policy' : 'collapse_root'
}
}
_reader = spark.read.format('iceberg')
_ActionsAuditData_load_path = 'dremio.actionsaudit'
_ActionsAuditData_input_data = {
"component": "ActionsAuditData",
"format": "iceberg",
"iceberg_catalog": """dremio""",
"table_name": """actionsaudit""",
}
try:
ActionsAuditData_df = _reader.load(_ActionsAuditData_load_path)
ActionsAuditData_df = ActionsAuditData_df.withColumn("ActionsAuditData_input_file", input_file_name())
# Force partition evaluation to surface lazy errors (e.g. glob matches 0 files)
ActionsAuditData_df.rdd.getNumPartitions()
except AnalysisException as e:
handle_analysis_error(
e,
component_name="ActionsAuditData",
message=f"Failed to load source 'ActionsAuditData' ({_ActionsAuditData_load_path}): {e!s}",
job_id=job_id, workspace=workspace, workflow=workflow, execution_environment=execution_environment,
extra_details={"format": "iceberg", "load_path": _ActionsAuditData_load_path},
input_data=_ActionsAuditData_input_data,
)
except Py4JJavaError as e:
handle_java_error(
e,
component_name="ActionsAuditData",
operation="load",
format_name="iceberg",
path=_ActionsAuditData_load_path,
job_id=job_id, workspace=workspace, workflow=workflow, execution_environment=execution_environment,
input_data=_ActionsAuditData_input_data,
)
ActionsAuditData_df, ActionsAuditData_observer = observe_metrics("ActionsAuditData_df", ActionsAuditData_df)
ActionsAuditData_df.createOrReplaceTempView('ActionsAuditData_df')
ActionsAuditData_dependency_key="ActionsAuditData"
ActionsAuditData_execute_status="SUCCESS"
except Exception as e:
ActionsAuditData_error = e
log_error(LOGGER, f"Component ActionsAuditData Failed", e, component_name="ActionsAuditData")
ActionsAuditData_execute_status="ERROR"
raise e
finally:
ActionsAuditData_end_time=time.time()
# %%
data_mapper__1_start_time=time.time()
data_mapper__1_fail_on_error=""
try:
_data_mapper__1_select_clause=[]
_data_mapper__1_expr = """DATE(action_date)""".replace("input_file_name()", "input_file")
_data_mapper__1_expr = _data_mapper__1_expr.replace("_dq_source_file", "input_file")
if "." in _data_mapper__1_expr:
_data_mapper__1_expr = rewrite_response_body_json_access(_data_mapper__1_expr)
_data_mapper__1_select_clause.append(f"{_data_mapper__1_expr} AS action_date")
_data_mapper__1_expr = """sub_category""".replace("input_file_name()", "input_file")
_data_mapper__1_expr = _data_mapper__1_expr.replace("_dq_source_file", "input_file")
if "." in _data_mapper__1_expr:
_data_mapper__1_expr = rewrite_response_body_json_access(_data_mapper__1_expr)
_data_mapper__1_select_clause.append(f"{_data_mapper__1_expr} AS service_type")
_data_mapper__1_expr = """action_count""".replace("input_file_name()", "input_file")
_data_mapper__1_expr = _data_mapper__1_expr.replace("_dq_source_file", "input_file")
if "." in _data_mapper__1_expr:
_data_mapper__1_expr = rewrite_response_body_json_access(_data_mapper__1_expr)
_data_mapper__1_select_clause.append(f"{_data_mapper__1_expr} AS action_count")
_data_mapper__1_mapping_sql = ("SELECT " + ', '.join(_data_mapper__1_select_clause) + " FROM ActionsAuditData_df").replace("{job_id}", f"'{job_id}'")
_data_mapper__1_input_data = {
"component": "data_mapper__1",
"datasource": "ActionsAuditData",
"include_existing_columns": False,
"to_schema_field_count": 3,
}
try:
data_mapper__1_df = spark.sql(_data_mapper__1_mapping_sql)
except AnalysisException as e:
handle_analysis_error(
e,
component_name="data_mapper__1",
error_code="TRF-MAP-002",
exception_class=SchemaMappingException,
message=f"Spark analysis error during data_mapper__1 mapping: {e!s}",
job_id=job_id, workspace=workspace, workflow=workflow, execution_environment=execution_environment,
extra_details={"retry_job_id": retry_job_id or None, "sql_preview": _data_mapper__1_mapping_sql[:2000]},
input_data=_data_mapper__1_input_data,
)
except Py4JJavaError as e:
handle_java_error(
e,
component_name="data_mapper__1",
operation="mapping SQL",
format_name="sql",
job_id=job_id, workspace=workspace, workflow=workflow, execution_environment=execution_environment,
override_class=ExpressionException, override_code="TRF-EXP-001",
extra_details={"retry_job_id": retry_job_id or None, "sql_preview": _data_mapper__1_mapping_sql[:2000]},
input_data=_data_mapper__1_input_data,
)
data_mapper__1_df, data_mapper__1_observer = observe_metrics("data_mapper__1_df", data_mapper__1_df)
data_mapper__1_df.createOrReplaceTempView("data_mapper__1_df")
data_mapper__1_dependency_key="data_mapper__1"
print(ActionsAuditData_dependency_key)
data_mapper__1_execute_status="SUCCESS"
except Exception as e:
data_mapper__1_error = e
log_error(LOGGER, f"Component data_mapper__1 Failed", e, component_name="data_mapper__1")
data_mapper__1_execute_status="ERROR"
raise e
finally:
data_mapper__1_end_time=time.time()
# %%
LatestServiceRequests_start_time=time.time()
print(data_mapper__1_df.columns)
LatestServiceRequests_fail_on_error=""
try:
_LatestServiceRequests_condition = rewrite_response_body_json_access_if_json(data_mapper__1_df, """action_date >= COALESCE((SELECT MAX(DATE(action_date)) FROM dremio.servicemetrics), (SELECT MIN(action_date) FROM data_mapper__1_df))""")
LatestServiceRequests_df = spark.sql(f"select * from data_mapper__1_df where {_LatestServiceRequests_condition}")
LatestServiceRequests_df, LatestServiceRequests_observer = observe_metrics("LatestServiceRequests_df", LatestServiceRequests_df)
LatestServiceRequests_df.createOrReplaceTempView('LatestServiceRequests_df')
LatestServiceRequests_dependency_key="LatestServiceRequests"
LatestServiceRequests_execute_status="SUCCESS"
except Exception as e:
LatestServiceRequests_error = e
log_error(LOGGER, f"Component LatestServiceRequests Failed", e, component_name="LatestServiceRequests")
LatestServiceRequests_execute_status="ERROR"
raise e
finally:
LatestServiceRequests_end_time=time.time()
# %%
aggregate__3_start_time=time.time()
aggregate__3_fail_on_error="True"
try:
_aggregate__3_group_cols = []
_aggregate__3_group_cols.append(expr(rewrite_response_body_json_access_if_json(LatestServiceRequests_df, """action_date""")))
_aggregate__3_group_cols.append(expr(rewrite_response_body_json_access_if_json(LatestServiceRequests_df, """service_type""")))
aggregate__3_df = LatestServiceRequests_df.groupBy(*_aggregate__3_group_cols).agg(
sum('action_count').alias("service_count")
)
aggregate__3_df, aggregate__3_observer = observe_metrics("aggregate__3_df", aggregate__3_df)
aggregate__3_df.createOrReplaceTempView('aggregate__3_df')
aggregate__3_dependency_key="aggregate__3"
print(LatestServiceRequests_dependency_key)
aggregate__3_execute_status="SUCCESS"
except Exception as e:
aggregate__3_error = e
log_error(LOGGER, f"Component aggregate__3 Failed", e, component_name="aggregate__3")
aggregate__3_execute_status="ERROR"
raise e
finally:
aggregate__3_end_time=time.time()
# %%
CheckpointOutput_start_time=time.time()
CheckpointOutput_df = aggregate__3_df.localCheckpoint()
aggregate__3_df.persist()
CheckpointOutput_df.createOrReplaceTempView("CheckpointOutput_df")
CheckpointOutput_end_time=time.time()
CheckpointOutput_dependency_key="CheckpointOutput"
print(aggregate__3_dependency_key)
# %%
ServiceRequestMetricsWriter_start_time=time.time()
ServiceRequestMetricsWriter_fail_on_error=""
try:
_ServiceRequestMetricsWriter_fields_to_update = CheckpointOutput_df.columns
_ServiceRequestMetricsWriter_set_clause=[]
_ServiceRequestMetricsWriter_unique_key_clause= []
for _key in ['action_date', 'service_type']:
_ServiceRequestMetricsWriter_unique_key_clause.append(f't.{_key} = s.{_key}')
for _field in _ServiceRequestMetricsWriter_fields_to_update:
if(_field not in _ServiceRequestMetricsWriter_unique_key_clause):
_ServiceRequestMetricsWriter_set_clause.append(f't.{_field} = s.{_field}')
_merge_query = '''
MERGE INTO dremio.servicemetrics t
USING CheckpointOutput_df s
ON ''' + ' AND '.join(_ServiceRequestMetricsWriter_unique_key_clause) + ''' WHEN MATCHED THEN
UPDATE SET ''' + ', '.join(_ServiceRequestMetricsWriter_set_clause) + ' WHEN NOT MATCHED THEN INSERT *'
spark.sql(_merge_query)
ServiceRequestMetricsWriter_dependency_key="ServiceRequestMetricsWriter"
print(CheckpointOutput_dependency_key)
ServiceRequestMetricsWriter_execute_status="SUCCESS"
except Exception as e:
ServiceRequestMetricsWriter_error = e
log_error(LOGGER, f"Component ServiceRequestMetricsWriter Failed", e, component_name="ServiceRequestMetricsWriter")
ServiceRequestMetricsWriter_execute_status="ERROR"
raise e
finally:
ServiceRequestMetricsWriter_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 (correlation_id={metrics['data'].get('correlation_id')}): {metrics['data']}")
finalize_end_time=time.time()
if os.getenv('EXECUTION_ENVIRONMENT'):
spark.stop()

File diff suppressed because one or more lines are too long