Add testws11/main.py.notebook, testws11/main.py, testws11/main.workflow
This commit is contained in:
115
testws11/main.py
Normal file
115
testws11/main.py
Normal file
@@ -0,0 +1,115 @@
|
||||
|
||||
__generated_with = "0.13.15"
|
||||
|
||||
# %%
|
||||
|
||||
import sys
|
||||
import time
|
||||
from pyspark.sql.utils import AnalysisException
|
||||
sys.path.append('/opt/spark/work-dir/')
|
||||
from workflow_templates.spark.udf_manager import bootstrap_udfs
|
||||
from util import get_logger, observe_metrics, collect_metrics, log_info, log_error, forgiving_serializer
|
||||
from pyspark.sql.functions import udf
|
||||
from pyspark.sql.functions import count, expr, lit
|
||||
from pyspark.sql.types import StringType, IntegerType
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from pyspark import SparkConf, Row
|
||||
from pyspark.sql import SparkSession
|
||||
from pyspark.sql.observation import Observation
|
||||
from pyspark import StorageLevel
|
||||
import os
|
||||
import pandas as pd
|
||||
import polars as pl
|
||||
import pyarrow as pa
|
||||
from pyspark.sql.functions import approx_count_distinct, avg, collect_list, collect_set, corr, count, countDistinct, covar_pop, covar_samp, first, kurtosis, last, max, mean, min, skewness, stddev, stddev_pop, stddev_samp, sum, var_pop, var_samp, variance,expr,to_json,struct, date_format, col, lit, when, regexp_replace, ltrim
|
||||
from functools import reduce
|
||||
from handle_structs_or_arrays import preprocess_then_expand
|
||||
import requests
|
||||
from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
from jinja2 import Template
|
||||
import json
|
||||
|
||||
|
||||
from secrets_manager import SecretsManager
|
||||
|
||||
from WorkflowManager import WorkflowDSL, WorkflowManager
|
||||
from KnowledgebaseManager import KnowledgebaseManager
|
||||
from gitea_client import GiteaClient, WorkspaceVersionedContent
|
||||
from FilesystemManager import FilesystemManager, SupportedFilesystemType
|
||||
from Materialization import Materialization
|
||||
|
||||
from dremio.flight.endpoint import DremioFlightEndpoint
|
||||
from dremio.flight.query import DremioFlightEndpointQuery
|
||||
|
||||
init_start_time=time.time()
|
||||
|
||||
LOGGER = get_logger()
|
||||
alias_str='abcdefghijklmnopqrstuvwxyz'
|
||||
workspace = os.getenv('WORKSPACE') or 'test_charan'
|
||||
workflow = 'testws11'
|
||||
execution_environment = os.getenv('EXECUTION_ENVIRONMENT') or 'CLUSTER'
|
||||
|
||||
job_id = os.getenv("EXECUTION_ID") or str(uuid.uuid4())
|
||||
retry_job_id = os.getenv("RETRY_EXECUTION_ID") or ''
|
||||
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}', Retry Job Id: '{retry_job_id}'")
|
||||
|
||||
sm = SecretsManager(os.getenv('SECRET_MANAGER_URL'), os.getenv('SECRET_MANAGER_NAMESPACE'), os.getenv('SECRET_MANAGER_ENV'), os.getenv('SECRET_MANAGER_TOKEN'))
|
||||
secrets = sm.list_secrets(workspace)
|
||||
|
||||
gitea_client=GiteaClient(os.getenv('GITEA_HOST'), os.getenv('GITEA_TOKEN'), os.getenv('GITEA_OWNER') or 'gitea_admin', os.getenv('GITEA_REPO') or 'tenant1')
|
||||
workspaceVersionedContent=WorkspaceVersionedContent(gitea_client)
|
||||
|
||||
filesystemManager = FilesystemManager.create(secrets.get('LAKEHOUSE_BUCKET'), storage_options={'key': secrets.get('S3_ACCESS_KEY'), 'secret': secrets.get('S3_SECRET_KEY'), 'region': secrets.get('S3_REGION')})
|
||||
if retry_job_id:
|
||||
logs = Materialization.get_execution_history_by_job_id(filesystemManager, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, retry_job_id, selected_components=['finalize']).to_dicts()
|
||||
if len(logs) == 1 and logs[0].get('metrics').get('execute_status') == 'SUCCESS':
|
||||
log_info(LOGGER, f"Workspace: '{workspace}', Workflow: '{workflow}', Execution Environment: '{execution_environment}', Job Id: '{job_id}' - Retry Job Id: '{retry_job_id}' was already successful. Hence exiting to forward processing to next in chain.")
|
||||
sys.exit(0)
|
||||
|
||||
_conf = SparkConf()
|
||||
_params = {
|
||||
"spark.hadoop.fs.s3a.access.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.secret.key": secrets.get(''),
|
||||
"spark.hadoop.fs.s3a.aws.region": secrets.get("S3_REGION") or "",
|
||||
"spark.sql.catalog.dremio.warehouse" : secrets.get('LAKEHOUSE_BUCKET'),
|
||||
"spark.hadoop.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.hadoop.fs.s3.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
|
||||
"spark.sql.catalog.dremio" : "org.apache.iceberg.spark.SparkCatalog",
|
||||
"spark.sql.catalog.dremio.type" : "hadoop",
|
||||
"spark.hadoop.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.s3.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem",
|
||||
"spark.hadoop.fs.gs.impl": "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem",
|
||||
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
|
||||
"spark.jars.packages": "com.amazonaws:aws-java-sdk-bundle:1.12.262,com.github.ben-manes.caffeine:caffeine:3.2.0,org.apache.iceberg:iceberg-aws-bundle:1.8.1,org.apache.iceberg:iceberg-common:1.8.1,org.apache.iceberg:iceberg-core:1.8.1,org.apache.iceberg:iceberg-spark:1.8.1,org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.11.901,org.apache.hadoop:hadoop-common:3.3.4,org.apache.hadoop:hadoop-cloud-storage:3.3.4,org.apache.hadoop:hadoop-client-runtime:3.3.4,org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.8.1,org.projectnessie.nessie-integrations:nessie-spark-extensions-3.5_2.12:0.103.2,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.2,za.co.absa.cobrix:spark-cobol_2.12:2.8.0,ch.cern.sparkmeasure:spark-measure_2.12:0.26"
|
||||
}
|
||||
|
||||
|
||||
|
||||
_conf.setAll(list(_params.items()))
|
||||
|
||||
spark = SparkSession.builder.appName(workspace).config(conf=_conf).getOrCreate()
|
||||
bootstrap_udfs(spark)
|
||||
|
||||
materialization = Materialization(spark, secrets.get('LAKEHOUSE_BUCKET'), workspace, workflow, job_id, retry_job_id, execution_environment, LOGGER)
|
||||
|
||||
init_dependency_key="init"
|
||||
|
||||
|
||||
init_end_time=time.time()
|
||||
|
||||
# %%
|
||||
|
||||
finalize_start_time=time.time()
|
||||
|
||||
metrics = {
|
||||
'data': collect_metrics(locals()),
|
||||
}
|
||||
materialization.materialized_execution_history({'finalize': {'execute_status': 'SUCCESS', 'fail_on_error': 'False'}, **metrics['data']})
|
||||
log_info(LOGGER, f"Workflow Data metrics: {metrics['data']}")
|
||||
|
||||
finalize_end_time=time.time()
|
||||
|
||||
spark.stop()
|
||||
Reference in New Issue
Block a user