Skip to content

luxorDB_dataloader.py

Camada de escrita/persistência padronizada no LuxorDB.

Este módulo contém a classe DataLoader, responsável por gerenciar o carregamento e a persistência de dados para o LuxorDB.

A classe DataLoader oferece funcionalidades para: - Carregar tabelas a partir de arquivos Excel ou dados em memória. - Monitorar arquivos e tabelas para detectar modificações e acionar recargas. - Normalizar colunas de texto para minúsculas. - Persistir dados em diferentes formatos (Excel, Parquet) no sistema de arquivos local. - Exportar dados para o Azure Blob Storage em formato Parquet.

O objetivo principal é garantir que os dados no LuxorDB estejam sempre atualizados e consistentes, facilitando a integração com outras partes do sistema.

De forma bem resumida: Este módulo é responsável por carregar, monitorar e persistir dados no LuxorDB, garantindo que as tabelas estejam sempre atualizadas e consistentes. Ele gerencia a leitura de arquivos Excel e dados em memória, normaliza colunas e exporta os dados para formatos como Excel, Parquet e Azure Blob Storage. Oferece funcoes para:

DataLoader

Source code in LuxorASAP/luxorDB_dataloader.py
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
class DataLoader:

    def __init__(self, luxorDB_directory = None):
        """Fornece uma forma padronizada de carregar tabelas para a luxorDB.
            1. Possui metodos para carregar tabelas que já estao carregadas na memória
                - Sao os metodos que possuem 'table' no nome
            2. Possui metodos para carregar arquivos de excel, com todas as suas abas
                Inclui metodo para checagem de alteracao de versao do arquivo
                - Sao os metodos que possuem 'file' no nome
        Args:
            luxorDB_directory (pathlib.Path, optional): Caminho completo ate o diretorio de destino dos dados.
        """
        self.luxorDB_directory = luxorDB_directory

        if self.luxorDB_directory is None:
            self.luxorDB_directory = Path(__file__).absolute().parent/"LuxorDB"/"tables"

        self.tracked_files = {}
        self.tracked_tables = {}

    def __persist_column_formatting(self, t):

        columns_to_persist = {"Name", "Class", "Vehicles", "Segment"}

        if len(set(t.columns).intersection(columns_to_persist)) > 0:
            # Vamos persistir a formatacao de algumas colunas
            columns_order = list(t.columns)
            columns_to_persist = list(set(t.columns).intersection(columns_to_persist))
            persistent_data = t[columns_to_persist].copy()

            columns_to_normalize = list(set(columns_order) - set(columns_to_persist))
            t = self.text_to_lowercase(t[columns_to_normalize])
            t.loc[:,columns_to_persist] = persistent_data
            return t[columns_order]

        # Nos outros casos, transformaremos tudo em lowercase
        return self.text_to_lowercase(t)


    def text_to_lowercase(self, t):
        """
        Converte todas as colunas de texto para lowercase
        Args:
            t (pd.DataFrame): pandas DataFrame
        Returns:
            pd.DataFrame
        """

        return t.map(lambda x: x.lower().strip() if isinstance(x, str) else x)


    def add_file_tracker(self, tracked_file_path, filetype="excel", sheet_names={}, 
            excel_size_limit = None,index=False, index_name="index",normalize_columns=False):
        """ Adiciona arquivo na lista para checar por alteracao
        Args:
            tracked_file_path (pathlib.Path): caminho completo ate o arquivo,
                    incluindo nome do arquivo e extensão.
            sheet_names (dict, optional): Caso seja uma planilha com varias abas, mapear
                    aqui o nome da aba para o nome do arquivo de saida desejado.
        """
        if tracked_file_path not in self.tracked_files:
            self.tracked_files[tracked_file_path] = {
                    "last_mtime" : dt.datetime.timestamp(dt.datetime(2000,1,1)),
                    "filetype" : filetype, "sheet_names": sheet_names,
                    "excel_size_limit" : excel_size_limit,
                    "index" : index,
                    "index_name" : index_name,
                    "normalize_columns" : normalize_columns,
                }


    def add_table_tracker(self, table_name:str):
        """ Adiciona tabela na lista para controle de alteracao."""

        if table_name not in self.tracked_tables:
            self.tracked_tables[table_name] = dt.datetime.timestamp(dt.datetime(2000,1,1))


    def remove_file_tracker(self, tracked_file_path):

        if tracked_file_path in self.tracked_files:
            del self.tracked_files[tracked_file_path]


    def remove_table_tracker(self, table_name:str):

        if table_name in self.tracked_tables:
            del self.tracked_tables[table_name]


    def is_file_modified(self, tracked_file_path: Path) -> {bool, float}:
        """ Checa se o arquivo foi modificado desde a ultima leitura.
        Returns:
            tuple(bool, float): (foi modificado?, timestamp da ultima modificacao)
        """

        file_data = self.tracked_files[tracked_file_path]

        last_saved_time = file_data["last_mtime"]
        file_last_update = tracked_file_path.stat().st_mtime
        return file_last_update > last_saved_time, file_last_update


    def set_file_modified_time(self, tracked_file_path, file_mtime):

        self.tracked_files[tracked_file_path]["last_mtime"] = file_mtime


    def load_file_if_modified(self, tracked_file_path, export_to_blob=False, blob_directory='enriched/parquet'):
        """Carrega arquivo no caminho indicado, carregando na base de dados caso modificado.
        Args:
            tracked_file_path (pathlib.Path): caminho ate o arquivo(cadastro previamente por add_file_tracker)
            type_map (_type_, optional): _description_. Defaults to None.
            filetype (str, optional): _description_. Defaults to "excel".
        """
        file_data = self.tracked_files[tracked_file_path]

        last_saved_time = file_data["last_mtime"]
        filetype = file_data["filetype"]
        file_sheets = file_data["sheet_names"]

        file_last_update = tracked_file_path.stat().st_mtime

        if file_last_update > last_saved_time: # Houve alteracao no arquivo
            if filetype == "excel":
                file_sheets = None if len(file_sheets) == 0 else list(file_sheets.keys())

                # tables sera sempre um dicionario de tabelas
                tables = None
                trials = 25
                t_counter = 1
                while trials - t_counter > 0:
                    try:
                        tables = pd.read_excel(tracked_file_path, sheet_name=file_sheets)
                        t_counter = trials # leitura concluida
                    except PermissionError:

                        logger.error(f"Erro ao tentar ler arquivo '{tracked_file_path}.\nTentativa {t_counter} de {trials};'.\nSe estiver aberto feche.")
                        time.sleep(10)
                        t_counter += 1

                for sheet_name, table_data in tables.items():

                    table_name = sheet_name if file_sheets is None else file_data["sheet_names"][sheet_name]

                    if table_name == "trades":
                        table_data["ID"] = table_data.index

                    self.__export_table(table_name, table_data, index=file_data["index"], index_name=file_data["index_name"],
                                            normalize_columns=file_data["normalize_columns"], export_to_blob=export_to_blob,
                                            blob_directory=blob_directory)
                self.tracked_files[tracked_file_path]["last_mtime"] = file_last_update


    def load_table_if_modified(self, table_name, table_data, last_update, index=False, index_name="index", normalize_columns=False,
                               do_not_load_excel=False, export_to_blob=False, blob_directory='enriched/parquet',
                               is_data_in_bytes=False, bytes_extension=".xlsx"):
        """
        Args:
            table_name (str): nome da tabela (sera o mesmo do arquivo a ser salvo)
            table_data (pd.DataFrame): tabela de dados
            last_update (timestamp): timestamp da ultima edicao feita na tabela
        """

        if table_name not in self.tracked_tables:
            self.add_table_tracker(table_name)


        last_update_time = self.tracked_tables[table_name]
        if last_update > last_update_time:

            self.tracked_tables[table_name] = last_update
            self.__export_table(table_name, table_data, index=index, index_name=index_name, normalize_columns=normalize_columns,
                                do_not_load_excel=do_not_load_excel, export_to_blob=export_to_blob, blob_directory=blob_directory,
                                is_data_in_bytes=is_data_in_bytes, bytes_extension=bytes_extension)


    def scan_files(self, export_to_blob=False, blob_directory='enriched/parquet'):
        """
            Para todos os arquivos cadastrados, vai buscar e carregar quando houver
            arquivo mais recente.
        """

        for file in self.tracked_files:

            self.load_file_if_modified(file, export_to_blob=export_to_blob, blob_directory=blob_directory)


    #def __load_bytes(self, content: bytes, extension=".xlsx") -> pd.DataFrame:
    #    if extension == ".xlsx" or extension == "xlsx" or extension == "xls":
    #        df = pd.read_excel(io.BytesIO(content), engine="openpyxl")
    #    
    #        return df
#
    #    raise ValueError(f'Extension {extension} not supported')

    def __load_bytes(self, content: bytes, extension: str) -> pd.DataFrame:
        extension = extension.lower()

        if extension in [".xlsx", ".xls", "xlsx", "xls"]:
            df = pd.read_excel(io.BytesIO(content), engine="openpyxl")
            return df

        if extension == ".csv":
            try:
                return pd.read_csv(io.BytesIO(content), encoding="utf-8")
            except UnicodeDecodeError:
                return pd.read_csv(io.BytesIO(content), encoding="latin1")

        if extension == ".parquet":
            df = pd.read_parquet(io.BytesIO(content))
            return df

        raise ValueError(f'Extension {extension} not supported')


    def __export_table(self, table_name, table_data, index=False, index_name="index", normalize_columns=False,
                       do_not_load_excel=False, export_to_blob=False, blob_directory='enriched/parquet',
                       is_data_in_bytes=False, bytes_extension=".xlsx"):

        dest_directory = self.luxorDB_directory
        #TODO -> formatar para index=False
        # Salvando em formato excel
        attempts = 10
        count_attempt = 0

        if is_data_in_bytes:
            table_data = self.__load_bytes(table_data, extension=bytes_extension)

        # Se o index tiver dados, vamos trata-los para virar uma coluna
        if index:
            # Tratando o nome do index, caso seja necessario transformar em coluna
            prev_index = table_data.index.name
            if prev_index is not None and index_name == "index":
                index_name = prev_index
            table_data.index.name = index_name

            table_data = table_data.reset_index()


        if normalize_columns:
            table_data = self.__persist_column_formatting(table_data)

        if not do_not_load_excel:
            while count_attempt < attempts:
                count_attempt += 1
                try:
                    if len(table_data) > 1_000_000:
                        table_data = table_data.tail(1_000_000)
                    table_data.to_excel(dest_directory/f"{table_name}.xlsx", index=False)
                    count_attempt = attempts # sair do loop

                except PermissionError:
                    logger.error(f"Erro ao tentar salvar arquivo {table_name}. Feche o arquivo. Tentativa {count_attempt} de {attempts}")
                    time.sleep(10 + count_attempt * 5)

        # Salvando em csv 
        # -> Salvar em csv foi descontinuado por falta de uso.
        #table_data.to_csv(dest_directory/"csv"/f"{table_name}.csv", sep=";", index=False)

        # Salvando em parquet (tudo como string... dtypes deverao ser atribuidos na leitura)
        table_data = table_data.astype(str)
        table_data.to_parquet(dest_directory/"parquet"/f"{table_name}.parquet", engine="fastparquet", index=False)

        if export_to_blob:
            # Definindo o Container e o Blob Name
            container_name = "luxorasap"
            blob_name = f"{blob_directory}/{table_name}.parquet"  #

            # Conversão para parquet em memória (sem precisar salvar local)
            table = pa.Table.from_pandas(table_data)
            parquet_buffer = io.BytesIO()
            pq.write_table(table, parquet_buffer)
            parquet_buffer.seek(0)  # Reseta o ponteiro para o início do buffer

            # Conectando ao Blob Storage
            connection_string = os.getenv('AZURE_STORAGE_CONNECTION_STRING')
            blob_service_client = BlobServiceClient.from_connection_string(conn_str=connection_string)

            # Criando um Blob Client
            blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)
            blob_client.upload_blob(parquet_buffer, overwrite=True)

__init__(luxorDB_directory=None)

Fornece uma forma padronizada de carregar tabelas para a luxorDB. 1. Possui metodos para carregar tabelas que já estao carregadas na memória - Sao os metodos que possuem 'table' no nome 2. Possui metodos para carregar arquivos de excel, com todas as suas abas Inclui metodo para checagem de alteracao de versao do arquivo - Sao os metodos que possuem 'file' no nome Args: luxorDB_directory (pathlib.Path, optional): Caminho completo ate o diretorio de destino dos dados.

Source code in LuxorASAP/luxorDB_dataloader.py
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
def __init__(self, luxorDB_directory = None):
    """Fornece uma forma padronizada de carregar tabelas para a luxorDB.
        1. Possui metodos para carregar tabelas que já estao carregadas na memória
            - Sao os metodos que possuem 'table' no nome
        2. Possui metodos para carregar arquivos de excel, com todas as suas abas
            Inclui metodo para checagem de alteracao de versao do arquivo
            - Sao os metodos que possuem 'file' no nome
    Args:
        luxorDB_directory (pathlib.Path, optional): Caminho completo ate o diretorio de destino dos dados.
    """
    self.luxorDB_directory = luxorDB_directory

    if self.luxorDB_directory is None:
        self.luxorDB_directory = Path(__file__).absolute().parent/"LuxorDB"/"tables"

    self.tracked_files = {}
    self.tracked_tables = {}

add_file_tracker(tracked_file_path, filetype='excel', sheet_names={}, excel_size_limit=None, index=False, index_name='index', normalize_columns=False)

Adiciona arquivo na lista para checar por alteracao Args: tracked_file_path (pathlib.Path): caminho completo ate o arquivo, incluindo nome do arquivo e extensão. sheet_names (dict, optional): Caso seja uma planilha com varias abas, mapear aqui o nome da aba para o nome do arquivo de saida desejado.

Source code in LuxorASAP/luxorDB_dataloader.py
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
def add_file_tracker(self, tracked_file_path, filetype="excel", sheet_names={}, 
        excel_size_limit = None,index=False, index_name="index",normalize_columns=False):
    """ Adiciona arquivo na lista para checar por alteracao
    Args:
        tracked_file_path (pathlib.Path): caminho completo ate o arquivo,
                incluindo nome do arquivo e extensão.
        sheet_names (dict, optional): Caso seja uma planilha com varias abas, mapear
                aqui o nome da aba para o nome do arquivo de saida desejado.
    """
    if tracked_file_path not in self.tracked_files:
        self.tracked_files[tracked_file_path] = {
                "last_mtime" : dt.datetime.timestamp(dt.datetime(2000,1,1)),
                "filetype" : filetype, "sheet_names": sheet_names,
                "excel_size_limit" : excel_size_limit,
                "index" : index,
                "index_name" : index_name,
                "normalize_columns" : normalize_columns,
            }

add_table_tracker(table_name)

Adiciona tabela na lista para controle de alteracao.

Source code in LuxorASAP/luxorDB_dataloader.py
116
117
118
119
120
def add_table_tracker(self, table_name:str):
    """ Adiciona tabela na lista para controle de alteracao."""

    if table_name not in self.tracked_tables:
        self.tracked_tables[table_name] = dt.datetime.timestamp(dt.datetime(2000,1,1))

is_file_modified(tracked_file_path)

Checa se o arquivo foi modificado desde a ultima leitura. Returns: tuple(bool, float): (foi modificado?, timestamp da ultima modificacao)

Source code in LuxorASAP/luxorDB_dataloader.py
135
136
137
138
139
140
141
142
143
144
145
def is_file_modified(self, tracked_file_path: Path) -> {bool, float}:
    """ Checa se o arquivo foi modificado desde a ultima leitura.
    Returns:
        tuple(bool, float): (foi modificado?, timestamp da ultima modificacao)
    """

    file_data = self.tracked_files[tracked_file_path]

    last_saved_time = file_data["last_mtime"]
    file_last_update = tracked_file_path.stat().st_mtime
    return file_last_update > last_saved_time, file_last_update

load_file_if_modified(tracked_file_path, export_to_blob=False, blob_directory='enriched/parquet')

Carrega arquivo no caminho indicado, carregando na base de dados caso modificado. Args: tracked_file_path (pathlib.Path): caminho ate o arquivo(cadastro previamente por add_file_tracker) type_map (type, optional): description. Defaults to None. filetype (str, optional): description. Defaults to "excel".

Source code in LuxorASAP/luxorDB_dataloader.py
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
def load_file_if_modified(self, tracked_file_path, export_to_blob=False, blob_directory='enriched/parquet'):
    """Carrega arquivo no caminho indicado, carregando na base de dados caso modificado.
    Args:
        tracked_file_path (pathlib.Path): caminho ate o arquivo(cadastro previamente por add_file_tracker)
        type_map (_type_, optional): _description_. Defaults to None.
        filetype (str, optional): _description_. Defaults to "excel".
    """
    file_data = self.tracked_files[tracked_file_path]

    last_saved_time = file_data["last_mtime"]
    filetype = file_data["filetype"]
    file_sheets = file_data["sheet_names"]

    file_last_update = tracked_file_path.stat().st_mtime

    if file_last_update > last_saved_time: # Houve alteracao no arquivo
        if filetype == "excel":
            file_sheets = None if len(file_sheets) == 0 else list(file_sheets.keys())

            # tables sera sempre um dicionario de tabelas
            tables = None
            trials = 25
            t_counter = 1
            while trials - t_counter > 0:
                try:
                    tables = pd.read_excel(tracked_file_path, sheet_name=file_sheets)
                    t_counter = trials # leitura concluida
                except PermissionError:

                    logger.error(f"Erro ao tentar ler arquivo '{tracked_file_path}.\nTentativa {t_counter} de {trials};'.\nSe estiver aberto feche.")
                    time.sleep(10)
                    t_counter += 1

            for sheet_name, table_data in tables.items():

                table_name = sheet_name if file_sheets is None else file_data["sheet_names"][sheet_name]

                if table_name == "trades":
                    table_data["ID"] = table_data.index

                self.__export_table(table_name, table_data, index=file_data["index"], index_name=file_data["index_name"],
                                        normalize_columns=file_data["normalize_columns"], export_to_blob=export_to_blob,
                                        blob_directory=blob_directory)
            self.tracked_files[tracked_file_path]["last_mtime"] = file_last_update

load_table_if_modified(table_name, table_data, last_update, index=False, index_name='index', normalize_columns=False, do_not_load_excel=False, export_to_blob=False, blob_directory='enriched/parquet', is_data_in_bytes=False, bytes_extension='.xlsx')

Parameters:

Name Type Description Default
table_name str

nome da tabela (sera o mesmo do arquivo a ser salvo)

required
table_data DataFrame

tabela de dados

required
last_update timestamp

timestamp da ultima edicao feita na tabela

required
Source code in LuxorASAP/luxorDB_dataloader.py
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
def load_table_if_modified(self, table_name, table_data, last_update, index=False, index_name="index", normalize_columns=False,
                           do_not_load_excel=False, export_to_blob=False, blob_directory='enriched/parquet',
                           is_data_in_bytes=False, bytes_extension=".xlsx"):
    """
    Args:
        table_name (str): nome da tabela (sera o mesmo do arquivo a ser salvo)
        table_data (pd.DataFrame): tabela de dados
        last_update (timestamp): timestamp da ultima edicao feita na tabela
    """

    if table_name not in self.tracked_tables:
        self.add_table_tracker(table_name)


    last_update_time = self.tracked_tables[table_name]
    if last_update > last_update_time:

        self.tracked_tables[table_name] = last_update
        self.__export_table(table_name, table_data, index=index, index_name=index_name, normalize_columns=normalize_columns,
                            do_not_load_excel=do_not_load_excel, export_to_blob=export_to_blob, blob_directory=blob_directory,
                            is_data_in_bytes=is_data_in_bytes, bytes_extension=bytes_extension)

scan_files(export_to_blob=False, blob_directory='enriched/parquet')

Para todos os arquivos cadastrados, vai buscar e carregar quando houver arquivo mais recente.

Source code in LuxorASAP/luxorDB_dataloader.py
222
223
224
225
226
227
228
229
230
def scan_files(self, export_to_blob=False, blob_directory='enriched/parquet'):
    """
        Para todos os arquivos cadastrados, vai buscar e carregar quando houver
        arquivo mais recente.
    """

    for file in self.tracked_files:

        self.load_file_if_modified(file, export_to_blob=export_to_blob, blob_directory=blob_directory)

text_to_lowercase(t)

Converte todas as colunas de texto para lowercase Args: t (pd.DataFrame): pandas DataFrame Returns: pd.DataFrame

Source code in LuxorASAP/luxorDB_dataloader.py
84
85
86
87
88
89
90
91
92
93
def text_to_lowercase(self, t):
    """
    Converte todas as colunas de texto para lowercase
    Args:
        t (pd.DataFrame): pandas DataFrame
    Returns:
        pd.DataFrame
    """

    return t.map(lambda x: x.lower().strip() if isinstance(x, str) else x)