__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()