Commit 6cf85f24 by Javier

Procesado de cdr moviles de fibra

parent a28997e9
...@@ -16,8 +16,14 @@ Le herramienta está desarrollada con Python. Se utilizan librerías de python p ...@@ -16,8 +16,14 @@ Le herramienta está desarrollada con Python. Se utilizan librerías de python p
### Importador MásMóvil ### Importador MásMóvil
Está compuesto de un script, `main.py`, se encarga de conectarse al sFTP de MásMóvil, y descargar los archivos de CDRs nuevos y guarda un registro en TB_CDR_Ficheros (tipo = 2) . Está compuesto de un script, `main.py`, se encarga de conectarse al sFTP de MásMóvil, y descargar los archivos de CDRs nuevos y guarda un registro en TB_CDR_Ficheros (tipo = 3) .
Posteriormente procederemos al procesado de dichos ficheros, guardando cada línea de estos archivos previamente descargados en las tablas `TB_Cdr_Fijos`. Posteriormente procederemos al procesado de dichos ficheros, guardando cada línea de estos archivos previamente descargados en las tablas `TB_Cdr_Fijos`.
En cada línea aparece la información de cada llamada, el teléfono que realiza la llamada, destino de llamada, duración, entre otros. En cada línea aparece la información de cada llamada, el teléfono que realiza la llamada, destino de llamada, duración, entre otros.
El script se encuentra alojado en la máquina `debian.utils` con la IP `192.168.2.63`. El script se encuentra alojado en la máquina `debian.utils` con la IP `192.168.2.63`.
### Descripción proceso
get_cdr_filename_download_fijos_fibra -> Obtiene fecha del último fichero subido procesado, en caso de que halla algún fichero sin procesar obtiene la fecha de creación dicho fichero.
extract_dwh_trafico_from_rar -> Para cada archivo .rar en `self.filename_fijo_download_sftp_path`, busca un archivo que contenga 'DWH_Trafico' en su nombre y lo descarga para procesarlo.
update_insert_date_fijo -> Inserto o actualizo en la tabla TB_CDR_ficheros los registros de los ficheros descargados de **fijos** , poniendo el campo procesado a null.
save_cdr_fijos_into_database -> Procesa los cdr fijos(tipo=3) que esten sin procesar(precessed_at = null) en la tabla TB_CDR_Ficheros y los almacena en TB_Cdr_Fibra
...@@ -8,10 +8,10 @@ load_dotenv() ...@@ -8,10 +8,10 @@ load_dotenv()
class ConectorDbWifi(): class ConectorDbWifi():
def __init__(self): def __init__(self):
sql_server_hostname = os.getenv("SQLSERVER_SERVER_TEST") sql_server_hostname = os.getenv("SQLSERVER_SERVER")
sql_server_database = os.getenv("SQLSERVER_DATABASE_TEST") sql_server_database = os.getenv("SQLSERVER_DATABASE")
sql_server_username = os.getenv("SQLSERVER_USERNAME_TEST") sql_server_username = os.getenv("SQLSERVER_USERNAME")
sql_server_password = os.getenv("SQLSERVER_PASSWORD_TEST") sql_server_password = os.getenv("SQLSERVER_PASSWORD")
self.cursor = None self.cursor = None
self.connection = None self.connection = None
......
from app.conector_sftp import SFTPConexion from app.conector_sftp import SFTPConexion
import os, time, csv, zipfile, shutil, patoolib import os, time
from datetime import datetime from datetime import datetime
from paramiko.ssh_exception import SSHException from paramiko.ssh_exception import SSHException
from log import logging from log import logging
...@@ -45,26 +45,43 @@ class ManagerCdrFijosFibra: ...@@ -45,26 +45,43 @@ class ManagerCdrFijosFibra:
with ConectorDbWifi() as cursor: with ConectorDbWifi() as cursor:
#Consultamos en la tabla la fecha del último fichero subido #Consultamos en la tabla la fecha del último fichero subido pero si hay algún fichero sin procesar lo obtiene.
cursor.execute("SELECT TOP 1 created_at FROM bdd_wifi.dbo.TB_CDR_Ficheros WHERE tipo = 3 ORDER BY created_at DESC") cursor.execute("""
row = cursor.fetchone() IF EXISTS (
last_cdr_fijos_fibra_table_ficheros = row[0] SELECT 1
FROM bdd_wifi.dbo.TB_CDR_Ficheros
# REcorro el directorio principal del sftp donde se almacena los da WHERE tipo = 3 AND processed_at IS NULL
)
BEGIN
SELECT TOP 1 fecha_origen
FROM bdd_wifi.dbo.TB_CDR_Ficheros
WHERE tipo = 3 AND processed_at IS NULL
ORDER BY fecha_origen ASC
END
ELSE
BEGIN
SELECT TOP 1 fecha_origen
FROM bdd_wifi.dbo.TB_CDR_Ficheros
WHERE tipo = 3
ORDER BY fecha_origen DESC
END
""")
row = cursor.fetchone()
last_cdr_fijos_fibra_table_ficheros = row[0] if row else None
# Recorro el directorio principal del sftp donde se almacena los da
item_list = self.cnx_sftp_fibra.listdir_attr('/') item_list = self.cnx_sftp_fibra.listdir_attr('/')
for item in item_list: for item in item_list:
fecha = time.localtime(item.st_mtime) # Fecha creación fichero en el SFTP fecha = time.localtime(item.st_mtime) # Fecha creación fichero en el SFTP
file_date_time = datetime(fecha[0], fecha[1], fecha[2], fecha[3], fecha[4], fecha[5]) file_date_time = datetime(fecha[0], fecha[1], fecha[2], fecha[3], fecha[4], fecha[5])
# Obtengo el nombre de los ficheros de fijos que la fecha sea posterior al último cdr ingresado en la tabla, y lo guardo en una lista, # Obtengo el nombre de los ficheros de fijos que la fecha sea posterior al último cdr ingresado en la tabla, y lo guardo en una lista,
# también guardo la ruta absoluta # también guardo la ruta absoluta
# y un diccionario con la fecha de creación del fichero en el sftp asociada al nombre del fichero # y un diccionario con la fecha de creación del fichero en el sftp asociada al nombre del fichero
if last_cdr_fijos_fibra_table_ficheros < file_date_time and 'mvno' not in item.filename and '0000' in item.filename: item_filename = item.filename
if last_cdr_fijos_fibra_table_ficheros < file_date_time and 'mvno' not in item_filename and '0000' in item_filename and '.rar' in item_filename:
item_filename = item.filename
self.filename_datetime_fijo[item_filename]=str(file_date_time) self.filename_datetime_fijo[item_filename]=str(file_date_time)
self.filename_fijo_download_sftp_path.append('/' + item_filename) self.filename_fijo_download_sftp_path.append('/' + item_filename)
...@@ -99,6 +116,7 @@ class ManagerCdrFijosFibra: ...@@ -99,6 +116,7 @@ class ManagerCdrFijosFibra:
# Usamos la librería rarfile para descomprimir el archivo # Usamos la librería rarfile para descomprimir el archivo
try: try:
with rarfile.RarFile(local_rar_path) as rar: with rarfile.RarFile(local_rar_path) as rar:
file_list = rar.namelist() file_list = rar.namelist()
logging.info(f"Archivos dentro del RAR: {file_list}") logging.info(f"Archivos dentro del RAR: {file_list}")
...@@ -147,7 +165,7 @@ class ManagerCdrFijosFibra: ...@@ -147,7 +165,7 @@ class ManagerCdrFijosFibra:
def update_insert_date_fijo(self): def update_insert_date_fijo(self):
""" """
Inserto o actualizo en la tabla TB_CDR_ficheros los ficheros de **fijos** , poniendo el campo procesado a null. Inserto o actualizo en la tabla TB_CDR_ficheros los registros de los ficheros descargados de **fijos** , poniendo el campo procesado a null.
""" """
# Recorro el nombre de los ficheros descargados y consulto en la tabla si ese fichero existe lo actualizo si no existe lo creo # Recorro el nombre de los ficheros descargados y consulto en la tabla si ese fichero existe lo actualizo si no existe lo creo
for cdr_filename in self.new_names_files: for cdr_filename in self.new_names_files:
...@@ -169,7 +187,7 @@ class ManagerCdrFijosFibra: ...@@ -169,7 +187,7 @@ class ManagerCdrFijosFibra:
name_dir_contenedor = self.nombre_carpetas_contenedore_txts.get(cdr_filename) name_dir_contenedor = self.nombre_carpetas_contenedore_txts.get(cdr_filename)
fecha_creacion_fichero = self.filename_datetime_fijo.get(name_dir_contenedor) fecha_creacion_fichero = self.filename_datetime_fijo.get(name_dir_contenedor)
fecha_creacion_fichero_formateada = datetime.strptime(fecha_creacion_fichero, '%Y-%m-%d %H:%M:%S') fecha_creacion_fichero_formateada = datetime.strptime(fecha_creacion_fichero, '%Y-%m-%d %H:%M:%S')
cursor.execute("UPDATE TB_CDR_Ficheros SET tipo = ?, fecha_origen = ?,Num_Reg_Txt = ?, Degradado=?, Mensual=? WHERE filename = ?", 3, fecha_creacion_fichero_formateada,lineas_fichero_txt, degradado, mensual, str(cdr_filename)) cursor.execute("UPDATE TB_CDR_Ficheros SET tipo = ?, fecha_origen = ?,Num_Reg_Txt = ? WHERE filename = ?", 3, fecha_creacion_fichero_formateada,lineas_fichero_txt, str(cdr_filename))
cursor.commit() cursor.commit()
message = f'Fichero {cdr_filename} actualizado en la tabla TB_CDR_Ficheros' message = f'Fichero {cdr_filename} actualizado en la tabla TB_CDR_Ficheros'
logging.info(message) logging.info(message)
...@@ -196,10 +214,11 @@ class ManagerCdrFijosFibra: ...@@ -196,10 +214,11 @@ class ManagerCdrFijosFibra:
logging.info(message) logging.info(message)
""" def save_cdr_fijos_into_database(self): def save_cdr_fijos_into_database(self):
"" """
Procesa los cdr fijos(tipo=3) que esten sin procesar(precessed_at = null) en la tabla TB_CDR_Ficheros Procesa los cdr fijos(tipo=3) que esten sin procesar(precessed_at = null) en la tabla TB_CDR_Ficheros y los almacena
"" en TB_Cdr_Fibra
"""
with ConectorDbWifi() as cursor: with ConectorDbWifi() as cursor:
try: try:
...@@ -218,12 +237,203 @@ class ManagerCdrFijosFibra: ...@@ -218,12 +237,203 @@ class ManagerCdrFijosFibra:
if len(filenames) != 0: if len(filenames) != 0:
# Obtengo una lista con las rutas de los ficheros a procesar # Obtengo una lista con las rutas de los ficheros a procesar
path_filenames = cdr_list(filenames) path_filenames = cdr_list(filenames)
#Crea el proceso de cdr de fijos introduciendolo en la tabla TB_Cdr_Fibra #Crea el proceso de cdr de fijos introduciendolo en la tabla TB_Cdr_Fibra
cdr_insert_update_cdr_fijos(path_filenames, cursor) cdr_insert_update_cdr_fijos(path_filenames, cursor)
except Exception as e: except Exception as e:
logging.error(f"An error occurred: {e}") """ logging.error(f"An error occurred: {e}")
def cdr_insert_update_cdr_fijos(path_filenames: list, cursor_wifi):
""" Inserta o actualiza cdr de fijos de fibra
Params:
path_filenames (list): Nombre de la ruta de fichero que hay que procesar
cursor_wifi (Cursor): Cursor DB_wifi
"""
try:
#print(f'Ficheros a procesar: {path_filenames}')
for path_file in path_filenames:
if os.path.exists(path_file):
with open(path_file, 'r', encoding='utf-8') as f:
lineas = f.readlines()
if not lineas:
print(f'⚠️ El fichero está vacío: {path_file}')
continue
nombre_fichero = os.path.basename(path_file)
#En caso de que se halla producido un error en el procesado de un fichero tendremos el num de registros procesados por el que se ha quedado guardado en la bd TB_CDR_ficheros
#con este proceso empezariamos desde el registro donde se quedo
cursor_wifi.execute("SELECT Num_Reg_Pro, Num_Reg_Dup, Num_Reg_Txt FROM TB_CDR_Ficheros WHERE filename = ?",
nombre_fichero,
)
Num_Reg_Pro_d =cursor_wifi.fetchone()
Num_Reg_Pro_db = int(Num_Reg_Pro_d[0]) + int(Num_Reg_Pro_d[1]) #Total de registros procesados de un fichero por si hay una excepción partir el procesado desde ese registro
Num_Reg_Txt = Num_Reg_Pro_d[2]
time_init = datetime.now()
lineas_save = 0
lines_traveled = 0 # Líneas recorridas
for linea in lineas:
lines_traveled += 1
line_parse = parse_linea_fibra(linea)
if line_parse is not None:
try:
#TODO Si existe lo actualiza si no lo inserta
# Definir los campos clave para identificar si el registro ya existe
campos_clave = (
line_parse["num_origen"],
line_parse["num_destino"],
line_parse["fecha_inicio"],
line_parse["fecha_fin"]
)
# Verificar si ya existe
cursor_wifi.execute("""
SELECT COUNT(*) FROM [bdd_wifi].[dbo].[TB_Cdr_Fibra]
WHERE num_origen = ? AND num_destino = ? AND fecha_inicio = ? AND fecha_fin = ?
""", campos_clave)
existe = cursor_wifi.fetchone()[0]
if existe:
# Si ya existe, actualizamos el registro
cursor_wifi.execute("""
UPDATE [bdd_wifi].[dbo].[TB_Cdr_Fibra]
SET campo1 = ?, campo2 = ?, campo5 = ?, destino = ?, duracion = ?,
campo8 = ?, campo9 = ?, impuestos = ?, campo11 = ?, campo12 = ?,
campo13 = ?, campo14 = ?, fecha_fin = ?, otra_fecha = ?, campo18 = ?
WHERE num_origen = ? AND num_destino = ? AND fecha_inicio = ? AND fecha_fin = ?
""", (
line_parse["campo1"], line_parse["campo2"], line_parse["campo5"], line_parse["destino"],
line_parse["duracion"], line_parse["campo8"], line_parse["campo9"], line_parse["impuestos"],
line_parse["campo11"], line_parse["campo12"], line_parse["campo13"], line_parse["campo14"],
line_parse["fecha_fin"], line_parse["otra_fecha"], line_parse["campo18"],
line_parse["num_origen"], line_parse["num_destino"], line_parse["fecha_inicio"], line_parse["fecha_fin"]
))
else:
# Query SQL para insertar en la base de datos
cursor_wifi.execute("""
INSERT INTO [bdd_wifi].[dbo].[TB_Cdr_Fibra] (campo1, campo2, num_origen, num_destino, campo5, destino, duracion,
campo8, campo9, impuestos, campo11, campo12, campo13, campo14,
fecha_inicio, fecha_fin, otra_fecha, campo18)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",(line_parse["campo1"], line_parse["campo2"], line_parse["num_origen"], line_parse["num_destino"],
line_parse["campo5"], line_parse["destino"], line_parse["duracion"], line_parse["campo8"],
line_parse["campo9"], line_parse["impuestos"], line_parse["campo11"], line_parse["campo12"],
line_parse["campo13"], line_parse["campo14"], line_parse["fecha_inicio"], line_parse["fecha_fin"],
line_parse["otra_fecha"], line_parse["campo18"]))
# Confirmar los cambios en la base de datos
cursor_wifi.commit()
lineas_save += 1
except Exception as e:
print(f"❌ Error insertando en DB: {e}")
time.sleep(1)
time_final = datetime.now()
time_procesado = time_final - time_init
time_procesado_fin = str(time_procesado)
message = f'Total registros procesados con costes: {lineas_save}.'
cursor_wifi.execute(
"UPDATE TB_CDR_Ficheros SET processed_at=?, Num_Reg_Pro=?, Time_Procesado=?,Observaciones=? WHERE filename=?",
datetime.now(),
lineas_save,
time_procesado_fin,
message,
nombre_fichero,
)
cursor_wifi.commit()
# Si las líneas recorridas son iguales a las líneas total del fichero se borra el fichero de local
print(f'path_file: {path_file} lines travelend: {lines_traveled} Num_Reg_Txt:{Num_Reg_Txt}')
if lines_traveled == Num_Reg_Txt:
if os.path.isfile(path_file):
os.remove(path_file)
else:
print(f'⚠️ El fichero no existe: {path_file}')
except Exception as e:
message = f'Total registros procesados con costes: {lineas_save}.'
cursor_wifi.execute(
"UPDATE TB_CDR_Ficheros SET processed_at=?, Num_Reg_Pro=?,Observaciones=? WHERE filename=?",
datetime.now(),
lineas_save,
message,
nombre_fichero,
)
cursor_wifi.commit()
def parse_linea_fibra(linea: str) -> dict:
"""
Parsea una línea del archivo de fibra en un diccionario con los campos correspondientes
siempre que el coste sea mayor a 0€.
Params:
linea (str): Línea del archivo a procesar.
Returns:
dict | None: Diccionario con los datos parseados o None si la línea está vacía.
"""
try:
if not linea or linea.strip() == "":
return None
campos = linea.strip().split("|")
if len(campos) > 9:
valor8 = float(campos[8].replace(",", ".")) if campos[8] else 0
valor9 = float(campos[9].replace(",", ".")) if campos[9] else 0
if valor8 > 0 or valor9 > 0:
output = {
"campo1": campos[0] if len(campos) > 0 and campos[0] else None,
"campo2": campos[1] if len(campos) > 1 and campos[1] else None,
"num_origen": campos[2] if len(campos) > 2 and campos[2] else None,
"num_destino": campos[3] if len(campos) > 3 and campos[3] else None,
"campo5": campos[5] if len(campos) > 5 and campos[5] else None,
"destino": campos[6] if len(campos) > 6 and campos[6] else None,
"duracion": campos[7] if len(campos) > 7 and campos[7] else None,
"campo8": campos[8] if len(campos) > 8 and campos[8] else None,
"campo9": campos[9] if len(campos) > 9 and campos[9] else None,
"impuestos": campos[10] if len(campos) > 10 and campos[10] else None,
"campo11": campos[11] if len(campos) > 11 and campos[11] else None,
"campo12": campos[12] if len(campos) > 12 and campos[12] else None,
"campo13": campos[13] if len(campos) > 13 and campos[13] else None,
"campo14": campos[14] if len(campos) > 14 and campos[14] else None,
"fecha_inicio": campos[15] if len(campos) > 15 and campos[15] else None,
"fecha_fin": campos[16] if len(campos) > 16 and campos[16] else None,
"otra_fecha": campos[17] if len(campos) > 17 and campos[17] else None,
"campo18": campos[18] if len(campos) > 18 and campos[18] else None,
}
else:
return None
# Convertir fechas al formato "d-m-Y H:i:s" si existen
def convertir_fecha(fecha):
if fecha:
try:
return datetime.strptime(fecha, "%d/%m/%Y %H:%M:%S").strftime("%d-%m-%Y %H:%M:%S")
except ValueError:
return fecha # Si hay un error en el formato, se deja tal cual
return None
output["fecha_inicio"] = convertir_fecha(output["fecha_inicio"])
output["fecha_fin"] = convertir_fecha(output["fecha_fin"])
output["otra_fecha"] = convertir_fecha(output["otra_fecha"])
return output
except Exception as e:
print(f'Excepcion en parse_linea_fibra: {e}')
def cdr_list(filenames:list) -> list: def cdr_list(filenames:list) -> list:
......
from app.manager_cdr_fijos import ManagerCdrFijos """ from app.manager_cdr_fijos import ManagerCdrFijos """
from app.manager_cdr_fijos_fibra import ManagerCdrFijosFibra from app.manager_cdr_fijos_fibra import ManagerCdrFijosFibra
if __name__ == '__main__': if __name__ == '__main__':
...@@ -20,4 +20,5 @@ if __name__ == '__main__': ...@@ -20,4 +20,5 @@ if __name__ == '__main__':
manager = ManagerCdrFijosFibra() manager = ManagerCdrFijosFibra()
manager.get_cdr_filename_download_fijos_fibra() manager.get_cdr_filename_download_fijos_fibra()
manager.extract_dwh_trafico_from_rar() manager.extract_dwh_trafico_from_rar()
manager.update_insert_date_fijo() manager.update_insert_date_fijo()
\ No newline at end of file manager.save_cdr_fijos_into_database()
\ No newline at end of file
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment