from loguru import logger import os import sys from dotenv import load_dotenv import config import pandas as pd from djagitit import db from djagitit import tableau load_dotenv() def setup_logging(): logger.add( f"{os.getenv('LOGDIR')}/{config.LogParams.log_file}", rotation=config.LogParams.log_rotation, retention=config.LogParams.log_retention, compression=config.LogParams.log_compression, level=config.LogParams.log_level ) def validate_job(job): if job not in config.Params.jobs_dict: raise ValueError( f"Invalid job: {job}. Available jobs: {list(config.Params.jobs_dict.keys())}" ) def get_hyper_filepath(job, step_index): return f"{os.getenv('DATADIR')}/{job}_{step_index + 1}.hyper" def read_query(query_file): return open(os.path.join("../sql/", query_file)).read() def extract_step_data(rdb, step, job, step_index, total_steps): query = read_query(step['query']) filepath = get_hyper_filepath(job, step_index) logger.info(f"Run query: {step['query']}") df = rdb.query(query) logger.info("Query completed") if not isinstance(df, pd.DataFrame): raise RuntimeError(f"Error running query: {step['query']}") if df.empty: logger.warning(f"No data found for step {step_index + 1}/{total_steps}") return None logger.info(f"Write hyper file: {filepath}") tableau.write_hyper(df, filepath) logger.info("Hyper file written") return filepath def publish_step_datasource(step, job, step_index): datasource_id = step['datasource_id'] filepath = get_hyper_filepath(job, step_index) logger.info(f"Publish datasource: {datasource_id}") tableau.publish_datasource( filepath, site=config.Params.site, datasource_id=datasource_id ) logger.info(f"Datasource published: {datasource_id}") logger.info(f"Delete hyper file: {filepath}") os.remove(filepath) logger.info("Hyper file deleted") def process_job(job): validate_job(job) steps = config.Params.jobs_dict[job] total_steps = len(steps) rdb = db.ReportingDB() logger.info("Extracting data") for i, step in enumerate(steps): extract_step_data(rdb, step, job, i, total_steps) logger.success("Data extracted") logger.info("Publishing datasources") for i, step in enumerate(steps): filepath = get_hyper_filepath(job, i) if os.path.exists(filepath): publish_step_datasource(step, job, i) logger.success("Datasources published") def parse_args(): if len(sys.argv) > 2: raise ValueError(f"Too many arguments: {sys.argv[1:]}") if len(sys.argv) == 2: return [sys.argv[1]] return list(config.Params.jobs_dict.keys()) def main(): setup_logging() logger.info("Starting NORWAY-JOBS-CEA...") logger.info(f"Arguments: {sys.argv[1:]}") try: jobs = parse_args() for job in jobs: logger.info(f"Processing job: {job}") process_job(job) logger.info("NORWAY-JOBS-CEA completed.") except ValueError as e: logger.error(str(e)) sys.exit(1) except RuntimeError as e: logger.error(str(e)) sys.exit(1) if __name__ == "__main__": main()