## REQUIRED MODULES AND PACKAGESimport osimport randomimport timefrom humanfriendly import format_timespanimport numpy as npimport pandas as pdfrom datetime import datetimeimport asyncioimport aiomysqlfrom pathlib import Pathimport json## PREFECT MODULESfrom prefect import flow, taskfrom prefect.blocks.system import Secret## DATABASE CONNECTION PACKAGESimport pymysqlfrom urllib.parse import quote_plusfrom sqlalchemy import create_enginefrom sshtunnel import SSHTunnelForwarder## CUSTOM FUNCTIONSfrom scheduled_update_extractors import connect_and_run_queriesfrom scheduled_update_transformers import transform_all_data, get_master_data, get_users_quiz_attemptsfrom scheduled_loader import insert_to_sheet## QUERIES TO RUN ON DATABASESoris_reps = """SELECT salespersons.*, users.status AS user_status, users.email AS email, users.email_confirmed AS email_confirmed, users.last_name AS last_name, users.first_name AS first_name, users.phone_number, users.gender, users.date_of_birth AS DOB, users.home_addressFROM salespersonsLEFT JOIN usersON salespersons.user_id = users.id;"""oris_client = """SELECT *FROM sales_clients;"""oris_offers = """SELECT offers.*, jobs.job_type AS job_type, jobs.experience_level AS experience_level, jobs.department AS department, users.email AS email, users.last_name AS last_name, users.first_name AS first_name, users.phone_number AS phone_number, users.formatted_phone_number AS formatted_phone_number, bank_accounts.account_name, bank_accounts.account_number, bank_accounts.bank_code, states.name AS state, sales_clients.company_name AS client_nameFROM offersLEFT JOIN job_applicationsON offers.job_application_id = job_applications.idLEFT JOIN jobsON job_applications.job_id = jobs.idLEFT JOIN salespersonsON offers.salesperson_id = salespersons.idLEFT JOIN users ON salespersons.user_id = users.idLEFT JOIN bank_accounts ON offers.bank_account_id = bank_accounts.idLEFT JOIN states ON offers.state_id = states.idLEFT JOIN sales_clientsON offers.oauth_client_id = sales_clients.oauth_client_id;"""oris_employees = """SELECT employees.*, sales_clients.company_name as client_name, users.email AS email, users.last_name AS last_name, users.first_name AS first_name, users.phone_number AS phone_number, users.formatted_phone_number AS formatted_phone_numberFROM employeesLEFT JOIN sales_clients ON employees.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_employment_records = """SELECT employment_records.*, employees.employee_code, employees.user_id, employees.is_active, employees.bank_account_id, employees.guarantor_id, employees.custom_fields, employees.statutory_exempt, sales_clients.company_name AS client_name, users.email, users.email_confirmed, users.last_name, users.first_name, users.status AS user_status, users.phone_number, users.gender, users.date_of_birth AS DOB, users.home_addressFROM employment_recordsLEFT JOIN employees ON employment_records.employee_uuid = employees.employee_uuidLEFT JOIN sales_clients ON employment_records.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_employee_status_periods = """SELECT employee_status_periods.*, employees.employee_code, employees.user_id, sales_clients.company_name as client_name, users.email AS email, users.last_name AS last_name, users.first_name AS first_nameFROM employee_status_periodsLEFT JOIN employees ON employee_status_periods.employee_uuid = employees.employee_uuidLEFT JOIN sales_clients ON employee_status_periods.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_employee_position_changes = """SELECT employee_position_changes.* , employees.employee_code, employees.user_id, sales_clients.company_name as client_name, users.email AS email, users.last_name AS last_name, users.first_name AS first_nameFROM employee_position_changesLEFT JOIN employees ON employee_position_changes.employee_uuid = employees.employee_uuidLEFT JOIN sales_clients ON employee_position_changes.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_employee_separations = """SELECT employee_separations.*, employees.employee_code, employees.user_id, sales_clients.company_name as client_name, users.email AS email, users.last_name AS last_name, users.first_name AS first_nameFROM employee_separationsLEFT JOIN employees ON employee_separations.employee_uuid = employees.employee_uuidLEFT JOIN sales_clients ON employee_separations.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_employee_payrolls = """SELECT employee_payrolls.*, payrolls.title AS payroll_title, employees.employee_code, employees.user_id, sales_clients.company_name as client_name, users.email AS email, users.last_name AS last_name, users.first_name AS first_nameFROM employee_payrollsLEFT JOIN payrolls ON employee_payrolls.payroll_id = payrolls.idLEFT JOIN employees ON employee_payrolls.employee_uuid = employees.employee_uuidLEFT JOIN sales_clients ON employee_payrolls.oauth_client_id = sales_clients.oauth_client_idLEFT JOIN users ON employees.user_id = users.id;"""oris_issues = """SELECT issues.*, sales_clients.company_name AS client_name FROM issuesLEFT JOIN sales_clientsON issues.oauth_client_id = sales_clients.oauth_client_id;"""oris_uploaded_achievement = """SELECT * FROM oris_sales_auth_db.uploaded_performance_records_tracker"""oris_uploaded_target = """SELECT *FROM uploaded_performance_tracker;"""db_and_queries = [ { "db_name": "oris_sales_auth_db", "queries": { "reps":oris_reps, "clients":oris_client, "offers":oris_offers, "employees":oris_employees, "employment_records":oris_employment_records, "employee_position_changes":oris_employee_position_changes, "employee_status_periods":oris_employee_status_periods, "employee_separations":oris_employee_separations, "employee_payrolls":oris_employee_payrolls, # "issues":oris_issues, # "uploaded_achievement": oris_uploaded_achievement, # "uploaded_target": oris_uploaded_target } }]# ------------------------------------------------------------------# PREFECT SETTINGS# ------------------------------------------------------------------PROJECT_PATH = Path.cwd()DATA_PATH = PROJECT_PATH / "data"CONFIG_PATH = PROJECT_PATH / "config"MAX_RETRIES = 3TIMEOUT = Nonedef load_settings(): """ Loads all Prefect Secret Blocks once and returns a configuration dictionary. """ db_secret = Secret.load("prod-database-variables").get() ssh_secret = Secret.load("prod-ssh-variables").get() gs_secret = Secret.load("google-sheet-variables").get() return { "ssh_host": ssh_secret["ssh_host"], "ssh_port": ssh_secret["ssh_port"], "ssh_username": ssh_secret["ssh_username"], "ssh_key_path": str(CONFIG_PATH / "aws_private_key.txt"), "db_host": db_secret["db_host"], "db_port": db_secret["db_port"], "db_username": db_secret["db_username"], "db_password": db_secret["db_password"], # "sheet_key": gs_secret["operational_db_sheet_key"], "max_retries": MAX_RETRIES, "timeout": TIMEOUT, "data_path": DATA_PATH, }# @task(log_prints=True, retries=2)def load_bank_data(config): return pd.read_csv( Path(config["data_path"]) / "bank_codes_and_names.csv" )@task(log_prints=True, retries=2)def extract_data(db_and_queries, config): print("Task started") return 'here' # return connect_and_run_queries( # local_host="127.0.0.1", # ssh_host=config["ssh_host"], # ssh_port=config["ssh_port"], # ssh_username=config["ssh_username"], # ssh_key_path=config["ssh_key_path"], # db_host=config["db_host"], # db_port=config["db_port"], # db_username=config["db_username"], # db_password=config["db_password"], # max_retries=config["max_retries"], # timeout=config["timeout"], # db_and_queries=db_and_queries, # )@task(log_prints=True, retries=2)@taskdef transform_data(query_results, bank_codes): config = { "banks_codes": bank_codes }# def transform_data(query_results, config): # config = { # "banks_codes": load_bank_data(config) # } # sales_db = query_results["oris_sales_auth_db"] # sales_results = transform_all_data( # sales_db, # config=config, # use_processes=False, # ) # return sales_results@flow( log_prints=True, name="Parent Flow", flow_run_name="EXPORT DATA TO GOOGLE SHEET",)def export_db(): print("Starting ETL") config = load_settings() bank_codes = load_bank_data(config) query_results = extract_data(db_and_queries, config) transformed = transform_data( query_results, bank_codes, ) return transformed