Loop to ensure data loading

Viewed 23

I've created a process that basically reads a csv, transforms it into a dataframe, and then loads data into a database. This process is triggered through nifi. The problem I have, is that in case there is a problem in the original csv file(this file is generally uploaded between 7:45 -8:00 am, and Nifi is triggered at 8:10 am), which is a daily file, which changes every day, and nifi is executed, but does not load the data, I will have lost the data of that day . I need suggestions on how I should modify the code of the process, in order to make sure that it loads, making sure that nifi is going to load, through a cycle that calls a backup file. Next, I share the code of the process:

import com.tchile.bigdata.hdfs.Hdfs
import com.tchile.bigdata.utilities.Utilities
import org.apache.spark.SparkConf
import org.apache.spark.sql.{SaveMode, SparkSession}
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

import java.sql.DriverManager
import java.util.concurrent.TimeUnit

class Process {
  def process(shell_variable: String): Unit = {

    val startTotal = TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis())
    println("[INFO] Inicio de proceso Logistica Carga Input Simple Data")

    // Productivo (Configuracion de Spark)
    val sparkConf = new SparkConf()
      .setAppName("Logistica_Carga_Input_SimpleData")
      .set("spark.sql.sources.partitionOverwriteMode", "dynamic")
      .set("spark.sql.parquet.output.committer.class", "org.apache.parquet.hadoop.ParquetOutputCommitter")
      .set("spark.sql.sources.commitProtocolClass", "org.apache.spark.sql.execution.datasources.SQLHadoopMapReduceCommitProtocol")
    val spark = SparkSession
      .builder()
      .config(sparkConf)
      .enableHiveSupport()
      .getOrCreate()
    val sc = spark.sparkContext

    import spark.implicits._

    val hdfs = new Hdfs
    val utils = new Utilities

    println("[INFO] Obteniendo variables desde shell/NIFI")

    // Valores extraídos desde NIFI
    val daysAgo = shell_variable.split(":")(0).toInt // Cantidad de días de reproceso: inicio
    val daysAhead = shell_variable.split(":")(1).toInt // Cantidad de días de reproceso: límite
    val repartition = shell_variable.split(":")(2).toInt

    // Variables auxiliares para reproceso en HDFS
    var pathToProcess = ""
    var pathToDelete = ""
    var deleteStatus = false

    println("[INFO] Obteniendo variables desde Parametros.conf")

    // Valores extraídos de Parametros.conf > exadata
    val driver_jdbc = sc.getConf.get("spark.exadata.driver_jdbc")
    val url_jdbc = sc.getConf.get("spark.exadata.url_jdbc")
    val user_jdbc = sc.getConf.get("spark.exadata.user_jdbc")
    val pass_jdbc = sc.getConf.get("spark.exadata.pass_jdbc")
    val table_name = sc.getConf.get("spark.exadata.table_name")
    val table_owner = sc.getConf.get("spark.exadata.table_owner")
    val table = sc.getConf.get("spark.exadata.table")

    // Valores extraídos de Parametros.conf > path

    //val pathCsv = sc.getConf.get("spark.path.Csv")

    // Cálculo de la fecha de los días atrás
    val startDate = utils.getCalculatedDate("yyyy-MM-dd", -daysAgo)

    // Cálculo de la fecha límite
    val endDate = utils.getCalculatedDate("yyyy-MM-dd", -daysAhead)

    // Se separan el año, el mes y el día a partir de la fecha_atras
    val startDateYear = startDate.substring(0, 4)
    val startDateMonth = startDate.substring(5, 7)
    val startDateDay = startDate.substring(8, 10)

    // Se separan el año, el mes y el día a partir del sub_day
    val endDateYear = endDate.substring(0, 4)
    val endDateMonth = endDate.substring(5, 7)
    val endDateDay = endDate.substring(8, 10)


    // Información para log
    println("[INFO] Reproceso de: " + daysAgo + " días")
    println("[INFO] Fecha inicio: " + startDate)
    println("[INFO] Fecha límite: " + endDate)

    try {

      // ================= INICIO LÓGICA DE PROCESO =================
      //Se crea df a partir de archivo diario input de acuerdo a la ruta indicada
      val df_csv = spark.read.format("csv").option("header","true").option("sep",";").option("mode","dropmalformed").load("/applications/recup_remozo_equipos/equipos_por_recuperar/output/agendamientos_sin_pet_2")

      val df_final = df_csv.select($"RutSinDV".as("RUT_SIN_DV"),
        $"dv".as("DV"),
        $"Agendado".as("AGENDADO"),
        to_date(col("Dia_Agendado"), "yyyy-MM-dd").as("DIA_AGENDADO"),
        $"Horario_Agendado".as("HORARIO_AGENDADO"),
        $"Nombre_Agendamiento".as("NOMBRE_AGENDAMIENTO"),
        $"Telefono_Agendamiento".as("TELEFONO_AGENDAMIENTO"),
        $"Email".as("EMAIL"),
        $"Region_Agendamiento".as("REGION_AGENDAMIENTO"),
        $"Comuna_Agendamiento".as("COMUNA_AGENDAMIENTO"),
        $"Direccion_Agendamiento".as("DIRECCION_AGENDAMIENTO"),
        $"Numero_Agendamiento".as("NUMERO_AGENDAMIENTO"),
        $"Depto_Agendamiento".as("DEPTO_AGENDAMIENTO"),
        to_timestamp(col("fecha_registro")).as("FECHA_REGISTRO"),
        to_timestamp(col("Fecha_Proceso")).as("FECHA_PROCESO")
      )
      // ================== FIN LÓGICA DE PROCESO ==================

      // Limpieza en EXADATA
      println("[INFO] Se inicia la limpieza por reproceso en EXADATA")

      val query_particiones = "(SELECT * FROM (WITH DATA AS (select table_name,partition_name,to_date(trim('''' " +
        "from regexp_substr(extractvalue(dbms_xmlgen.getxmltype('select high_value from all_tab_partitions " +
        "where table_name='''|| table_name|| ''' and table_owner = '''|| table_owner|| ''' and partition_name = '''" +
        "|| partition_name|| ''''),'//text()'),'''.*?''')),'syyyy-mm-dd hh24:mi:ss') high_value_in_date_format " +
        "FROM all_tab_partitions WHERE table_name = '" + table_name + "' AND table_owner = '" + table_owner + "')" +
        "SELECT partition_name FROM DATA WHERE high_value_in_date_format > DATE '" + startDateYear + "-" + startDateMonth + "-" + startDateDay + "' " +
        "AND high_value_in_date_format <= DATE '" + endDateYear + "-" + endDateMonth + "-" + endDateDay + "') A)"

      Class.forName(driver_jdbc)
      val db = DriverManager.getConnection(url_jdbc, user_jdbc, pass_jdbc)
      val st = db.createStatement()

      try {
        val consultaParticiones = spark.read.format("jdbc")
          .option("url", url_jdbc)
          .option("driver", driver_jdbc)
          .option("dbTable", query_particiones)
          .option("user", user_jdbc)
          .option("password", pass_jdbc)
          .load()
          .collect()
        for (partition <- consultaParticiones) {
          st.executeUpdate("call " + table_owner + ".DO_THE_TRUNCATE_PARTITION('" + table + "','" + partition.getString(0) + "')")
        }
      } catch {
        case e: Exception =>
          println("[ERROR TRUNCATE] " + e)
      }

      st.close()
      db.close()

      println("[INFO] Se inicia la inserción en EXADATA")

      df_final.filter($"DIA_AGENDADO" >= "2022-08-01")
        .repartition(repartition).write.mode("append")
        .jdbc(url_jdbc, table, utils.jdbcProperties(driver_jdbc, user_jdbc, pass_jdbc))

      println("[INFO] Inserción en EXADATA completada con éxito")
      println("[INFO] Proceso Logistica Carga Input SimpleData")

      val endTotal = TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) - startTotal
      println("[INFO] TIEMPO TOTAL EJECUCIÓN: " + utils.secondsToMinutes(endTotal))
      }
    catch {
      case e: Exception =>
        println("[EXCEPTION] " + e)
        val endTotal = TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) - startTotal
        println("[INFO] TIEMPO TOTAL EJECUCIÓN (CON ERROR): " + utils.secondsToMinutes(endTotal))
        throw e
    }
  }
}

I appreciate suggestions on how I should modify or add to the code, to ensure the daily data upload.

0 Answers
Related