diff --git a/service_request_metrics/main.py b/service_request_metrics/main.py new file mode 100644 index 0000000..bb0af07 --- /dev/null +++ b/service_request_metrics/main.py @@ -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() \ No newline at end of file diff --git a/service_request_metrics/main.py.notebook b/service_request_metrics/main.py.notebook new file mode 100644 index 0000000..cd333fc --- /dev/null +++ b/service_request_metrics/main.py.notebook @@ -0,0 +1,196 @@ +import marimo + +__generated_with = "0.13.15" +app = marimo.App() + + +@app.cell +def init(): + + import sys + sys.path.append('/opt/spark/work-dir/') + from workflow_templates.spark.udf_manager import bootstrap_udfs + from pyspark.sql.functions import udf + from pyspark.sql.functions import lit + from pyspark.sql.types import StringType, IntegerType + import uuid + from pathlib import Path + from pyspark import SparkConf, Row + from pyspark.sql import SparkSession + import os + import pandas as pd + import polars as pl + import pyarrow as pa + from pyspark.sql.functions import expr,to_json,col,struct + from functools import reduce + from handle_structs_or_arrays import preprocess_then_expand + import requests + 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 dremio.flight.endpoint import DremioFlightEndpoint + from dremio.flight.query import DremioFlightEndpointQuery + + + alias_str='abcdefghijklmnopqrstuvwxyz' + workspace = os.getenv('WORKSPACE') or 'exp360cust' + + job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4()) + + 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) + conf = SparkConf() + params = { + "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": "us-west-1", + "spark.sql.catalog.dremio.warehouse" : 's3://'+ (secrets.get('LAKEHOUSE_BUCKET') or ''), + "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.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions", + "spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem", + "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" + } + + + + conf.setAll(list(params.items())) + + spark = SparkSession.builder.appName(workspace).config(conf=conf).getOrCreate() + bootstrap_udfs(spark) + return expr, job_id, lit, preprocess_then_expand, reduce, spark + + +@app.cell +def ActionsAuditData(spark): + + + + ActionsAuditData_df = spark.read.table('dremio.actionsaudit') + ActionsAuditData_df.createOrReplaceTempView('ActionsAuditData_df') + return (ActionsAuditData_df,) + + +@app.cell +def data_mapper__1(ActionsAuditData_df, job_id, spark): + + _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") + + + data_mapper__1_df=spark.sql(("SELECT " + ', '.join(_data_mapper__1_select_clause) + " FROM ActionsAuditData_df").replace("{job_id}",f"'{job_id}'")) + data_mapper__1_df.createOrReplaceTempView("data_mapper__1_df") + return (data_mapper__1_df,) + + +@app.cell +def LatestServiceRequests(data_mapper__1_df, spark): + + print(data_mapper__1_df.columns) + 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))") + LatestServiceRequests_df.createOrReplaceTempView('LatestServiceRequests_df') + return (LatestServiceRequests_df,) + + +@app.cell +def aggregate__3( + LatestServiceRequests_df, + expr, + lit, + preprocess_then_expand, + reduce, +): + + + + + + + + + + _params = { + "datasource": "LatestServiceRequests", + "selectFunctions" : [{'fieldName': 'service_count', 'aggregationFunction': 'SUM(action_count)'}] + } + + _df_flat, _grouping_specs, _rewritten_selects = preprocess_then_expand( LatestServiceRequests_df, + group_expression="action_date, service_type", + cube="", + rollup="", + grouping_set="", + select_functions=[{'fieldName': 'service_count', 'aggregationFunction': 'SUM(action_count)'}] + ) + + _agg_exprs = [expr(f["aggregationFunction"]).alias(f["fieldName"]) + for f in _rewritten_selects + ] + + _all_group_cols = list({c for gs in _grouping_specs for c in gs}) + + _partials = [] + for _gs in _grouping_specs: + _gdf = _df_flat.groupBy(*_gs).agg(*_agg_exprs) + for _col in _all_group_cols: + if _col not in _gs: + _gdf = _gdf.withColumn(_col, lit(None)) + _partials.append(_gdf) + + + aggregate__3_df = reduce(lambda a, b: a.unionByName(b), _partials) + + aggregate__3_df.createOrReplaceTempView('aggregate__3_df') + + + + + return (aggregate__3_df,) + + +@app.cell +def ServiceRequestMetricsWriter(aggregate__3_df, spark): + + + + + _ServiceRequestMetricsWriter_fields_to_update = aggregate__3_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 aggregate__3_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) + + + return + + +if __name__ == "__main__": + app.run() diff --git a/service_request_metrics/main.workflow b/service_request_metrics/main.workflow new file mode 100644 index 0000000..cc04ac5 --- /dev/null +++ b/service_request_metrics/main.workflow @@ -0,0 +1 @@ +{"version":"v1alpha","kind":"VisualBuilder","metadata":{"name":"service_request_metrics","description":"Visual builder workflow for service request metrics","runtime":"spark"},"spec":{"ui":{"edges":[{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-reader__0","target":"data-mapper__1","id":"xy-edge__data-reader__0-data-mapper__1"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"data-mapper__1","target":"filter__2","id":"xy-edge__data-mapper__1-filter__2"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"filter__2","target":"aggregate__3","id":"xy-edge__filter__2-aggregate__3"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"aggregate__3","target":"code-transform__0","id":"xy-edge__aggregate__3-code-transform__0"},{"style":{"stroke":"var(--color-primary-lighter)","strokeWidth":1},"source":"code-transform__0","target":"data-writer__4","id":"xy-edge__code-transform__0-data-writer__4"}],"nodes":[{"id":"data-reader__0","type":"workflowNode","position":{"x":589,"y":131},"data":{"nodeType":"data-reader","id":"data-reader__0"},"measured":{"width":240,"height":144},"selected":false,"dragging":false,"style":{}},{"id":"data-mapper__1","type":"workflowNode","position":{"x":1044.5,"y":131.5},"data":{"nodeType":"data-mapper","id":"data-mapper__1"},"measured":{"width":240,"height":144},"selected":false,"dragging":false,"style":{}},{"id":"filter__2","type":"workflowNode","position":{"x":1521.25,"y":125.75},"data":{"nodeType":"filter","id":"filter__2"},"measured":{"width":240,"height":168},"selected":false,"dragging":false,"style":{}},{"id":"aggregate__3","type":"workflowNode","position":{"x":2003.25,"y":121.75},"data":{"nodeType":"aggregate","id":"aggregate__3"},"measured":{"width":240,"height":172},"selected":false,"dragging":false,"style":{}},{"id":"data-writer__4","type":"workflowNode","position":{"x":2819.27194734189,"y":74.1320621195083},"data":{"nodeType":"data-writer","id":"data-writer__4"},"measured":{"width":240,"height":200},"selected":false,"dragging":false,"style":{}},{"id":"code-transform__0","type":"workflowNode","position":{"x":2411.932611772461,"y":81.33940728526409},"data":{"nodeType":"code-transform","id":"code-transform__0"},"measured":{"width":240,"height":144},"selected":true,"dragging":false,"style":{}}],"nodesData":{"data-reader__0":{"typeLabel":"Spark","isDefault":false,"name":"ActionsAuditData","type":"SparkReader","format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"actionsaudit","connectedComponents":[],"columns":[]},"data-mapper__1":{"name":"data_mapper__1","type":"DataMapping","fromDataReader":"ActionsAuditData","includeExistingColumns":false,"toSchema":[{"fieldName":"action_date","valueExpression":"DATE(action_date)"},{"fieldName":"service_type","valueExpression":"sub_category"},{"fieldName":"action_count","valueExpression":"action_count"}],"additionalData":{"isGlossaryAssisted":false,"selectedSourceGlossary":"","selectedTargetGlossary":"","manualMappings":[{"id":"mapping-1753983206926","newFieldName":"action_date","mappingType":"expression","value":"DATE(action_date)"},{"id":"mapping-1753983229372","newFieldName":"service_type","mappingType":"sourceColumn","value":"sub_category"},{"id":"mapping-1753983263800","newFieldName":"action_count","mappingType":"sourceColumn","value":"action_count"}],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"after":"","totalSourceTerms":0,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[]},"isDefault":false,"datasource":"ActionsAuditData","connectedComponents":["ActionsAuditData"]},"filter__2":{"name":"LatestServiceRequests","type":"Filter","datasource":"data_mapper__1","condition":"action_date >= COALESCE((SELECT MAX(DATE(action_date)) FROM dremio.servicemetrics), (SELECT MIN(action_date) FROM data_mapper__1_df))","isDefault":false},"aggregate__3":{"name":"aggregate__3","type":"SQLAggregation","datasource":"LatestServiceRequests","groupByParams":{"group_expression":"action_date, service_type"},"selectFunctions":[{"fieldName":"service_count","aggregationFunction":"sum('action_count')"}],"isDefault":false,"materialization_strategy":"NONE","fail_on_error":true,"connectedComponents":["LatestServiceRequests"]},"data-writer__4":{"name":"ServiceRequestMetricsWriter","type":"IcebergWriter","iceberg_catalog":"dremio","warehouse_directory":"servicemetrics","datasource":"CheckpointOutput","mode":"merge","typeLabel":"Iceberg (Legacy)","unique_id":["action_date","service_type"],"isDefault":false,"connectedComponents":["CheckpointOutput"]},"code-transform__0":{"name":"CheckpointOutput","type":"CodeTransform","language":"python","datasource":"aggregate__3","code":"CheckpointOutput_df = {{datasource}}_df.localCheckpoint()\n{{datasource}}_df.persist()\nCheckpointOutput_df.createOrReplaceTempView(\"CheckpointOutput_df\")","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["aggregate__3"]}}},"blocks":[{"name":"ActionsAuditData","type":"SparkReader","options":{},"typeLabel":"Spark","isDefault":false,"format":"iceberg","credentials":{"accessKey":"S3_ACCESS_KEY","secretKey":"S3_SECRET_KEY"},"iceberg_catalog":"dremio","region":"us-west-1","table_name":"actionsaudit","connectedComponents":[],"columns":[]},{"name":"data_mapper__1","type":"DataMapping","options":{},"fromDataReader":"ActionsAuditData","includeExistingColumns":false,"toSchema":[{"fieldName":"action_date","valueExpression":"DATE(action_date)"},{"fieldName":"service_type","valueExpression":"sub_category"},{"fieldName":"action_count","valueExpression":"action_count"}],"additionalData":{"isGlossaryAssisted":false,"selectedSourceGlossary":"","selectedTargetGlossary":"","manualMappings":[{"id":"mapping-1753983206926","newFieldName":"action_date","mappingType":"expression","value":"DATE(action_date)"},{"id":"mapping-1753983229372","newFieldName":"service_type","mappingType":"sourceColumn","value":"sub_category"},{"id":"mapping-1753983263800","newFieldName":"action_count","mappingType":"sourceColumn","value":"action_count"}],"confirmedMapping":[],"mappingBatchKeysAfterConfirm":[],"after":"","totalSourceTerms":0,"selectedSourceSystem":"","selectedTargetSystem":"","selectedSourceLayout":"","selectedTargetLayout":"","selectedTargetLayoutFile":"","manualTargetMappings":[],"confirmedGlossaryMappings":[],"confirmedManualTargetMappings":[],"confirmedSSEMappings":[],"glossaryAssistedMappings":[],"additionalFieldMappings":[]},"isDefault":false,"datasource":"ActionsAuditData","connectedComponents":["ActionsAuditData"]},{"name":"LatestServiceRequests","type":"Filter","options":{},"datasource":"data_mapper__1","condition":"action_date >= COALESCE((SELECT MAX(DATE(action_date)) FROM dremio.servicemetrics), (SELECT MIN(action_date) FROM data_mapper__1_df))","isDefault":false},{"name":"aggregate__3","type":"SQLAggregation","options":{},"datasource":"LatestServiceRequests","groupByParams":{"group_expression":"action_date, service_type"},"selectFunctions":[{"fieldName":"service_count","aggregationFunction":"sum('action_count')"}],"isDefault":false,"materialization_strategy":"NONE","fail_on_error":true,"connectedComponents":["LatestServiceRequests"]},{"name":"ServiceRequestMetricsWriter","type":"IcebergWriter","options":{},"iceberg_catalog":"dremio","warehouse_directory":"servicemetrics","datasource":"CheckpointOutput","mode":"merge","typeLabel":"Iceberg (Legacy)","unique_id":["action_date","service_type"],"isDefault":false,"connectedComponents":["CheckpointOutput"]},{"name":"CheckpointOutput","type":"CodeTransform","options":{},"language":"python","datasource":"aggregate__3","code":"CheckpointOutput_df = {{datasource}}_df.localCheckpoint()\n{{datasource}}_df.persist()\nCheckpointOutput_df.createOrReplaceTempView(\"CheckpointOutput_df\")","materialization_strategy":"NONE","fail_on_error":true,"isDefault":false,"connectedComponents":["aggregate__3"]}]}} \ No newline at end of file