scieee AI-readable full text Open interactive document viewer

Data Lakes em ambientes híbridos Cloud/Edge

Costa, Daniel Vilar da

Abstract

A análise dos dados tem sido, tradicionalmente, realizada em servidores na nuvem, onde a capacidade de armazenamento e de processamento são quase ilimitadas. Em contrapartida, os dispositivos periféricos têm severas limitações tanto de armazenamento como de processamento. No entanto, estes dispositivos encontram-se mais próximos do local onde os dados são gerados. Por causa disso, estes são, usualmente, utilizados para cargas de trabalho transacionais onde a confiabilidade e interatividade são fulcrais. Devido às limitações dos dispositivos periféricos, os dados são, geralmente, extraídos periodicamente para a nuvem onde são depois armazenados e processados. De modo a permitir a análise exploratória de dados heterogéneos, é comum utilizar uma infraestrutura Data Lake que permite gerir dados em formato bruto de múltiplas fontes. No entanto, transferir todos os dados coletados para a nuvem é inviável devido à limitada capacidade da rede que não tem conseguido acompanhar o crescimento do volume de dados coletados. Esta dissertação ultrapassa estes desafios ao implementar um componente middleware capaz de armazenar os dados previamente transmitidos na nuvem e propaga partes da interrogação para a periferia. Deste modo, consegue-se reduzir o volume de dados transferido ao enviar, idealmente, apenas uma vez os dados necessários para responder aos pedidos. Além disso, esta solução equilibra o impacto na rede e o custo computacional na periferia de modo a minimizar o tempo de execução.

Full text

Universidade do Minho Escola de Engenharia Daniel Vilar da Costa Data Lakes em ambientes híbridos Cloud/Edge Fevereiro, 2022 Universidade do Minho Escola de Engenharia Daniel Vilar da Costa Data Lakes em ambientes híbridos Cloud/Edge Dissertação de Mestrado Mestrado Integrado em Engenharia Informática Trabalho efetuado sob a orientação do(a) Ricardo Manuel Pereira Vilaça José Orlando Roque Nascimento Pereira Fevereiro, 2022 ii DIREITOS DE AUTOR E CONDIÇÕES DE UTILIZAÇÃO DO TRABALHO POR TERCEIROS Este é um trabalho académico que pode ser utilizado por terceiros desde que respeitadas as regras e boas práticas internacionalmente aceites, no que concerne aos direitos de autor e direitos conexos. Assim, o presente trabalho pode ser utilizado nos termos previstos na licença abaixo indicada. Caso o utilizador necessite de permissão para poder fazer um uso do trabalho em condições não previstas no licenciamento indicado, deverá contactar o autor, através do RepositóriUM da Universidade do Minho. Licença concedida aos utilizadores deste trabalho Creative Commons Atribuição-NãoComercial-CompartilhaIgual 4.0 Internacional CC BY-NC-SA 4.0 https://creativecommons.org/licenses/by-nc-sa/4.0/deed.pt iii DECLARAÇÃO DE INTEGRIDADE Declaro ter atuado com integridade na elaboração do presente trabalho académico e confirmo que não recorri à prática de plágio nem a qualquer forma de utilização indevida ou falsificação de informações ou resultados em nenhuma das etapas conducente à sua elaboração. Mais declaro que conheço e que respeitei o Código de Conduta Ética da Universidade do Minho. Agradecimentos Parcialmente financiado pelo projeto AIDA – Adaptive, Intelligent and Distributed Assurance Platform (POCI-01-0247-FEDER-045907), cofinanciado pelo Fundo Europeu de Desenvolvimento Regional (FEDER) através do Programa Operacional da Competitividade e Internacionalização (COMPETE 2020) e pela Fundação para a Ciência e Tecnologia (FCT) no âmbito do CMU Portugal. iv Resumo Data Lakes em ambientes híbridos Cloud/Edge A análise dos dados tem sido, tradicionalmente, realizada em servidores na nuvem, onde a capacidade de armazenamento e de processamento são quase ilimitadas. Em contrapartida, os dispositivos periféricos têm severas limitações tanto de armazenamento como de processamento. No entanto, estes dispositivos encontram-se mais próximos do local onde os dados são gerados. Por causa disso, estes são, usualmente, utilizados para cargas de trabalho transacionais onde a confiabilidade e interatividade são fulcrais. Devido às limitações dos dispositivos periféricos, os dados são, geralmente, extraídos periodicamente para a nuvem onde são depois armazenados e processados. De modo a permitir a análise exploratória de dados heterogéneos, é comum utilizar uma infraestrutura Data Lake que permite gerir dados em formato bruto de múltiplas fontes. No entanto, transferir todos os dados coletados para a nuvem é inviável devido à limitada capacidade da rede que não tem conseguido acompanhar o crescimento do volume de dados coletados. Esta dissertação ultrapassa estes desafios ao implementar um componente middleware capaz de armazenar os dados previamente transmitidos na nuvem e propaga partes da interrogação para a periferia. Deste modo, consegue-se reduzir o volume de dados transferido ao enviar, idealmente, apenas uma vez os dados necessários para responder aos pedidos. Além disso, esta solução equilibra o impacto na rede e o custo computacional na periferia de modo a minimizar o tempo de execução. Palavras-chave: Ambiente Cloud/Edge, Sincronização, Replicação, Federação de dados, Análise de dados exploratória v Abstract Data Lakes in hybrid Cloud/Edge environments Data analysis has traditionally been performed on dedicated servers in the cloud, where storage and processing capabilities are almost unlimited, in contrast to edge devices. Nonetheless, these devices are closer to where data is generated. Because of this, they have, usually, a transactional workload, where reliability and interactivity are essential. Due to the limitations of edge devices, generally, data is extracted periodically to the cloud to be stored and processed. In order to allow exploratory data analysis, the heterogeneous data is stored in a Data Lake infrastructure that manages data in raw format from multiple data sources. Nonetheless, transferring all collected data to the cloud is unfeasible because the increase in the volume of collected data has surpassed the network capabilities. This thesis overcomes these challenges by employing a middleware component capable of storing previously transmitted data in the cloud and pushing down query fragments to the edge. Consequently, the volume of data transmitted to the cloud is reduced by uploading, ideally, only once the required data. Furthermore, the solution balances the impact on the network and the computational effort in the edge in order to minimize execution time. Keywords: Cloud/Edge environment, Synchronization, Replication, Data federation, Exploratory data analysis vi Índice Lista de Figuras x Lista de Tabelas xi Glossário xii Siglas xiii Símbolos xv 1 Introdução 1 1.1 Contextualização .................................. 1 1.2 Motivação ..................................... 2 1.3 Objetivos ...................................... 2 1.4 Contribuições .................................... 3 1.4.1 Publicações ................................ 3 1.5 Organização do documento ............................. 4 2 Trabalho relacionado 5 2.1 Técnicas de sincronização de réplicas ........................ 5 2.1.1 Políticas de atualização ........................... 6 2.1.2 Métodos de replicação ........................... 9 2.1.3 Comparação das técnicas de sincronização ................. 15 2.2 Sistemas multi-store ................................ 15 2.2.1 Linguagem de interrogação ......................... 16 2.2.2 Arquitetura ................................. 17 2.2.3 CloudMdsQL ................................ 19 2.2.4 Trino .................................... 20 2.2.5 Apache Drill ................................ 20 2.2.6 Dremio .................................. 21 vii ÍNDICE 2.2.7 Delta Lake ................................. 22 2.2.8 PostgreSQL FDW .............................. 22 2.2.9 Comparação dos sistemas multi-store .................... 23 2.3 Discussão ..................................... 23 3 Design e implementação 25 3.1 Arquitetura ..................................... 25 3.1.1 Conector .................................. 27 3.1.2 Cache ................................... 27 3.2 Implementação ................................... 28 3.2.1 Processo .................................. 28 3.2.2 Resolução de colisões ........................... 30 3.2.3 Tolerância a faltas ............................. 30 3.2.4 Transações e isolamento .......................... 31 3.2.5 Outras funcionalidades ........................... 32 3.2.6 Limitações ................................. 33 3.3 Algoritmo de sincronização ............................. 33 3.3.1 Suporte a seleções ............................. 34 3.3.2 Suporte a projeções ............................ 39 3.3.3 Implementação .............................. 40 4 Avaliação 41 4.1 Avaliação dos métodos de replicação ........................ 41 4.1.1 Microbenchmark .............................. 42 4.1.2 Ambiente de testes ............................. 42 4.1.3 Resultados ................................. 42 4.1.4 Análise / Discussão de resultados ..................... 46 4.2 Avaliação do algoritmo de sincronização adaptativo ................. 47 4.2.1 Microbenchmark .............................. 48 4.2.2 Ambiente de testes ............................. 48 4.2.3 Resultados ................................. 48 4.2.4 Análise / Discussão de resultados ..................... 51 4.3 Benchmark ..................................... 51 4.3.1 CH-benCHmark .............................. 52 4.3.2 Visão global ................................ 52 4.3.3 Esquema ................................. 53 4.3.4 Transações ................................. 54 viii Símbolos 𝛼Fator de sobrecarga 𝑏𝑟Número máximo estimado de bytes que satisfazem o filtro 𝑐Número de condições 𝑐𝑏Custo por byte transferido 𝑐𝑐Custo por condição 𝑐𝑒Custo de estimar o número de linhas 𝑐𝑓Custo total de realizar a filtragem 𝑐𝑟Número de condições do filtro 𝑐𝑡Custo total para transferência dos dados 𝑓Número de filtros não estimados 𝑟Número de linhas na fonte de dados 𝑟𝑓Número marginal estimado de linhas filtradas pelo filtro 𝑟𝑟Número de linhas estimado que satisfazem o filtro 𝑤Tamanho médio de cada linha em bytes xv Capítulo 1 Introdução Este capítulo apresenta o contexto onde a dissertação se insere na Seção 1.1. Do mesmo modo, a Seção 1.2 descreve quais os problemas e oportunidades que motivaram o desenvolvimento de uma nova solução. A Seção 1.3 enumera os principais objetivos desta dissertação. A Seção 1.4 exibe as contribuições feitas por este. 1.1 Contextualização Os dispositivos de Internet das Coisas têm limitações tanto no processamento como no armazenamento quando comparado com a capacidade abundante da nuvem. De modo a colmatar restrições dos dispositivos periféricos, a abordagem mais comum consiste em processar e armazenar os dados na nuvem [5]. No entanto, esta abordagem tem ficado limitada pelos custos e pela capacidade da rede que, apesar do progresso nestas tecnologias, não tem conseguido acompanhar os volumes de dados transferidos para a nuvem [29]. Em contrapartida, os dispositivos de periferia têm cada vez maior capacidade de processamento e armazenamento que podem ser explorados [37]. O crescimento do volume de dados não estruturado levou à criação de Data Lakes capazes de armazenar elevadas quantidades de dados heterogéneos, que podem ser depois analisados. Os dados são coletados e armazenados no formato bruto e processados apenas quando for realizada a consulta dos dados. Deste modo, a decisão sobre quais dados devem ser extraídos pode ser realizada quando necessário ao invés de pré-determinado. Esta opção permite maior flexibilidade visto que nem sempre é possível saber de antemão quais os dados relevantes. No entanto, ao contrário dos armazéns de dados, que processam os dados quando coletados, os Data Lakes demonstram piores desempenhos nas consultas porque estes realizam o processamento durante a leitura [16]. 1 CAPÍTULO 1. INTRODUÇÃO A heterogeneidade das fontes de dados e dos seus dados apresentam obstáculos à sua análise [7]. Para resolver estes problemas, foram criadas linguagens de consulta comum que possam ser utilizadas pelas diferentes categorias de dados. Estas linguagens oferecem uma abstração sobre a fonte de dados e, por consequência, fornecem maior flexibilidade e facilitam a extração e processamento dos dados. 1.2 Motivação A crescente utilização de dispositivos periféricos conduziu ao surgimento de novos desafios que precisam de ser solucionados. Do mesmo modo, a nuvem é cada vez mais utilizada para completar o processamento dos dados recolhidos. A proximidade dos dispositivos periféricos à geração de dados e o aumento da capacidade de armazenamento e processamento permite reduzir o volume de dados enviados para a nuvem e, assim, reduzir custos e a latência. As vantagens dos dispositivos periféricos devem ser aproveitadas de modo a melhorar o desempenho e reduzir os custos do sistema. A necessidade de recolha de dados tem crescido, motivada, em grande parte, pelo sucesso da aprendizagem automática. Do mesmo modo, o surgimento de diferentes bases de dados, nomeadamente NoSQL, conduziram a uma diversificação do modo como os dados são armazenados e consultados. A necessidade de aproveitar as caraterísticas e funcionalidades de cada uma das fontes de dados levam a que seja necessário desenvolver uma linguagem poliglota que facilite a consulta de cada fonte de dados. Os sistemas na nuvem necessitam de suportar exploração dos dados e interrogações ad hoc. Nestas cargas de trabalho não é exequível determinar, de antemão, quais interrogações serão realizadas. Em consequência, não é possível determinar quais os dados que serão relevantes. Além disso, é recorrente a necessidade de sistemas capazes de responder a interrogações analíticas que necessitam de estar atualizados. 1.3 Objetivos O objetivo deste trabalho consiste em desenvolver um mecanismo de consulta para uma infraestrutura Data Lake utilizando um modelo de consulta comum. Este mecanismo considera os recursos disponíveis da fonte de dados e faz uma distribuição das tarefas entre a periferia e a nuvem de modo a minimizar a transferência e o tempo de processamento dos dados. Por um lado, pretende-se ultrapassar as limitações da rede transferindo idealmente apenas uma única vez os dados necessários. Por outro, espera-se que o mecanismo desenvolvido consiga responder às interrogações com os dados mais atuais. Além disso, o sistema deve ser capaz de ser utilizado para fontes de dados heterogéneas. Neste trabalho também pretende-se encontrar uma nova solução que consiga reduzir o volume de dados transferido dos dispositivos periféricos para a nuvem. A solução deve estar preparada para responder eficientemente a interrogações ad hoc, onde não é possível determinar quais os dados são relevantes. 2 CAPÍTULO 1. INTRODUÇÃO Por fim, pretende-se perceber se a solução desenvolvida pode ser utilizada para reduzir os tempos de execução de interrogações analíticas e o volume de dados transferidos dos dispositivos periféricos para a nuvem. 1.4 Contribuições Esta dissertação foi desenvolvido no âmbito do projeto AIDA, cofinanciado pelo FEDER, através do COMPETE e pela FCT. Nesta dissertação realizaram-se as seguintes contribuições: • Arquitetura Data Lakes para um ambiente Cloud/Edge. A arquitetura realiza consultas sobre os dados que se encontram na periferia e mantém os dados numa cache na nuvem para evitar retransmissões. • Algoritmo de sincronização com suporte de push-down de filtros e projeções. Este permite realizar a sincronização apenas dos dados necessários para responder às interrogações, reduzindo deste modo o volume de dados transmitido para a nuvem. • Algoritmo adaptativo de sincronização que tem em conta a capacidade de processamento do dispositivo remoto, da largura de banda da rede e da carga de trabalho de modo a balancear o tempo de processamento na periferia e o volume de dados transmitido. • Implementação da arquitetura e de algoritmos em PostgreSQL. Utilizou-se Foreign Data Wrappers para criar um middleware capaz de intercetar as interrogações feitas à fonte de dados. • Adaptação do CH-benCHmark a uma carga de trabalho orientada ao evento. Este benchmark realiza apenas operações de inserção e assemelha-se a uma carga de trabalho típica dos dispositivos periféricos. • Avaliação dos diferentes algoritmos de replicação existentes. Realizaram-se medições de tempo de execução, volume de dados transferido e do uso da memória auxiliar. • Avaliação do desempenho da solução proposta. Mediu-se o desempenho das interrogações analíticas na nuvem da solução proposta. Comparou-se também o impacto da solução no débito de transações nos dispositivos periféricos. 1.4.1 Publicações O artigo curto intitulado ”Adaptive Database Synchronization for an Online Analytical Cloud-to-Edge Continuum” foi aceite para ser publicado na 36.ª conferência virtual ACM/SIGAPP Symposium On Applied 3 CAPÍTULO 1. INTRODUÇÃO Computing. Neste pequeno artigo é descrito o algoritmo de sincronização adaptativo proposto nesta dissertação. O artigo intitulado ”AIDA-DB: A Data Management Architecture for the Edge and Cloud Continuum” foi aceite para ser publicado no 1.º workshop internacional em Secure FunctiON ChAining and FederaTed AI - SONATAI 2022 co-located with IEEE 19th Annual Consumer Communications Networking Conference (CCNC). Neste artigo é descrito uma arquitetura para gestão dos dados num ambiente Cloud/Edge. O middleware proposto faz parte integrante desta arquitetura. 1.5 Organização do documento O resto deste documento encontra-se organizado do seguinte modo: OCapítulo 2 descreve o trabalho relacionado com os principais métodos de sincronização de réplicas, bem como as suas políticas de atualização. Do mesmo modo, neste capítulo enumeram-se as principais arquiteturas e motores de consulta para fontes de dados heterogéneos. OCapítulo 3 descreve a arquitetura e a sua implementação. Neste capítulo descreve-se também quais os componentes envolvidos na implementação, o processo e a descrição das funcionalidades suportadas. Neste capítulo também é descrito o algoritmo de sincronização proposto. Por fim, no Capítulo 4 avaliou-se o desempenho de métodos de sincronização existentes. Além disso, comparou-se um destes métodos com o método de sincronização proposto. Por último, avaliou-se o desempenho do sistema desenvolvido nesta dissertação. 4 Capítulo 2 Trabalho relacionado A replicação é usualmente utilizada para obter alta disponibilidade, tolerância a faltas e aumentar o desempenho [34] [21]. Neste trabalho, a replicação dos dados é utilizada para melhorar o desempenho das interrogações. Os dados armazenados na réplica permitem reduzir a distância entre o local onde os dados são armazenados e onde estes são processados. A localidade dos dados [8] permite melhorar os desempenhos nas interrogações [22]. Os dados guardados localmente não necessitam de ser retransmitidos da fonte de dados. Esta propriedade permite reduzir o volume de dados transferidos e assim melhorar o desempenho do sistema. Quando se usa um armazém de dados, estes são guardados num repositório centralizado. Além disso, os dados encontram-se geralmente transformados para serem posteriormente consultados. Estas transformações permitem remover informação duplicada, normalizar ou remover dados [40]. Do mesmo modo, os repositórios encontram-se geralmente próximos dos servidores onde estes são processados para melhorar o desempenho. Esta abordagem permite tempos de acesso bastante reduzidos devido à proximidade dos dados ao seu processamento. Contudo, necessita que as transformações realizadas aos dados extraídos sejam planeados antes de estes serem consultados. Neste sentido, para extrair apenas os dados necessários, é necessário um conhecimento prévio sobre quais os dados serão consultados. Além disso, os armazéns de dados sincronizam periodicamente as réplicas. Durante o processo de sincronização, as diferenças entre a réplica e a fonte de dados são transmitidas para a nuvem. 2.1 Técnicas de sincronização de réplicas A redundância dos dados entre a fonte de dados e a réplica exige um sistema de sincronização. Os dados que sofrem modificações ao longo do tempo são consideradas mutáveis. Nesta dissertação, considera-se 5 CAPÍTULO 2. TRABALHO RELACIONADO que os dados podem sofrer modificações e essas devem ser refletidas na réplica. Em oposição, dados que não sofrem modificações são considerados imutáveis. O sistema de sincronização é responsável por detetar as mudanças nos dados da fonte e transferi-las para a nuvem. Neste caso, existem diversas técnicas de sincronização de réplicas, cada uma com as suas vantagens e desvantagens. A comparação das técnicas de sincronização deve considerar diferentes fatores como as operações suportadas, tamanho da mensagem, o tempo de execução e a memória utilizada na fonte de dados. Geralmente, a capacidade dos dispositivos periféricos e a largura de banda da rede são limitados. Por estes motivos, tanto a memória como o tempo de processamento devem ser mínimos para poderem ser realizados em dispositivos periféricos. Do mesmo modo, o volume de dados transferidos dos dispositivos periféricos para a nuvem deve ser também mínimos para superar as limitações da capacidade da rede. 2.1.1 Políticas de atualização As réplicas de dados mutáveis necessitam de ser atualizadas com as fontes de dados. A frequência da atualização desses dados depende dos requisitos. A atualização das réplicas pode ser efetuada periodicamente, quando existe uma alteração na fonte de dados, sempre que é realizada uma consulta ou manualmente [4]. 2.1.1.1 Periodicamente A atualização periódica realiza a operação de sincronização,repetidamente, após um período pré-definido [31]. Também é comum atribuir um tempo de vida à replicação. Este é utilizado para determinar quando a réplica deixa de ser válida e que por isso não pode ser utilizada. Normalmente, a taxa de atualização e o tempo de vida escolhidos dependem dos requisitos, da efemeridade dos dados e da frequência de alterações destes. Esta política é bastante utilizada pela sua simplicidade e independência da implementação da fonte de dados. A fonte de dados é apenas consultada periodicamente para realizar a atualização. Entre as atualizações, todas as consultas são realizadas apenas sobre os dados presentes na réplica evitando, assim, a latência das mensagens enviadas para a fonte de dados. Por outro lado, reduz o número de mensagens enviadas e recebidas da fonte de dados. Ainda assim, esta solução faz com que, entre sincronizações, os dados possam ficar desatualizados e os resultados das interrogações sejam incoerentes com a fonte de dados. 2.1.1.2 Ao alterar (On commit) A operação de sincronização pode ser realizada quando existe uma alteração na fonte de dados [28]. Esta operação de replicação pode ser síncrona, ou seja, a escrita na fonte de dados e na réplica são realizadas em simultâneo, ou assíncrona, onde se altera na fonte de dados primeiro e posteriormente na réplica [19]. 6 CAPÍTULO 2. TRABALHO RELACIONADO A replicação síncrona realiza as alterações na base de dados original e nas réplicas em simultâneo. Neste sentido, garante-se que os dados escritos na fonte de dados estão também presentes na réplica. Nesta política, a fonte de dados apenas efetiva as alterações quando tem conhecimento que todas as réplicas receberam a informação sobre a atualização. Isto garante que em caso de falha da fonte de dados, não são perdidas as atualizações. No entanto, o desempenho deteriora-se devido à latência da comunicação entre a fonte e a réplica. Além disso, a fonte e a réplica ficam acoplados, ou seja, a fonte necessita de conhecer as réplicas e a modificação dos dados fica dependente destas. Consequentemente, a falha de uma das réplicas tem impacto no desempenho da fonte de dados. A replicação assíncrona não necessita da confirmação das réplicas para realizar a modificação. Neste sentido, a fonte de dados apenas transmite informação das modificações realizadas para as réplicas após realizar a atualização. As réplicas, ao receberem a informação sobre a alteração na fonte, modificam os seus dados para terem a cópia mais recente. O tempo de resposta é menor por não ser necessário esperar pela resposta das réplicas. Todavia, em caso de falha, não é garantido que a réplica tenha a cópia mais atualizada dos dados. No entanto, a réplica necessita de ter conhecimento de todas as réplicas. Utilizando esta solução, as interrogações podem ser realizadas apenas à réplica sem necessidade de consultar diretamente a fonte de dados. Deste modo, evita-se a latência da interrogação. Esta abordagem envia mensagens por cada alteração dos dados. O custo adicional por cada mensagem pode ser prejudicial principalmente quando o número de atualizações é substancialmente superior às consultas. Este custo é consequência de exigir comunicação com a réplica sempre que é realizada uma atualização [31]. Esta solução não é adequada, portanto, para fontes de dados com mudanças frequentes. No entanto, para fontes de dados onde o número de consultas é superior ao número de atualizações, esta política evita consultas diretas à fonte de dados. Esta política tem a desvantagem de transferir todas as atualizações. Ou seja, pode transferir atualizações de dados que nunca serão lidos e pode transferir atualizações que não são necessárias como é o caso de inserções seguidas de remoções antes de uma consulta. 2.1.1.3 Ao consultar (On demand) A sincronização da réplica pode ser realizada sempre que existe uma consulta à fonte de dados [31]. Deste modo, sempre que é realizada uma consulta este integra os dados da réplica local com os novos dados da fonte. Desta forma, os dados transferidos da fonte para a réplica não refletem os estados intermédios. Além disso, o volume de dados transferido da fonte de dados é reduzido porque apenas as diferenças relevantes para responder à interrogação são enviadas. Este método garante também que a consulta contém dados coerentes com a fonte de dados. Apesar disso, esta solução necessita de consultar a fonte de dados em todas as interrogações. Se o número de interrogações for superior ao de atualizações, então muitos dos pedidos não retornaram novos dados. Este custo deve-se ao facto de apenas existir comunicação com a fonte de dados quando é realizada uma consulta. Além disso, esta solução tem maior latência devido à necessidade de consultar a 7 CAPÍTULO 2. TRABALHO RELACIONADO fonte de dados. Esta solução, no entanto, obriga a atualizar a réplica em todas as leituras que retornarem alterações aumentando o tempo de execução da interrogação. 2.1.1.4 Discussão Diferentes políticas de atualização tem diferentes características que devem ser consideradas para cada problema. No contexto desta dissertação, pretende-se não só reduzir o impacto na rede e na fonte de dados como também manter a réplica atualizada com os dados mais recentes. Por esse motivo, as políticas de atualização são avaliadas tendo em conta diferentes requisitos. Primeiramente, deve conseguir reduzir o volume de dados transferido da fonte de dados para a réplica. Esta propriedade permite reduzir o impacto que a transferência dos dados tem no tempo de sincronização. Além disso, permite reduzir custos e evita sobrecarregar a rede. Por outro lado, a réplica deve permanecer coerente com a fonte de dados. Esta propriedade garante que os dados lidos da réplica representam o estado da fonte de dados. Além disso, o impacto da sincronização na fonte de dados deve ser mínimo. Neste sentido, a política de atualização deve evitar operações de sincronização de modo a reduzir o impacto nos dispositivos periféricos. Por fim, de modo a tornar o sistema de replicação compatível com uma grande variedade de fontes de dados, a técnica de replicação não deve estar dependente da implementação interna. Assim, a implementação pode ser genérica e capaz de ser utilizada para múltiplas fontes de dados. Tendo em conta os requisitos da dissertação, a política periódica não garante que a réplica esteja sempre coerente com a base de dados e por consequência não permite que sejam aplicadas interrogações em tempo real. No entanto, esta solução pode ser aplicada para diferentes fontes de dados e reduz significativamente a latência e o impacto na fonte de dados quando o número de consultas é superior ao número de atualizações realizadas. A política de atualizar quando existe uma alteração na fonte de dados não pode ser adaptada a todas as bases de dados porque depende da implementação interna. Todavia, esta solução pode ser aplicada a um subconjunto de fontes de dados que suportem esta política. Esta política é especialmente vantajosa quando as consultas são mais frequentes do que as alterações dos dados. Esta solução teria, no entanto, de utilizar uma sincronização assíncrona para reduzir a latência de escrita nas fontes de dados. A atualização ao consultar cumpre todos os requisitos. Contudo, esta abordagem favorece a consistência dos dados em detrimento de latência de resposta. Por este motivo, esta política pode ter piores desempenhos que as restantes alternativas. Neste sentido, esta política tem maiores latências, mas compensa pela flexibilidade de poder ser utilizada em diversas fontes de dados e por garantir que os dados estão sempre atualizados. Esta política é especialmente vantajosa quando o número de alterações na fonte de dados é superior ao número de consultas realizadas. 8 CAPÍTULO 2. TRABALHO RELACIONADO Tabela 1: Tabela comparativa das políticas de atualização. Política Consistência Latência Compatibilidade Carga de trabalho ideal Periódica ✘Baixa Extensa Núm. consultas > Núm. atualizações Ao alterar ✔Baixa Limitada Núm. consultas > Núm. alterações Ao consultar ✔Alta Extensa Núm. consultas < Núm. alterações Em suma, as diferentes políticas de atualização adequam-se a diferentes condições. Como representado na Tabela 1, a política de atualização depende das condições e da carga de trabalho. Neste sentido, atualizar a réplica apenas quando é realizada a consulta adequa-se ao ambiente descrito porque permite ter os dados sempre coerentes. Do mesmo modo, espera-se que o número de alterações à fonte de dados seja superior ao número de consultas realizadas. Apesar disso, esta abordagem necessita de consultar a fonte de dados no momento da consulta resultando em maior latência. 2.1.2 Métodos de replicação Os métodos de replicação são utilizados para extrair dados das fontes de dados de modo a atualizar as réplicas. Estes utilizam algoritmos para determinar quais os dados que não estão atualmente presentes na réplica. Os métodos abordados fazem um balanço entre o espaço de armazenamento extra, o tempo de processamento necessário e a precisão. Como diferentes métodos são adequados a diferentes sistemas, enumerou-se alguns dos métodos de replicação e as suas vantagens e desvantagens. 2.1.2.1 Completa O método de replicação completa consiste em atualizar totalmente a réplica quando existe uma sincronização. Este substitui a réplica pelo último resultado da interrogação. Este método pode ser utilizado quando existe uma política de atualização periódica. Geralmente, este método é utilizado para obter os dados pela primeira vez da fonte de dados porque não adiciona complexidade e envia o volume mínimo de informação necessário à sincronização [4]. Este método é o mais simples e pode ser implementado genericamente para qualquer fonte de dados. Esta abordagem não necessita de manter informação adicional sobre os dados nem de processamento extra para detetar diferenças entre a réplica e a fonte de dados. Apesar de ser simples e requerer menos espaço em disco, esta abordagem solicita dados que já se encontram replicados. Assim, quando o número de alterações é inferior ao tamanho dos dados, este método envia muitos dados que podiam ser evitados. Este método é também muito ineficiente para grandes volumes de dados visto que estes necessitam de ser repetidamente enviados por completo. Devido ao facto desta operação ser extensa para grandes volumes de dados, esta pode gerar latência nas consultas, dependendo da política de atualização. 9 CAPÍTULO 2. TRABALHO RELACIONADO Tabela 3: Tabela comparativa da complexidade computacional dos métodos de sincronização. Método Computação Memória persistente auxiliar Memória volátil auxiliar Mensagem Completo 𝑂(𝑛)𝑂(1)𝑂(1)𝑂(𝑛) Chave 𝑂(𝑛)*/𝑂(𝑑+log𝑛)** 𝑂(𝑛)𝑂(1)𝑂(𝑑) Log 𝑂(𝑢)𝑂(𝑢)𝑂(1)𝑂(𝑑) Hash 𝑂(𝑛+𝑑)𝑂(𝑛)𝑂(𝑛)𝑂(𝑑) CPI 𝑂(𝑛+𝑑3)*** 𝑂(1)𝑂(𝑑)𝑂(𝑑) BCH 𝑂(𝑛+𝑑2)*** 𝑂(1)𝑂(𝑑)𝑂(𝑑) IBLT 𝑂(𝑛+𝑑)*** 𝑂(1)𝑂(𝑑)𝑂(𝑑) 𝑛Número de linhas dos dados na fonte. 𝑑Diferença entre os dados na fonte e na réplica. 𝑢Número de alterações à fonte de dados. *Pior caso quando a chave não seja indexada. ** Melhor caso quando a chave seja indexada. *** Amortizado. dados comum e uma única linguagem para consultar as fontes de dados [24]. As principais arquiteturas são armazéns de dados, sistemas de bases de dados federadas e sistemas sem esquema [33]. Os diferentes sistemas podem ser comparados pela linguagem de interrogação que suportam, as fontes de dados que podem suportar e otimizações realizadas para reduzir o tempo de execução. 2.2.1 Linguagem de interrogação Os sistemas multi-store usam normalmente uma linguagem textual para consultar os dados provenientes da fonte de dados. Grande parte dos sistemas utiliza a linguagem Structured Query Language (SQL) ou semelhante. Esta linguagem é bastante comum visto ser por norma utilizada para consulta de fontes de dados relacionais. Ainda assim, existem outros sistemas que seguem outras abordagens para superar as limitações existentes com a linguagem SQL. A linguagem SQL é a linguagem mais popular para consulta de dados relacionais [33]. Esta linguagem oferece diversas funcionalidades desde as mais básicas como projeções e filtros até agregações e janelas. Para além de consultas, esta linguagem suporta funcionalidades de manipulação de dados conferindo em alguns sistemas a possibilidade de modificar a fonte de dados. Apesar de a linguagem SQL ser bastante flexível, esta tem diversas limitações. Nos últimos anos, aumentou o número de bases de dados não relacionais. Em geral, estas bases de dados utilizam uma linguagem própria e especializada para tirar o máximo proveito das capacidades da fonte de dados. Em especial, as bases de dados NoSQL são conhecidas por não terem esquema rígido, não terem os dados normalizados e suportarem estruturas como arrays e dicionários. Esta flexibilidade de armazenamento dos dados não é muitas vezes compatível com a linguagem SQL. Por este motivo, grande parte dos 16 CAPÍTULO 2. TRABALHO RELACIONADO sistemas multi-store suporta funções especiais que permitem converter estruturas em tabelas [2]. Esta funcionalidade permite converter modelos orientados ao objeto ou ao documento num modelo relacional. Por outro lado, diferentes fontes de dados tem requisitos e funcionalidades diferentes. Apesar de a linguagem SQL ter um leque abrangente de funções que permitem transformar os dados, estas são limitadas. Em especial, as bases de dados orientadas ao grafo necessitam de funções especiais para efetuar travessias. A linguagem ANSI SQL não suporta estas funcionalidades, mas podem ser utilizadas funções para superar estas limitações. Estas funções utilizam o conhecimento sobre a fonte de dados e conseguem tirar proveito das funcionalidades expostas da base de dados [24]. Em suma, a linguagem SQL é muito popular facilitando, deste modo, a sua adoção e o seu uso. Contudo, a linguagem foi desenvolvida para um modelo de dados relacional e que não está adaptado à flexibilidade das bases de dados não relacionais. Para superar estas limitações é normalmente utilizado uma extensão ao ANSI SQL com funções que convertem estruturas complexas em tabelas. Para suportar funcionalidades específicas de cada fonte de dados, cria-se funções que encapsulam as capacidades da fonte de dados [2]. 2.2.1.1 Comparação entre subconsultas nativas e funções Alguns sistemas suportam subconsultas nativas [24], ou seja, permitem fazer consultas diretamente à fonte de dados utilizando a Application Programming Interface (API) da fonte de dados. Esta funcionalidade permite maior flexibilidade e desempenho. Além disso, não só utiliza uma linguagem comum como, por exemplo, SQL para consultar múltiplas fontes de dados como também permite realizar interrogações diretamente à fonte de dados utilizando uma linguagem de programação. No entanto, esta abordagem torna-se mais difícil de usar com o aumento do número de fontes de dados. Apesar disso, o código da consulta direta à fonte de dados está presente na interrogação facilitando a depuração de erros e a sua alteração. Apesar disso, reduz a legibilidade da interrogação. Em alternativa, alguns sistemas preferem utilizar funções que escondem a complexidade da interrogação nativa [2]. Esta abordagem permite encapsular o conhecimento sobre a fonte de dados e expõe uma versão simplificada, permitindo que seja mais fácil utilizar posteriormente. Em suma, a utilização de funções que encapsulam consultas nativas às fontes de dados é geralmente a melhor solução para melhorar a legibilidade. Ainda assim, para sistemas que utilizam poucas fontes de dados e necessitam de realizar consultas nativas frequentemente, é vantajoso utilizar as subconsultas nativas. 2.2.2 Arquitetura Um armazém de dados guarda centralmente os dados como representado na Figura 2. Neste caso, existe uma réplica dos dados da fonte num repositório centralizado. Os dados são previamente extraídos 17 CAPÍTULO 2. TRABALHO RELACIONADO Consumidor analítico Interrogação Fonte de Dados Fonte de Dados Fonte de Dados Extração, Transformação e Carregamento Réplica centralizada Figura 2: Arquitetura de um armazém de dados. e transformados para poderem ser descritos por um esquema pré-definido. Geralmente os dados são periodicamente atualizados para evitar que os dados estejam desatualizados. Neste caso, as consultas são realizadas diretamente à réplica em vez da fonte de dados. Aceder apenas à réplica tem grandes vantagens no desempenho das interrogações porque reduz a latência e tira proveito da localidade dos dados. Apesar disso, não garante que os dados estejam sempre atualizados. Os sistemas federados permitem realizar interrogações a diferentes fontes de dados utilizando uma arquitetura mediador-adaptador, como representado na Figura 3. Esta arquitetura consiste em ter um mediador responsável por gerir as fontes de dados e o adaptador converte os dados heterogéneo num formato comum. No caso de subconsultas nativas, existe um adaptador à linguagem de programação em vez de ser diretamente à fonte de dados. Esta abordagem tem diversas vantagens, nomeadamente tirar proveito das funcionalidades de cada fonte de dados e facilita a extensão de novas bases de dados [24]. Por outro lado, a abordagem sem esquema obtém o esquema no momento em que estes dados são consultados. Esta permite que a evolução do esquema dos dados seja feita facilmente. Apesar disso, esta abordagem obriga a ter maior conhecimento sobre a localização e formato dos dados e tem pior desempenhos porque é necessário obter o esquema frequentemente [2]. As diferentes arquiteturas apresentadas cumprem diferentes requisitos. Os armazéns de dados oferecem o melhor desempenho porque os dados já se encontram no formato comum e estão guardados localmente, oferecendo, consequentemente, ótima localidade. No entanto, esta arquitetura não permite que os dados estejam sempre atualizados. Por outro lado, um sistema federado, que utiliza a arquitetura mediador-adaptador, permite suportar 18 CAPÍTULO 2. TRABALHO RELACIONADO Consumidor analítico Interrogação Fonte de Dados Fonte de Dados Fonte de Dados Adaptador Adaptador Adaptador Mediador Figura 3: Arquitetura de um sistema federado. uma grande variedade de fontes de dados. Enquanto o adaptador fornece a capacidade de suportar a múltiplas fontes de dados e tirar proveito das suas capacidades, o mediador utiliza todo o conhecimento fornecido pelos adaptadores para otimizar das interrogações permitindo assim bom desempenho. Em contrapartida, a abordagem sem esquema facilita a utilização de fontes de dados não relacionais onde o esquema muda frequentemente. Contudo, esta abordagem tem impacto no desempenho das interrogações. Por este motivo, uma abordagem híbrida é geralmente utilizada, onde a atualização do esquema é realizada apenas quando necessário. Esta solução evita a descoberta constante do esquema e mantém a facilidade de uso de uma abordagem sem esquema. 2.2.3 CloudMdsQL CloudMdsQL é uma linguagem baseada em SQL para consulta de fontes de dados heterogéneas. Esta utiliza uma abordagem federada e sem esquema [24]. Este permite efetuar interrogações a múltiplas fontes de dados utilizando a combinação de interrogações SQL e blocos embebidos em Python que consultam abase de dados. Estas características permitem tirar proveito de todas as funcionalidades de fontes de dados NoSQL. O uso de subconsultas nativas exclui a necessidade de utilizar adaptadores facilitando assim a implementação. Este suporta um catálogo que contém toda a informação sobre as diferentes fontes de dados, nomeadamente índices e estatísticas. Este catálogo permite otimizar as interrogações. Na própria linguagem é 19 CAPÍTULO 2. TRABALHO RELACIONADO possível especificar estatísticas sobre os dados. Este catálogo é utilizado pelo planeador para escolher o plano mais eficiente. Esta linguagem suporta push-down de filtros, interrogações aninhadas e utiliza bind joins para juntar os dados de diferentes fontes de dados, incluindo interrogações nativas. Os filtros são propagados para as subconsultas de modo a reduzir o número de tuplos retornados. Do mesmo modo, o sistema tira partido das interrogações aninhadas e de bind joins, convertendo-os em filtros, para reduzir o volume de dados transferido. Estas otimizações permitem reduzir o tempo de execução e tirar proveito de otimizações das fontes de dados. 2.2.4 Trino Trino, anteriormente denominado Presto, é um motor de consulta distribuído que usa vários adaptadores para juntar diversas fontes de dados. Permite efetuar interrogações em diferentes fontes de dados e processar os dados em diferentes trabalhadores. A sua natureza distribuída permite que o processamento de dados seja escalável e eficiente [36]. A arquitetura é bastante semelhante ao CloudMdsQL, mas os dados podem ser processados por múltiplos trabalhadores que estão desacoplados da fonte de dados. Esta natureza distribuída obriga a que o coordenador tenha de gerir os múltiplos trabalhadores. Trino tem capacidade de obter metadados sobre as fontes de dados que são depois utilizados para otimizar as interrogações. O motor de consulta suporta ANSI SQL com pequenas extensões como funções lambda para facilitar o processamento de dados. Os adaptadores facilitam a integração com múltiplas fontes de dados e tira proveito dos adaptadores para paralelizar e reduzir o volume de dados enviados. A arquitetura permite realizar interrogações com baixa latência sendo indicado para interrogações muito complexas e demoradas. 2.2.5 Apache Drill Apache Drill é um motor de consulta distribuído que usa um adaptador para juntar diferentes fontes de dados [2]. Baseia-se no Dremel e utiliza a representação colunar dos dados aninhados e uma arquitetura em árvore para executar as consultas [26]. A cada fonte de dados está associado um trabalhador, denominado Drillbit. Cada Drillbit pode ser utilizado para agrupar dados de múltiplos trabalhadores. O Drillbit pode ser usado pelo cliente para consultar as fontes de dados. Quando é feita uma interrogação a um dos Drillbit, este torna-se coordenador do trabalho e otimiza a execução das tarefas. Estas tarefas são depois distribuídas para os restantes trabalhadores. Devido a esta natureza distribuída e não ter um coordenador central, os metadados são descentralizados. A arquitetura permite tirar proveito da localidade dos dados. O facto de o processamento dos dados estar mais perto da fonte reduz o volume de dados enviado entre trabalhadores e melhora o desempenho. 20 CAPÍTULO 2. TRABALHO RELACIONADO O otimizador utiliza os meta-dados para distribuir o trabalho com base na proximidade, capacidade e desempenho de cada um dos trabalhadores. A arquitetura dos trabalhadores segue um esquema em árvore. O coordenador é responsável por atribuir tarefas, receber as consultas e responder ao cliente. Existem trabalhadores intermediários responsáveis por processar os dados provenientes das folhas. As folhas são responsáveis por obter e processar os dados da fonte de dados. O sistema paraleliza não só o trabalho entre múltiplos trabalhadores como também suporta execução paralela dentro de cada trabalhador. O facto de ser altamente paralelizável e distribuído permite reduzir o tempo de execução quando o volume de dados é muito grande. Drill trata todos os dados como se estivessem em formato de tabela. Todavia, fornece sistemas de descoberta de esquema e funções que permitem facilmente converter dados aninhados para o formato tabelar. Este foi também desenvolvido para ser facilmente extensível a novas otimizações e permite adicionar novos adaptadores. 2.2.6 Dremio Dremio é um motor de consulta baseado no Apache Drill. Tem muitas das vantagens do Apache Drill, mas é mais eficiente devido a algumas otimizações. Este tira proveito do Apache Arrow para reduzir o volume de dados enviados e tempo de serialização e desserialização da fonte de dados para o motor de consulta. Além disso, os dados são vetorizados de modo a tirar proveito das vantagens do Single Instruction, Multiple Data (SIMD) e das capacidades do Graphics Processing Unit (GPU) [13]. Por outro lado, permite guardar, em formato colunar, os dados em bruto ou processados de modo a evitar consultas a fonte de dados. Esta materialização dos dados, denominadas pelo Dremio de Data Reflections, permite otimizar interrogações analíticas a ficheiros e fontes de dados com formatos pouco eficientes. A materialização de certas interrogações evita reprocessamento sendo utilizadas pelo otimizador para múltiplas interrogações. Além disso, os dados materializados são indexados pelos diferentes valores para melhorar o desempenho. O Dremio tem suporte a diferentes políticas de atualização que permitem refrescar periodicamente os dados materializados. Atualizar os dados materializados pode ser realizado tanto com o método de replicação completa ou de replicação incremental baseada em chave. Estas otimizações permitem melhorar significativamente o desempenho das interrogações quando comparado com os outros sistemas multi-store que apenas paralelizam o processamento dos dados em diferentes nodos. No entanto, a materialização, apesar de ser muito eficiente, pode ficar rapidamente incoerente com as fontes de dados principalmente com grande volume de alteração na fonte de dados. Por outro lado, o suporte a fontes de dados é reduzido quando comparado com outros motores. Estender a novas fontes de dados está limitada apenas a fontes de dados que suportam a interface SQL. 21 CAPÍTULO 2. TRABALHO RELACIONADO 2.2.7 Delta Lake Delta Lake é um armazenamento tabelar que suporta as propriedades Atomicity, Consistency, Isolation, Durability (ACID) [3]. Este armazena os dados num formato eficiente e adequado para consultas analíticas. Além disso, armazena logs com informação sobre os dados modificados, permitindo assim o suporte a transações e a realizar interrogações sobre estados passados. O suporte de transações permite que os dados sejam modificados sem o risco de, em caso de falha, os dados ficarem num estado inconsistente. Por outro lado, a mudança de esquema pode ser realizada transacionalmente. Os dados armazenados no Delta Lake podem ser acedidos através do motor analítico Apache Spark. Este motor consegue processar grandes volumes de dados e oferece uma interface que permite realizar operações analíticas em sistemas de dados distribuídos [41]. Esta abstrai grande parte da complexidade de fazer o processamento num ambiente distribuído. Por outro lado, o motor analítico oferece uma interface simples que pode ser implementada para diferentes fontes de dados. Além disso, oferece métodos que permitem otimizar a consulta de dados da fonte. 2.2.8 PostgreSQL FDW PostgreSQL é uma base de dados objeto-relacional que suporta fonte de dados externas através de Foreign Data Wrapper (FDW) [30]. O FDW permite abstrair a comunicação com uma fonte de dados externa. Deste modo, ao realizar uma consulta a uma fonte de dados externa não é preciso conhecer detalhes sobre a sua implementação. Deste modo o PostgreSQL pode ser utilizado como um sistema federado de bases de dados. Na base de dados é possível definir um servidor responsável por consultar e modificar a fonte de dados. Do mesmo modo, para cada servidor é possível associar múltiplas tabelas virtuais que representam os dados contidos na fonte de dados externa. Estas tabelas não são armazenadas no PostgreSQL. Se for necessário consultar os dados de uma ou mais tabelas de um servidor, então parte do plano de execução é propagado para o Foreign Data Wrapper Handler. Este é responsável por executar o plano de execução na fonte de dados remota. Por outro lado, a extensão Multicorn permite implementar facilmente Foreign Data Wrappers utilizando uma linguagem de alto nível. Utilizando uma interface mais simples do que a implementação nativa, é possível criar tabelas remotas para qualquer fonte de dados. Esta extensão recebe partes do plano de execução proveniente do PostgreSQL e chama métodos de uma classe escrita em Python. Por fim, o Multicorn permite o suporte de push-down de seleções, projeções e de ordenações. Estas operações podem ser aproveitadas pela implementação para melhorar o desempenho. PostgreSQL suporta execução paralela e assíncrona que permitem reduzir o tempo de execução. Além disso, pode ser utilizado como cluster e permitir que vários nodos executem a interrogação em paralelo. 22 CAPÍTULO 2. TRABALHO RELACIONADO Por outro lado, cada nodo pode ter vários trabalhadores em simultâneo. Por ser uma base de dados, podem-se armazenar dados diretamente no PostgreSQL. 2.2.9 Comparação dos sistemas multi-store Os diferentes sistemas multi-store utilizam diferentes abordagens para permitir a consulta de múltiplas fontes de dados. Nesta dissertação pretende-se que o sistema multi-store utilizado permita consultar múltiplas fontes de dados heterogéneas. Deste modo, e para facilitar a adoção deste sistema este deve suportar uma linguagem única para acessar os diferentes dados. Esta linguagem deve ser popular e fácil de usar. Por outro lado, o suporte a novas fontes de dados deve ser fácil de ser implementado. Além disso, os conectores à fonte de dados devem suportar push-down das interrogações de modo a consultar à fonte apenas os dados necessários. Por fim, o sistema multi-store deve também conseguir otimizar as interrogações de modo a tirar máximo proveito das capacidades da nuvem para reduzir o tempo de execução. Todos os sistemas multi-store enumerados suportam interrogações utilizando a linguagem SQL permitindo que seja facilmente adotável. No entanto, alguns sistemas permitem maior facilidade de adoção. Alguns sistemas têm suporte a uma grande variedade de fontes de dados, como é o caso do Trino, PostgreSQL FDW e do Apache Drill. Outros permitem facilmente implementar o suporte a novas fontes de dados como é o caso do CloudMdsQL e Delta Lake. Todos os sistemas mencionados permitem, com ou menor dificuldade, o push-down das interrogações. Do mesmo modo todos eles realizam otimizações que permitem reduzir o tempo de execução tendo em conta os dados da fonte. 2.3 Discussão Nesta dissertação pretende-se que seja possível tirar proveito do push-down e da materialização dos dados recolhidos de fontes heterogéneas para reduzir o volume de dados transferidos. Do mesmo modo, pretende-se que o sistema consiga responder a interrogações analíticas exploratórias em tempo real onde não é possível conhecer previamente quais os dados que serão necessários. Neste sentido, o sistema desenvolvido deve conseguir transferir apenas os dados necessários e evitar o reenvio destes. Tendo em conta os requisitos desta dissertação, a réplica deve ser atualizada a cada interrogação e deve transferir apenas os dados necessários. Deste modo, a réplica deve ser atualizada sempre que existe uma consulta à fonte de dados. Esta é a única política que garante que os dados se encontram atualizados e que apenas consulta os dados necessários. 23 CAPÍTULO 2. TRABALHO RELACIONADO Por outro lado, os métodos de sincronização têm diferentes limitações. O método de sincronização deve reduzir o volume de dados transferido, o que não é o caso da sincronização completa. Por se tratar de fontes de dados heterogénea, o método não deve ser dependente da implementação interna, como é o caso do método de replicação incremental utilizando log. Por fim, o método deve ser eficiente a detetar as novas diferenças. O método de replicação incremental, ao contrário dos outros métodos, permite detetar os novos tuplos sem a necessidade de percorrer todos os dados no caso de a chave estar indexada. Além disso, este método permite que se tire proveito da capacidade de push-down de filtros para efetuar a sincronização. Por fim, os sistemas multi-store enumerados permitem realizar eficientemente interrogações a fontes de dados heterogéneas e tirar partido das capacidades push-down das fontes. Todos eles suportam a linguagem SQL e consequentemente permitem fácil adoção. Neste sentido, o sistema deve suportar uma grande variedade de fontes de dados e deve permitir ser facilmente estendido. Por estes motivos, o sistema PostgreSQL FDW oferece maiores garantias devido à sua maturidade, pela lista de conectores disponíveis e pela facilidade de uso. 24 Capítulo 3 Design e implementação O objetivo desta dissertação consiste em desenvolver um mecanismo capaz de responder a interrogações analíticas exploratórias na nuvem sobre dados heterogéneos que se encontram na periferia. Além disso, pretende-se que este mecanismo reduza o tempo de execução. Para isso, o mecanismo utiliza uma cache e um sistema de sincronização para reduzir o volume de dados transferidos. Neste sentido, construiuse uma arquitetura capaz de aplicar filtros push-down na fonte de dados e guardar os tuplos recebidos da periferia de modo a evitar a sua retransmissão. A combinação destas duas funcionalidades permite transferir apenas os dados necessários e, idealmente, transferir apenas uma vez cada tuplo. O sistema desenvolvido é composto por duas partes. A Seção 3.1 descreve o funcionamento do sistema de armazenamento responsável por responder às interrogações analíticas sobre os dados contidos em dispositivos periféricos. Este mecanismo contém uma cache que armazena os tuplos previamente transferidos para evitar retransmissões. Na Seção 3.3 descreve-se o algoritmo de sincronização implementado. Este não só evita a retransmissão de tuplos como também propaga os filtros para a fonte de dados de modo a transferir apenas os dados necessários. Nesta secção, descreve-se também o mecanismo adaptativo que permite balancear o volume de dados transferido e o tempo de execução na fonte tendo em conta a carga de trabalho, a capacidade da rede e do dispositivo periférico. 3.1 Arquitetura O sistema desenvolvido é composto por dois módulos, como apresentado na Figura 4. O módulo de sincronização é responsável por solicitar da fonte de dados apenas os dados que não se encontram na nuvem. 25 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO nova consulta, os tuplos retornados não devem diferir da primeira leitura. Do mesmo modo, se durante a transação for retornado um conjunto de tuplos, em interrogações posteriores, o conjunto deve ser o mesmo independentemente se outras transações tenham inserido ou removido tuplos. Se o conector suportar este nível de isolamento, consegue-se garantir que, por um lado, os dados presentes na cache são coerentes com a fonte de dados e, por outro lado, várias leituras à mesma tabela não retornam resultados diferentes. 3.2.5 Outras funcionalidades O módulo da cache suporta realizar transformações aos dados provenientes da fonte antes de inserir na cache. Geralmente no processo de copiar os dados dos dispositivos periféricos para a nuvem é necessário realizar um processo de transformação dos dados. Por exemplo, pode ser necessário projetar ou renomear as colunas. Em certos casos pode ser necessário ter colunas com valores derivados ou codificar valores. Do mesmo modo, permite transformar os dados num formato mais próprio para consulta. Também é possível realizar agregações ou desduplicação dos dados. Antes de os dados serem inseridos na cache, é possível fornecer uma interrogação SQL que fará a transformação dos dados. Os dados provenientes da fonte são armazenados numa tabela temporária para serem posteriormente inseridos na cache dos dados. Por omissão, os dados não sofrem nenhuma modificação. Se uma interrogação SQL for fornecida, essa interrogação é utilizada para fazer a transformação dos dados antes de estes serem inseridos na cache. Esta funcionalidade permite não só materializar as transformações para evitar reprocessamento, como também manter este sempre atualizado. Além disso, o sistema suporta a importação de esquema. Esta funcionalidade permite facilmente criar uma tabela para cada um dos conectores. Este utiliza informação do PostgreSQL para determinar dados relevantes sobre as colunas da tabela como o nome e o tipo. Esta funcionalidade evita ter de criar um esquema manualmente para cada Foreign Data Wrapper Por fim, além de ser necessário ter uma cache para tirar proveito da localidade, em certas situações pode ser vantajoso poder modificar também os dados presentes na fonte de dados remota. Nestes casos, o nosso sistema serve de intermediário entre o motor de consulta e a base de dados remota. O sistema propaga os comandos de inserção, remoção e atualização para os dispositivos remotos. O módulo de sincronização ao receber estes comandos recebe informação de quais os tuplos foram inseridos, atualizados ou removidos e propaga essa informação para a fonte de dados. Se o módulo de sincronização conseguir atualizar a informação sobre o estado da cache ao receber informação sobre inserção, atualização ou remoção, então pode-se atualizar a cache. Caso contrário, as inserções e atualizações podem ser lidas posteriormente quando forem consultadas. Atualizar a cache evita que os dados inseridos através do motor de consulta na nuvem sejam retransmitidos. 32 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO 3.2.6 Limitações A extensão Multicorn permite reduzir a complexidade da implementação de um Foreign Data Wrapper para o motor de consulta PostgreSQL. No entanto, esta extensão sofre de algumas limitações que limitam o escopo da implementação. Primeiramente, o Multicorn apenas suporta a versão 12 do PostgreSQL. Esta limitação faz com que não seja possível tirar proveito das últimas atualizações. Entre estas funcionalidades está a possibilidade de execução assíncrona das interrogações. O Multicorn cria uma instância por cada processo e por cada tabela. Por este motivo, não é possível reaproveitar as conexões para múltiplas tabelas. Este pode ser um fator limitador principalmente quando o servidor suporta um número limitado de conexões e utilizadas múltiplas tabelas. Por outro lado, não existe suporte a execução paralela de interrogações. O suporte de execução paralela permite executar vários planos em simultâneo e assim reduziro tempo de execução em vezde executar sequencialmente cada plano. Apesar desta limitação, esta funcionalidade é facilmente implementável na própria extensão. Do mesmo modo, o Multicorn não suporta disjunções. Todavia, este suporta conjunções de filtros e por este motivo a implementação foca-se apenas em conjunções de predicados. De modo a superar esta limitação a interrogação pode ser reescrita para realizar a união de várias consultas. Por fim, o Multicorn não suporta junções, funções nem agregações. Por este motivo, não é possível propagá-las para a fonte de dados nem na consulta da cache. Deste modo, o middleware desenvolvido apenas retorna os tuplos em cache. 3.3 Algoritmo de sincronização O algoritmo de sincronização proposto oferece três grandes vantagens. Em primeiro lugar, permite realizar a sincronização incrementalmente, isto é, evita a retransmissão de tuplos previamente transferidos. Por outro lado, suporta push-down de seleções e projeções de modo a sincronizar apenas os dados necessários. O suporte de seleções permite tirar proveito de índices ou de outras otimizações que reduzem o tempo de execução. Além disso, o algoritmo funciona para todas as fontes de dados com suporte de seleções e projeções sem necessitar de conhecer a sua implementação interna. Além disso, o algoritmo apresentado balanceia conceitos como, por um lado, apenas transferir os dados necessários para responder à interrogação e, por outro, remover alguns filtros de modo a reduzir o tempo de execução no dispositivo periférico. ASeção 3.3.1 descreve a relevância do suporte de seleções como modo para reduzir o volume de dados transferidos. A Seção 3.3.1.1 apresenta um método que tira partido dos filtros previamente realizados e das marcas temporais para minimizar a transferência de dados ao mínimo necessário. A Seção 3.3.1.2 33 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO exibe uma abordagem adaptativa para reduzir a complexidade das interrogações. Por fim, a Seção 3.3.2 descreve abordagens complementares que permitem o suporte do push-down de projeções. 3.3.1 Suporte a seleções O algoritmo proposto suporta o push-down de seleções na fonte de dados. O suporte de seleções permite transferir apenas os tuplos necessários e, por consequência, reduzir o volume de dados. Dependendo das interrogações que serão realizadas ao sistema, alguns dados podem nunca ser necessários. Assim, o algoritmo de sincronização deve preferencialmente transferir apenas os tuplos que satisfazem as condições da interrogação e que não foram previamente transferidos. Neste sentido, desenvolveram-se técnicas que permitissem consultar apenas os dados necessários. Deste modo, pode-se aproveitar as capacidades de filtragem das fontes de dados para obter apenas os dados necessários. Neste contexto, o sistema deve propagar os filtros das interrogações para a fonte de dados quando vantajoso. Tendo por princípio a replicação incremental baseada na chave, pode-se estender o seu funcionamento para suportar a propagação de filtros. Esta abordagem permite tirar proveito das vantagens inerentes a este algoritmo. Além disso, consegue-se evitar o envio de dados que não são posteriormente lidos. 3.3.1.1 Minimizar a transferência de dados De modo a minimizar a transferência de dados da fonte, o sistema guarda os filtros aplicados bem como a maior marca temporal recebida para esse filtro. Quando o sistema questiona a fonte de dados, este aplica filtros que excluem os dados já presentes na nuvem. Deste modo, reduz-se ao mínimo o volume de dados enviados da fonte para a nuvem. A título de exemplo, considere uma tabela denominada Demo. Esta tabela é constituída pela coluna Id (chave primária), duas colunas, AeBe a coluna Ts, que indica quando a linha foi inserida. AFigura 6 representa os diferentes estados da cache na nuvem após a realização de sucessivas interrogações. Inicialmente a cache está vazia como representado na Figura 6a. Quando o servidor realiza a interrogação SELECT * FROM Demo WHERE A = 1. Como não existe nenhuma linha em cache, a interrogação é propagada à fonte de dados como apresentada. Neste caso, o dispositivo enviará as linhas com 𝐼𝑑 ∈ {3,4}como representado na Figura 6b. Além disso, o sistema guarda, também, a maior marca temporal recebida (4) e associa-se ao predicado realizado (𝐴=1). Deste modo, é guardado o predicado 𝐴=1∧𝑇𝑠 ≤4. A negação deste predicado é utilizado em interrogações futuras. Depois, o servidor realiza a interrogação SELECT * FROM Demo WHERE B = 1. Ao observar a Figura 6b é possível observar que as linhas cujo 𝐼𝑑 ∈ {3,4}já foram envidas. Por esse motivo não é necessário reenviá-los. Para evitar enviar estas linhas é aplicado um filtro extra que combina o predicado atual 𝑃com a negação dos filtros previamente realizados 𝑄usando a expressão 𝑃∧ ¬𝑄. Neste caso, 34 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO Id Ts A B 1 1 0 0 2 2 0 1 3 3 1 0 4 4 1 1 (a) Estado inicial da tabela Demo. A cache encontra-se inicialmente vazia. Id Ts A B 1 1 0 0 2 2 0 1 3 3 1 0 4 4 1 1 (b) Estado da tabela Demo depois da 1.ª interrogação:’SELECT * FROM Demo WHERE A = 1’. Id Ts A B 1 1 0 0 2 2 0 1 3 3 1 0 4 4 1 1 (c) Estado da tabela Demo depois da 2.ª interrogação:’SELECT * FROM Demo WHERE B = 1’. Id Ts A B 1 1 0 0 2 2 0 1 3 3 1 0 4 4 1 1 5 5 1 1 (d) Estado da tabela Demo depois de ser inserido a linha ⟨5,5,1,1⟩e da 3.ª interrogação:’SELECT * FROM Demo WHERE B = 1’. Figura 6: Exemplo do estado da tabela Demo após sucessivas interrogações. As novas linhas adicionadas à cache são representadas com fundo amarelo. As linhas representadas com fundo cinzento já se encontram na cache. Por fim, as linhas com fundo branco não se encontram ainda em cache. ainterrogação realizada no dispositivo na borda é SELECT * FROM Demo WHERE B = 1 AND NOT (A = 1 AND Ts <= 4). Executando a interrogação, é retornada a linha cujo 𝐼𝑑 =2como demonstra a Figura 6c. Mais uma vez, nós guardamos o filtro aplicado (𝐵=1) em conjunto com a maior marca temporal transferida (2). Finalmente, o resultado do predicado de exclusão é (𝐴=1∧𝑇𝑠 ≤4)∨(𝐵=1∧𝑇𝑠 ≤2) Por fim, o sistema realiza a interrogação SELECT * FROM Demo WHERE B = 1. Entretanto, no dispositivo remoto foi inserido uma nova linha ⟨5,5, 𝐹, 1,1⟩na tabela Demo. Uma vez mais, a interrogação é transformada para excluir as linhas que foram previamente requisitadas. A interrogação passada é equivalente a SELECT * FROM Demo WHERE B = 1 AND NOT ((A = 1 AND Ts <= 4) OR (B = 1 AND Ts <= 2)). Desta vez, o sistema retorna a linha cujo 𝐼𝑑 =5. Apesar de a linha com 𝐼𝑑 =4satisfazer tanto os predicados 𝐴=1e𝐵=1, esta não é transferida porque tem marca temporal inferior ou igual 4. 3.3.1.2 Simplificação dos filtros O número de filtros cresce com o número de interrogações realizadas pelo servidor. Apesar de os otimizadores conseguirem remover filtros redundantes, existe na mesma um problema de armazenar um número crescente de filtros. Além disso, a complexidade de transferir e interpretar uma interrogação aumenta proporcionalmente com o número de filtros. Para mitigar estes problemas, o algoritmo desenvolvido apresenta algumas estratégias para simplificar os filtros realizados. Utilizando um exemplo prático, considere novamente a tabela Demo, cujo esquema está representado na Figura 6 mas com dados diferentes. Quando a interrogação SELECT * FROM Demo WHERE A = 1 AND B = 1 é executada e determina-se que a maior marca temporal recebida foi de 1, o sistema guarda a condição 𝐴=1∧𝐵=1∧𝑇𝑠 ≤1. Na Figura 7a é possível observar como são guardados os predicados. Quando é executado a segunda interrogação SELECT * FROM Demo WHERE A = 1 AND C = 1, a 35 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO Ts≤1A=1{} B=1 (a) Estado após a 1.ª interrogação:SELECT * FROM Demo WHERE A = 1 AND B = 1 Ts≤1 A=1{} B=1 Ts≤3C=1 (b) Estado após a 2.ª interrogação:SELECT * FROM Demo WHERE A = 1 AND B = 1 Ts≤6A=1{} (c) Estado após a 3.ª interrogação:SELECT * FROM Demo WHERE A = 1 Ts≤6{} (d) Estado após a 4.ª interrogação:SELECT * FROM Demo Figura 7: Evolução do estado dos predicados após várias interrogações à tabela Demo. maior marca temporal recebida é 3. Neste caso, o predicado 𝐴=1∧𝐶=1∧𝑇𝑠 ≤3é guardado. Contudo, armazenar 𝐴=1∧𝐵=1∧𝑇𝑠 ≤1e𝐴=1∧𝐶=1∧𝑇𝑠 ≤3não é ótimo, visto que o predicado 𝐴=1é guardado duas vezes. Em alternativa, uma representação mais eficiente consiste em armazenar 𝐴=1∧ ((𝐵=1∧𝑇𝑠 ≤1)∨(𝐶=1∧𝑇𝑠 ≤3)), como mostra a Figura 7b. Portanto, o primeiro passo para reduzir o número de filtros consiste em remover os filtros duplicados. Considere que foi executado a terceira interrogação SELECT * FROM Demo WHERE A = 1 e a maior marca temporal foi 6. Como o predicado 𝐴=1∧𝑇𝑠 ≤6contém todos os dados previamente retornados pelos dois predicados anteriores, é possível simplificar os filtros e armazenar este predicado como representado na Figura 7c. Do mesmo modo, se executarmos a interrogação SELECT * FROM Demo e não for retornado nenhuma linha então pode-se remover todos os outros filtros como mostra a Figura 7d. Assim sendo, a segunda otimização consiste em combinar os filtros quando estes são subconjuntos de outros filtros. A partir destas duas técnicas é possível reduzir a complexidade da interrogação, reduzir o volume de dados armazenado, os dados enviados na rede e o custo de analisar a interrogação. 3.3.1.3 Limpeza de filtros As otimizações propostas na Seção 3.3.1.2 são úteis para remover filtros redundantes. No entanto, a complexidade da interrogação aumenta com cada interrogação. Se a interrogação se tornar demasiado complexa, têm-se o caso onde o tempo de execução não compensa as vantagens obtidas pela redução do volume de dados. Por este motivo, torna-se importante considerar a hipótese de se remover algumas das condições guardadas. Apesar de a remoção dos filtros levar a um acréscimo no volume de dados transferidos, a redução do tempo de execução pode ser favorável. 36 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO Deste modo, o algoritmo balanceia os custos de filtrar os dados nos dispositivos remotos e o custo de transferir os dados. Consequentemente, para fazer esta comparação é necessário definir alguns parâmetros. Estes parâmetros definem os custos relativos de cada operação. Em primeiro lugar, o custo de aplicar um filtro pode ser definido como 𝑐𝑓=𝑐𝑐×𝑐×𝑟onde 𝑐𝑓 é o custo total de realizar a filtragem, 𝑐𝑐denota o custo por cada condição, 𝑐representa o número de condições e 𝑟é o número de estimado linhas na fonte de dados. Esta fórmula é válida para o pior caso onde é feita uma verificação completa da tabela. Também se assume nesta fórmula que todas as condições têm o mesmo custo e que o seu custo é independente das outras condições. Neste caso o custo de filtrar a tabela aumenta com o número de filtros e o número de linhas. Por exemplo, considere que previamente fora realizada a seguinte interrogação:SELECT * FROM Demo WHERE A = 1. Retornou-se 50 tuplos e a maior marca temporal lida foi de 100. Neste caso, a diferença entre a tabela na fonte de dados e a tabela presente na cache pode ser consultada com a seguinte interrogação SELECT * FROM Demo WHERE NOT (A = 1 AND TS <= 100) ou SELECT * FROM Demo WHERE A <> 1 AND TS > 100. Considere que na tabela Demo existem 100 linhas e o custo de cada condição é de 100 unidades. Neste caso, são aplicadas 2 condições (𝐴≠1e𝑇𝑆 >100) à tabela. No pior caso, estas condições necessitam de ser testadas para todos os tuplos da tabela. Assim, o custo de se aplicar o filtro (𝑐𝑓) é igual a 100 ×2×100 =20000 unidades. Por outro lado, o custo de transferir os dados pode ser definido como 𝑐𝑡=𝑐𝑏×𝑟×𝑤onde 𝑐𝑡 representa o custo total relativo à transferência dos dados, 𝑐𝑏representa o custo por cada byte transferido, 𝑟expressa o número estimado de linhas na fonte e 𝑤é o tamanho médio estimado de cada linha em bytes. Esta fórmula não considera o custo adicional decorrente do protocolo de transmissão nem da serialização dos dados. Além disso, considera que não existe custo adicional por cada linha. Neste caso, o custo da transferência dos dados aumenta com o número de linhas e o tamanho médio de cada linha. Utilizando o mesmo exemplo, considere que sabendo que cada byte custa 1 unidade, cada tuplo ocupa em média 100 bytes. Também se sabe que 50 tuplos já se encontram na cache e que por isso apenas 50 tuplos ainda não foram transferidos. Assim, o custo da transferência (𝑐𝑡) é igual a 1×100 ×50 =5000 unidades. No caso de o custo de filtrar os dados ser superior ao necessário para transferir os dados, então podese considerar que existem filtros supérfluos. Este é geralmente o caso quando o número de linhas que satisfazem a condição é bastante reduzido. Por este motivo, filtros que removem poucas linhas são fortes candidatos a poderem ser removidos. Seguindo os exemplos previamente usados, o custo de transferir os dados é de 5000 unidades, enquanto o custo de filtrar os dados é de 20000 unidades. Neste caso é possível observar que o custo de filtragem é superior ao custo dos dados transferidos. Sendo assim, deve-se considerar, para este exemplo, a hipótese de serem removidos filtros. Uma possível abordagem para remover os filtros consiste em avaliar os custos de cada filtro. O custo de cada filtro pode ser avaliado tendo em conta o número de condições e o número de bytes que podem 37 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO ser retransmitidos se ocorrer a sua remoção. Neste sentido, procura-se determinar quantos bytes dos presentes na cache devem-se ao filtro. Para estimar o número de bytes retornados por um filtro é necessário estimar o número de bytes que satisfazem o filtro. Estimar o número de bytes da cache que respondem ao filtro permite perceber o limite máximo de bytes que necessitam de ser retransmitidos no caso de este ser removido. O número de máximo de bytes retransmitidos no caso da sua remoção pode ser expresso como 𝑏𝑟=𝑟𝑟×𝑤onde 𝑏𝑟representa o número máximo estimado de bytes que serão retransmitidos se este filtro for removido, 𝑟𝑟denota o número estimado de linhas que satisfazem o filtro na cache e 𝑤 representa o tamanho médio em bytes de cada linha. Seguindo o exemplo anterior, pretende-se saber quantos bytes da cache respeitam a condição 𝐴=1∧𝑇𝑆 ≤100. Neste caso se sabe que esta condição é verdade (𝑟𝑟) para 50 tuplos da cache e cada tuplo ocupa em média (𝑤) 100 bytes. Sendo assim, o número de bytes esperado que o filtro ocupa na cache (𝑏𝑟) é de 50 ×100 =5000 bytes. A remoção do filtro 𝐴=1∧𝑇𝑆 ≤100 pode no máximo obrigar à transferência de 5000 bytes numa futura interrogação. Após estimar todos os filtros, testa-se o filtro com menor potencial de número de bytes retransmitidos. Neste caso, procura-se pelo filtro com menor número de bytes esperado (𝑏𝑟). Depois, testa-se a hipótese de a remoção deste filtro reduzir o custo total do sistema. Este é o caso se 𝑐𝑟×𝑐𝑐×𝑟<𝑐𝑏×𝑟𝑓×𝑤 onde 𝑐𝑟é o número de condições do filtro que pode ser removido, 𝑐𝑐representa o custo por condição, 𝑟 é o número estimado de linhas na tabela da fonte de dados, 𝑐𝑏denota o custo por cada byte transferido, 𝑟𝑓representa o número marginal estimado de linhas filtradas pelo filtro e 𝑤simboliza o tamanho médio em bytes de cada linha. O número marginal estimado das linhas filtradas pelo filtro representa o número de linhas que necessitarão de ser retransmitidas no caso deste filtro ser removido. Se esta condição for verdadeira, então a remoção do filtro é vantajosa e o filtro é removido. Posteriormente, é realizado o mesmo procedimento para o próximo filtro candidato. Em oposição, o processo termina se não houver mais filtros ou se a remoção do filtro aumentar o custo total. Utilizando o exemplo anterior, seria testado se a remoção do filtro 𝐴=1∧𝑇𝑆 ≤100 seria vantajoso. Sendo este filtro composto por duas condições (𝐴=1e𝑇𝑆 ≤100), o custo por condição (𝑐𝑐) ser igual a 100 unidades e o número de tuplos na fonte de dados (𝑟) ser 100 unidades, então, neste caso, o custo de realizar a filtragem só com este filtro é de 2×100 ×100 =20000 unidades. Do mesmo modo, visto que apenas existe um filtro guardado, então o número de linhas da cache que respondem ao filtro (𝑟𝑟) são 50 e cada linha ocupa em média (𝑤) 100 bytes. Sendo o custo por byte transferido (𝑐𝑏) de 1 unidade, o custo da retransmissão caso este filtro seja removido será de 1×50 ×100 =5000 unidades. Sendo assim, como o custo de realizar a filtragem do filtro (20 000) é superior ao custo da retransmissão (5 000), então o filtro pode ser removido. Como não existe mais nenhum filtro, o processo termina. Uma das desvantagens deste processo é o custo computacional. Por exemplo, estimar o número de linhas é um processo computacionalmente intensivo. Por este motivo é imperativo apenas executar a filtragem se o sistema está confiante que consegue melhorar o desempenho. Neste caso, a remoção de filtros é executada apenas se 𝑐𝑓>𝑐𝑡+𝑐𝑒×𝑓onde 𝑐𝑓é o custo de filtrar toda a tabela, 𝑐𝑡denota o 38 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO custo estimado de transferir os dados em falta e 𝑐𝑒representa o custo de estimar o número de linhas e 𝑓representa o número de filtros que ainda não foram estimados. Se esta condição for verdadeira, então deve-se ponderar a remoção de filtros e, por consequência, estimar o número de bytes para cada filtro ainda não estimado. Caso contrário, nenhum filtro é removido. 3.3.2 Suporte a projeções Em complemento à abordagem anterior, pode ser necessário suportar projeçõesde modo a evitar transferir colunas que nunca serão lidas. O suporte de projeções oferece maior granularidade na sincronização dos dados. O algoritmo descrito previamente pode ser estendido para suportar projeções. Para suportar projeções, é necessário guardar o registo dos predicados previamente realizados para cada coluna. Com isto, as seleções previamente realizadas são independentes de cada coluna. Ou seja, apenas as colunas transferidas da fonte de dados atualizam a lista dos predicados previamente realizados. Esta solução tem algumas restrições nomeadamente a necessidade de transferir sempre o identificador da linha, a marca temporal e as colunas envolvidas nos predicados. O identificador é essencial para ser possível juntar os novos dados com os dados presentes na cache. A marca temporal é necessária para determinar o maior valor da chave lido. Por fim, as colunas utilizadas nos predicados necessitam de ser também transferidos para ser possível verificar as condições na cache. Um dos problemas desta abordagem é que diferentes colunas podem conter predicados diferentes. Neste sentido, no mesmo tuplo, algumas colunas podem já ter sido previamente transferidos e noutros casos não. Assim sendo, existem diversas estratégias que podem ser concretizadas para suportar projeções. A primeira estratégia consiste em utilizar o predicado menos restritivo de todas as colunas e não realizar nenhuma transformação. Neste caso, de todos os predicados utilizados deve-se aplicar à fonte de dados os filtros que são subconjuntos de outros filtros. Apesar de esta abordagem poder ser utilizada com qualquer fonte de dados, não evita a retransmissão das colunas previamente transferidas. Outra abordagem possível consiste em retornar nulo para os valores que já se encontram na cache. Isto quer dizer que para cada coluna é feito uma interrogação com o respetivo predicado. À semelhança da estratégia anterior é aplicado à fonte de dados uma seleção com o predicado menos restritivo. Todavia, as colunas cujo predicado seja mais restritivo aplicam uma transformação utilizando uma expressão condicional. Neste caso, se o valor já estiver presente na cache, então este pode ser substituído por um valor nulo. No módulo de sincronização, os valores nulos não atualizam a cache. Os valores nulos podem ser utilizados para reduzir o volume de dados transferidos. No entanto, esta abordagem obriga a que a fonte de dados suporte expressões condicionais. Por outro lado, a complexidade da interrogação cresce se todas as colunas tiverem associadas a predicados diferentes. Como prova de conceito, esta abordagem foi implementada. 39 CAPÍTULO 3. DESIGN E IMPLEMENTAÇÃO A terceira abordagem consiste em juntar várias seleções. Neste caso, para cada coluna que não seja o identificador, é realizada uma consulta à fonte de dados com os seus respetivos filtros. Em cada consulta são projetados apenas o identificador do tuplo e a respetiva coluna. Deste modo, consegue-se aplicar as seleções a cada uma das colunas. As colunas com o mesmo predicado podem ser projetadas na mesma interrogação. Por fim, as colunas são combinadas para responder à interrogação. Esta solução tem a vantagem de apenas necessitar que a fonte suporte projeções e seleções visto que a junção pode ser realizada na nuvem. Contudo, esta operação não é eficiente e por esse motivo pode ter grande impacto no tempo de execução da interrogação. 3.3.3 Implementação O algoritmo de sincronização proposto foi implementado em Python e utiliza a biblioteca psycopg2 para comunicar com a fonte de dados remota. O módulo de sincronização recebe as projeções e as seleções provenientes do módulo da cache. Estes são depois utilizados para responder à interrogação. Para comunicar com a fonte de dados, desenvolveu-se um componente responsável por transformar os filtros guardados numa interrogação SQL. Esta interrogação era posteriormente utilizada para consultar os dados remotos. Após retornar os tuplos da execução da interrogação para o módulo da cache, os filtros da interrogação são utilizados para atualizar o registo de filtros previamente executados. O módulo de sincronização guarda os filtros e as marcas temporais numa estrutura SetTrie [35]. Esta estrutura permite reduzir o tempo de execução para a procura de subconjuntos e superconjuntos de filtros. Estas operações são necessárias para realizar a simplificação dos filtros. Do mesmo modo, esta estrutura evita duplicação de filtros e por esse motivo é mais compacta. 40 Capítulo 4 Avaliação Este capítulo apresenta os resultados e a análise das medições de desempenho realizadas. Primeiramente, na Seção 4.1 comparou-se o volume de dados transferidos, a memória utilizada e o tempo de execução das interrogações de dois algoritmos de sincronização. A Seção 4.2 apresenta os resultados do algoritmo proposto de sincronização quando comparado com o algoritmo de sincronização incremental baseado na chave. A Seção 4.3 apresenta uma adaptação do CH-benCHmark tornando-o orientado ao evento, adequando-o a uma carga de trabalho típica dos dispositivos periféricos. Na Seção 4.4 empregase esse benchmark para avaliar os benefícios do uso da cache para reduzir o tempo de execução das interrogações. 4.1 Avaliação dos métodos de replicação Na Seção 2.1.2 foram enumerados diferentes métodos de sincronização. Destes métodos, a replicação incremental utilizando a estrutura IBLT e a baseada em chave destacaram-se pela baixa complexidade computacional e não necessitarem de conhecer a implementação interna. Neste sentido, conduziram-se experiências sobre os dois métodos com o intuito de medir o volume de dados enviado, a memória utilizada e, por fim, o tempo de execução de cada um deles. Estas métricas permitem comparar o desempenho dos dois métodos e determinar a base para o algoritmo desenvolvido. A estrutura IBLT é uma estrutura probabilística onde existe probabilidade de haver colisão e que, por isso, pode ocorrer a retransmissão de tuplos. Para avaliar o seu desempenho, realizaram-se medições sobre uma estrutura com fator de sobrecarga de 𝛼= 1,2 e 𝛼= 1,3. Por exemplo, se a sincronização tiver 100 diferenças e o fator de sobrecarga é 𝛼= 1,5, então a estrutura é composta por 150 baldes. O aumento do número de baldes reduz a probabilidade de ocorrer colisões, mas ocupa mais espaço em 41 CAPÍTULO 4. AVALIAÇÃO 4.2.1 Microbenchmark Para avaliar o desempenho do algoritmo quando comparado com as alternativas foi realizado um microbenchmark iterativo implementado em Python. Em cada iteração, o servidor na nuvem interroga o dispositivo remoto e adiciona 100 novas linhas. Depois, executa o comando ANALYZE para atualizar as estatísticas usadas pelo planeador [20]. O microbenchmark considera apenas as operações relacionadas com a replicação dos dados. Por exemplo, o tempo necessário para ler os dados dos dispositivos e transferir através da rede. No total, o microbenchmark executa 250 iterações. O esquema da base de dados está representada na Tabela 4. Inicialmente a relação encontra-se vazia. Omicrobenchmark interroga os dados aplicando filtros com uma determinada seletividade. Por exemplo, se forem aplicados dois filtros com 50% de seletividade então um possível expressão é (e.g., WHERE q0 < 0.5 AND q1 < 0.5). Este microbenchmark simula um trabalho analítico ad hoc. Por este motivo, obenchmark gera filtros estocásticos para se assemelhar ao trabalho esperado. No benchmark foram definidas duas cargas de trabalho: •Simples — executa interrogações de leitura com, em média, 1 filtro. •Complexo — executa interrogações de leitura com, em média, 10 filtros. 4.2.2 Ambiente de testes Os testes realizados foram executados num único CPU quad-core de 1,6 GHz, com 8 GB de memória RAM disponível e armazenamento SSD. Os testes foram realizados no sistema operativo Ubuntu 20.04 LTS. Para simular uma tipologia típica de rede, foi adicionado artificialmente, utilizando ferramentas ao nível do sistema operativo, uma latência de 50 ms e limitou-se a largura de banda entre os servidores na nuvem e na borda a 50 MB/s. O benchmark foi implementado com a linguagem Python3.7. Utilizou-se a base de dados PostgreSQL para armazenar os dados e a biblioteca psycopg2 para fazer a ligação à base de dados. 4.2.3 Resultados Os testes realizados permitem perceber o tempo de execução e do volume de dados transferido para a nuvem dos diferentes algoritmos. O principal objetivo dos algoritmos de sincronização consiste em reduzir o tempo de execução. Neste sentido, mediu-se o tempo de execução de cada um dos métodos, tanto numa carga de trabalho simples como complexa, para perceber o seu desempenho. Do mesmo modo, devido às limitações da capacidade da rede, pretende-se comparar a redução do volume de dados dos diferentes algoritmos. Algoritmos com menor tráfego têm menor impacto na rede. 48 CAPÍTULO 4. AVALIAÇÃO Chave Sem Limpeza Com Limpeza Adaptativo Interrogações simples 0.0 2.5 5.0 7.5 10.0 12.5 15.0 17.5 20.0 Tempo de execução (s) Chave Sem Limpeza Com Limpeza Adaptativo Interrogações complexas Figura 11: Tempo de execução baseado na carga de trabalho e o algoritmo de sincronização. 4.2.3.1 Tempo de execução AFigura 11 revela o tempo de execução tanto das interrogações simples e complexas para os diferentes algoritmos. Esta métrica permite perceber qual o desempenho dos algoritmos quando a complexidade das interrogações muda. Os resultados do tempo de execução da carga de trabalho simples, presentes na Figura 11, mostram, como era esperado, que o tempo de execução de aplicar a limpeza de filtros tem piores resultados que não aplicar a limpeza. O algoritmo adaptativo proposto (Adaptativo) realizou as 250 interrogações simples à fonte de dados em 17,5 s. Em comparação com o método de replicação baseado na chave (Chave) que demorou 19,1 s para completar todas as interrogações. Assim sendo, o algoritmo adaptativo reduziu o tempo de execução em cerca de 8%. Contudo, Sem limpeza demorou 17,3 s para completar todas as interrogações e Limpeza demorou 18,2 s. Neste sentido, a realizar sempre a limpeza aumentou o tempo de execução em cerca de 5%. O Adaptativo demorou cerca de 4% menos para completar as mesmas interrogações do que o Limpeza Em contrapartida, para interrogações complexas obteve-se um resultado diferente. O tempo de execução do método baseado na chave (Chave) é o mesmo que o simples visto que os filtros não são propagados. Em oposição aos resultados obtidos para as interrogações simples, não realizar a limpeza dos filtros teve piores resultados que o método baseado na chave (Chave), demorando 20,1 s para completar o teste. Este resultado mostrou um aumento de cerca de 5% em relação ao algoritmo onde não é propagado os filtros (Chave). Em oposição, o algoritmo de Limpeza eAdaptativo tiveram os melhores resultados: 14,6 s e 14,4 s respetivamente. O algoritmo adaptativo proposto conseguiu assim uma redução de cerca de 49 CAPÍTULO 4. AVALIAÇÃO Chave Sem Limpeza Com Limpeza Adaptativo Interrogações simples 0 10 20 30 40 50 Tráfego (MB) Chave Sem Limpeza Com Limpeza Adaptativo Interrogações complexas Figura 12: Volume de dados transferidos pela rede dependendo da carga de trabalho e o algoritmo de sincronização. 25% no tempo de execução em relação ao tempo do método baseado na chave e 28% em comparação ao algoritmo onde nunca é realizado a limpeza dos filtros. Em suma, o algoritmo adaptativo proposto consegue obter tempos de execução muito próximos do ótimo, evitando o custo da limpeza dos filtros quando as interrogações são simples e aplicando a limpeza destes em interrogações mais complexas onde o custo de processamento no dispositivo remoto é mais expressivo. 4.2.3.2 Tráfego AFigura 12 representa o volume de dados transferidos pela rede da fonte de dados para a nuvem. Esta métrica revela o número de tuplos transferidos dos diferentes algoritmos e para diferentes complexidades das interrogações. O método de replicação incremental baseado na chave transferiu 47,7 MB de dados. Como o algoritmo de replicação incremental não propaga os filtros para a fonte de dados, o volume de dados transferido é equivalente ao volume presente na fonte. Os resultados obtidos demonstram que o algoritmo proposto consegue efetivamente reduzir o volume de dados transferido. Para a carga de trabalho onde foram realizadas interrogações mais simples, em comparação com o método de replicação incremental baseado na chave (Chave), o algoritmo adaptativo reduziu 19% o volume de dados transferido. O Adaptativo transferiu a mesma quantidade de dados (38,4 MB) que o Sem limpeza e menos 282 kB que o Limpeza. Esta diferença deve-se aos filtros realizados sobre a fonte de dados. Nestes algoritmos, apenas os dados necessários para responder à interrogação 50 CAPÍTULO 4. AVALIAÇÃO são transmitidos. Deste modo, evita-se a transmissão de dados não utilizados e consequentemente uma redução no volume de dados transferido. Na carga de trabalho com interrogações mais complexas, o algoritmo adaptativo conseguiu reduzir o volume de dados transferidos para 21,7 MB, uma redução de 55% em comparação com o método baseado na chave (Chave). No entanto, este valor é superior ao volume de dados transferido pelo algoritmo Sem Limpeza que apenas transferiu 19,7 MB, o mínimo necessário para responder à interrogação. O algoritmo Limpeza transferiu um total de 22,2 MB. Neste sentido, é possível observar que o algoritmo adaptativo realiza retransmissões para cargas de trabalho com interrogações mais complexas. O algoritmo Limpeza transferiu 12% mais dados do que os necessários para responder à interrogação. 4.2.4 Análise / Discussão de resultados Os resultados apresentados demonstram que o número de tuplos transferidos bem como o número de filtros aplicados têm impacto no desempenho do algoritmo proposto. Neste sentido, o algoritmo adaptativo proposto demonstrou conseguir evitar custos adicionais devido à filtragem e reduzir o tempo de execução para cargas de trabalho com muitos filtros. O algoritmo adaptativo proposto mostrou reduzir não só o volume de dados transferido para a nuvem como também o tempo de execução relativamente ao método de replicação incremental baseado na chave. O algoritmo adaptativo consegue detetar quando a limpeza dos filtros é desnecessária e evitar esse custo. Assim, o tempo de execução é muito próximo do tempo de execução onde não é realizada a limpeza dos filtros. Este resultado é observável nos resultados para interrogações simples. Por outro lado, para interrogações mais complexas o sistema adaptativo realiza a limpeza de filtros. Esse custo adicional de transferir mais dados é compensado pela redução do número de filtros aplicados no dispositivo remoto. O algoritmo adaptativo percebe que a remoção de filtros pode ser vantajosa e remove filtros que considera desnecessários. Deste modo, consegue oferecer tempos de execução muito perto dos tempos de estar ativamente a realizar a limpeza de filtros. Por fim, é importante mencionar que os testes realizados foram realizados num ambiente controlado onde existe uma conexão estável entre o dispositivo remoto e a nuvem. Isto significa que o volume de dados transferidos pode ter maior impacto no tempo de execução num ambiente real. Em casos onde a largura de banda seja limitada, a redução do volume de dados transferido pode ser ainda mais relevante. 4.3 Benchmark Com o intuito de determinar o desempenho da solução é importante determinar um modelo de benchmark que permita identificar se existe melhoria em relação à abordagem comum. 51 CAPÍTULO 4. AVALIAÇÃO De modo a avaliar o desempenho da cache, utilizou-se uma variação do CH-benCHmark [10]. Este benchmark realiza interrogações transacionais e analíticas. As interrogações transacionais são utilizadas para gerar dados nas fontes de dados relacionais e não relacionais. Este permitirá, não só, ter uma fonte de dados em constante mutação, mas também, perceber o impacto que as interrogações analíticas têm na fonte de dados. Por outro lado, as interrogações analíticas serão aplicadas apenas no sistema multi-store para perceber o desempenho dos diferentes sistemas. Obenchmark utilizado será uma adaptação do CH-benCHmark orientado ao evento. Esta caraterística permite verificar o desempenho do sistema a uma base de dados só com inserção. O facto de não ter operações de atualização e remoção adequa o benchmark a uma carga de trabalho típica dos dispositivos periféricos. Este benchmark substitui todas as operações de atualização e remoção por inserções de eventos. Esta abordagem permite reduzir o tempo de ingestão dos dados, tendo impacto nas interrogações. Nese caso, considera-se que o dispositivo remoto guarda eventos de um armazém. 4.3.1 CH-benCHmark Tipicamente, os dados são utilizados em aplicações transacionais ou para realizar consultas analíticas. Usualmente as bases de dados estão otimizadas para Online Transaction Processing (OLTP) ou Online Analytical Processing (OLAP).Hybrid Transactional/Analytical Processing (HTAP) está otimizado tanto para OLTP como também para OLAP. Deste modo, consegue obter melhor desempenho quando ambas as cargas de trabalho são necessárias. O CH-benCHmark avalia o desempenho de sistemas HTAP, permitindo simultaneamente testar o comportamento da base de dados para interrogações transacionais e analíticas. Este usa Transaction Processing Performance Council Benchmark C (TPC-C) para avaliar o desempenho das interrogações transacionais. Do mesmo modo, adapta as interrogações do Transaction Processing Performance Council Benchmark H (TPC-H) para avaliar o desempenho analítico do sistema. 4.3.2 Visão global Obenchmark desenvolvido é composto por cinco transações. Cada uma destas transações é composta por operações de só de leitura ou de leitura e inserção. Estas transações foram alteradas de modo a não realizarem nem remoções, nem atualizações. Para isso, foram seguidos alguns princípios. O primeiro princípio utilizado consiste em remover as colunas cujos valores são derivados. Assim, os valores são calculados na leitura em vez de na escrita. Esta abordagem evita atualizações. Os valores derivados estão mencionados nas condições de consistência mencionados no documento que descreve o benchmark TPC-C. Além disso, as remoções foram convertidas em eventos. Ou seja, em vez de remover os tuplos da tabela, são inseridos eventos. A informação de que os tuplos foram removidos pode ser derivada destes tuplos. Nos casos onde os valores derivados eram muito complexos, é mantido uma tabela com o 52 CAPÍTULO 4. AVALIAÇÃO histórico desses valores. Estes valores guardam um marcador temporal para ser possível obter o último valor inserido. Por fim, é necessário interrogar a base de dados com interrogações analíticas. Foram definidas vistas que realizam a transformação dos dados das novas tabelas para corresponder com as tabelas definidas no TPC-C. Deste modo, as interrogações analíticas não sofrem alterações e as consultas são realizadas às vistas em vez das novas tabelas. 4.3.3 Esquema 10 21k+ 100k Warehouse [T] W 3k District [T] W*10 1+ 1+ 1+ Customer [T/A] W*30k History [T] W*30k+ 5-15 0-1 Order [T/A] W*30k+ 1+ 3+ Stock [T/A] W*100k W Item [T/A] 100k Order-line [T/A] W*300k+ 0-10 Delivery [T/A] W*21k+ Delivery Order [T/A] W*210k+ Customer History [T] W*30k+ Stock History [T/A] W*100k+ W*10 Supplier [A] 10k 1+ Nation [A] 62 1+ Region [A] 5 Figura 13: Modelo de entidade-relacionamento do benchmark. 1. O texto a negrito identifica o nome da tabela. 2. Tabelas usadas para interrogações transacionais são marcadas com um T e as analíticas com um A. As tabelas utilizadas em ambas as cargas de trabalho são marcadas com T/A. 3. Os números nos blocos representa a cardinalidade das relações. 4. Os números nas ligações representam a cardinalidade da relação. 5. O símbolo de adição (+) simboliza que ocorre inserções nesta tabela. Primeiramente, o benchmark adiciona duas novas relações: delivery edelivery_orders. Estas relações são responsáveis por armazenar eventos relativos às entregas. Além disso, os dados presentes na relação new_order podem ser derivados de outras relações. Neste caso, os pedidos que ainda não foram entregues são denominados novos pedidos. Assim, pode-se utilizar a relação orders edelivery_orders para computar quais os pedidos que ainda não foram entregues. Deste modo, a relação new_order pode ser removida. 53 CAPÍTULO 4. AVALIAÇÃO Além disso, algumas colunas foram removidas porque os dados nelas contidos podiam ser calculados através de outros dados presentes na base de dados. Desta forma, os valores contidos nestas colunas não precisam de ser atualizados. Por fim, as tabelas stock_history ecustomer_history podem armazenar informação de dados computacionalmente intensivos como é o caso da quantidade de estoque e informação sobre o cliente. Nestes casos, é guardado o identificador da tabela correspondente, uma marca temporal e o valor computacional intensivo. O tuplo cuja marca temporal seja a maior representa o valor atual. 4.3.4 Transações OTPC-C é constituído por cinco transações, das quais três modificam o estado da base de dados e as restantes realizam apenas consultas. Esta secção descreve todas as mudanças realizadas nas transações de modo adequarem-se às novas relações. 4.3.4.1 Entrega No TPC-C, a transação de ”Entrega” remove um pedido de cada distrito e atualiza as tabelas order e customer. Esta abordagem não é compatível com os requisitos do benchmark. No entanto, é possível cumprir estes requisitos usando uma abordagem orientada ao evento. Esta abordagem usa uma tabela para armazenar todas as entregas e os seus respetivos pedidos. Portanto, no primeiro ele guarda a marca temporal e o identificador do operador. Além disso, a marca temporal e o identificador do armazém permitem identificar cada entrega. A relação deliver_orders armazena o identificador do pedido onde foi entregue. Primeiramente, as novas entregas são inseridas com a informação sobre o armazém, uma marca temporal e o identificador do operador. Depois, os pedidos a serem entregues são selecionados. Neste caso, seleciona-se o pedido mais antigo de cada distrito que ainda não foi entregue. Para determinar quais os pedidos que ainda não foram entregues, verifica-se o último pedido entregue de cada distrito. Por fim, no caso de existir um novo pedido, insere-se a nova entrega. Em contraste com o TPC-C original, a relação order_line não é atualizada nem o balanço do cliente. Estes valores podem ser facilmente calculados. Para isso são utilizadas as condições de consistência que aparecem no documento que descreve o TPC-C. 4.3.4.2 Novo pedido A transação ”Novo pedido” é a operação mais frequente do TPC-C. Esta transação insere um novo pedido e atualiza o estoque e o contador de pedidos de cada distrito. Em oposição ao TPC-C original, os novos pedidos inseridos não contêm informação acerca da entrega. Esta informação é armazenada na relação delivery_order. Primeiramente, à semelhança do TPC-C 54 CAPÍTULO 4. AVALIAÇÃO original, é verificado se todos os itens são válidos. Caso contrário, a transação é abortada. Além disso, é coletada informação sobre os itens como a quantidade total e a sua localidade. Do mesmo modo, a transação obtém informação sobre os impostos do armazém e do distrito. O identificador do novo pedido é o sucessor do último pedido do distrito. Esta operação é eficiente porque tira proveito de um índice. Além disso, o estoque é atualizado com base no último registo no stock_history. Por fim, adiciona-se o registo da nova quantidade de estoque. Em contraste com o TPC-C, algumas colunas como o s_ytd pertencente à relação stock não precisam de ser atualizadas porque podem ser computadas com os dados da tabela order_line. Finalmente, insere-se na tabela order_line informação como a quantidade, a quantia e dados sobre a distribuição. 4.3.4.3 Pagamento A transação de pagamento atualiza o balanço do cliente. Além disso, atualiza estatísticas sobre o armazém e o distrito. Mais uma vez estas estatísticas podem ser derivadas através de outros dados presentes na base de dados. Nesta transação, insere-se o evento de um novo pagamento. Este evento tem associado informação como marca temporal e a quantia. Esta informação permite calcular o balanço do cliente e outras estatísticas. Deste modo o balanço do cliente, do distrito e do armazém não é atualizado, mas pode ser derivado a partir desta informação. Clientes com mau crédito precisam de atualizar os seus dados. À semelhança de outras transações, é utilizado o último registo do cliente para determinar o valor atual. Os novos valores são inseridos na tabela customer_history. 4.3.4.4 Estado do pedido Existem duas transações que não modificam o estado da base de dados. Esta transação retorna o estado do último pedido efetuado pelo cliente. O último pedido do cliente é obtido interrogando a base de dados pelo pedido com maior identificador. O estado do pedido é determinado com base nos eventos de entrega. Se para o pedido já existir um evento de entrega então o pedido já foi concluído. Caso contrário o pedido ainda se encontra pendente. Esta transação retorna o balanço do cliente. Este é calculado ao subtrair a quantia de todos os pedidos entregues com a soma de todos os pagamentos como mencionado nas condições de consistência do TPCC. 4.3.4.5 Nível de estoque Por fim, a transação ”Nível de estoque” conta o número de itens dos últimos 20 pedidos que estão abaixo do limiar de estoque. 55 CAPÍTULO 4. AVALIAÇÃO Por um lado, é necessário interrogar pelos itens dos últimos 20 pedidos. Por outro, o benchmark utiliza a relação stock_history para determinar o estoque de cada item. Assim, são selecionados os últimos 20 pedidos e as suas linhas. Posteriormente, determina-se o stock de cada um dos itens. O estoque de cada item é o último valor guardado na tabela stock_history de cada item. No final, são filtrados todos os itens cujo estoque é superior ao limiar previamente definido. 4.3.5 Implementação Obenchmark desenvolvido utiliza Python TPC-C e OLTPBench. O primeiro serviu de base para a implementação da parte transacional. O segundo benchmark foi necessário modificar para corresponder com as novas tabelas. Desta feita, foram criadas vistas que faziam o mapeamento entre as novas tabelas e as antigas. Deste modo, não foi necessário modificar as interrogações. Além disso, foram adicionados índices para melhorar o desempenho das consultas. Os programas desenvolvidos foram implementados para PostgreSQL e MongoDB. O MongoDB utiliza Aggregation Framework para responder a interrogações mais complexas. 4.4 Avaliação do uso da cache Por fim, é importante perceber se a utilização da cache oferece uma redução efetiva no tempo de execução das interrogações. Neste sentido, é importante determinar se o custo de realizar a sincronização e a atualização da cache permitem efetivamente reduzir o tempo de execução das interrogações na nuvem. Do mesmo modo, mediu-se o número de transações realizadas na fonte de dados de modo a perceber o impacto que a sincronização tem neste. 4.4.1 Benchmark Para realizar as medições de desempenho com a cache e sem esta, foi utilizado o CH-benCHmark orientado ao evento, como descrito na Seção 4.3.1. Este benchmark realiza transações na fonte de dados e interrogações analíticas na nuvem. Estas transações inserem novos tuplos que são posteriormente consultados pela nuvem para responder às interrogações. Utilizando este benchmark consegue-se perceber o desempenho do motor de consulta na nuvem e medir o impacto que a sincronização tem na fonte de dados. No entanto, algumas interrogações foram excluídas por um de dois motivos: o tempo de execução destas eram proibitivos ou realizavam operações que não eram suportadas pelo conector desenvolvido. Por estes motivos, apenas foram realizadas algumas das interrogações. 56 CAPÍTULO 4. AVALIAÇÃO 4.4.2 Ambiente de testes Os testes foram realizados em dois servidores em localizações diferentes. O servidor que simulava um dispositivo remoto, e realizava uma carga de trabalho transacional, encontrava-se em Eemshaven, Países Baixos. O servidor que simulava a nuvem, e executada interrogações analíticas, encontrava-se em Varsóvia, Polônia. Ambos os servidores tinham um CPU dual-core de 2,2 GHz, 2 GB memória RAM disponível e armazenamento SSD. Também, em ambos os servidores foi utilizado o sistema operativo Debian. As interrogações transacionais foram realizadas com uma variação do programa OLTPBench [11]. As interrogações analíticas utilizaram uma variação do programa PyTPCC [32]. A latência entre os dois servidores foi medida em 25 ms e a largura de banda limitada a 50 Mbit/s. O benchmark foi executado durante 10 minutos e com 3 minutos de aquecimento. As medições foram realizadas sobre duas bases de dados PostgreSQL. Uma das instâncias realizou interrogações analíticas e outra transacionais. Utilizou-se Foreign Data Wrappers para consultar os dados presentes no servidor remoto. Para medir o desempenho do sistema sem cache, foi criado um FDW que aplica os filtros na fonte de dados remota. Por outro lado, para medir o desempenho da cache foi utilizado oFDW desenvolvido, utilizando o algoritmo de sincronização proposto. 4.4.3 Resultados As experiências realizadas pretendem medir tanto o tempo de execução das interrogações analíticas, bem como o impacto no desempenho do dispositivo periférico. Neste sentido, pretende-se perceber se o uso da cache reduz efetivamente o tempo de execução das interrogações analíticas na nuvem. Do mesmo modo, pretende-se compreender se o uso da cache tem impacto no desempenho transacional do dispositivo periférico. Sendo as operações transacionais o principal foco destes dispositivos, o impacto das interrogações analíticas deve ser reduzido. 4.4.3.1 Tempo de execução das interrogações analíticas De modo a comparar o tempo de execução das interrogações analíticas com e sem o uso da cache foram realizadas um conjunto de interrogações. Os tempos médios de execução de cada interrogação está presente na Figura 14. Os resultados obtidos demonstraram que o tempo de execução com a cache geralmente resulta em melhores tempos de execução. As interrogações Q7,Q13 eQ17 foram as que obtiveram maior redução no tempo de execução, tendo reduzido respetivamente 57%, 64% e 70%. A Q7 realiza junção de múltiplas tabelas. Por outro lado, a interrogação Q13 eQ17 realizavam múltiplas consultas à mesma relação. Outras interrogações que realizam múltiplas consultas à mesma relação também obtiveram reduções superiores a 50%. Do mesmo modo, interrogações que realizavam consultas a tabelas estáticas, como é o caso da relação item, obtiveram bons resultados, como, por exemplo, a interrogação Q16. Regra geral, a cache 57 BIBLIOGRAFIA [12] Y. Dodis, L. Reyzin e A. Smith. “Fuzzy Extractors: How to Generate Strong Keys from Biometrics and Other Noisy Data”. Em: Advances in Cryptology - EUROCRYPT 2004. 2004, pp. 523–540. [13] Dremio Data Reflections Overview & Best Practices. url: https://hello.dremio.com/wpdata-reflections-best-practice-and-overview.html (acedido em 18/12/2020). [14] D. Eppstein, M. T. Goodrich, F. Uyeda e G. Varghese. “What’s the Difference? Efficient Set Reconciliation without Prior Context”. Em: SIGCOMM Comput. Commun. Rev. 41.4 (2011), pp. 218–229. [15] Extracting data: Stitch Documentation. url: https://www.stitchdata.com/docs/replication (acedido em 04/01/2021). [16] H. Fang. “Managing data lakes in big data era: What’s a data lake and why has it became popular in data management ecosystem”. Em: 2015 IEEE International Conference on Cyber Technology in Automation, Control, and Intelligent Systems (CYBER). 2015, pp. 820–824. [17] J. Goldstein e P. Åke Larson. “Optimizing Queries Using Materialized Views: A Practical, Scalable Solution”. Em: SIGMOD Rec. 30.2 (2001), pp. 331–342. [18] M. T. Goodrich e M. Mitzenmacher. “Invertible bloom lookup tables”. Em: 2011 49th Annual Allerton Conference on Communication, Control, and Computing (Allerton). 2011, pp. 792–799. [19] J. Gray, P. Helland, P. O’Neil e D. Shasha. “The dangers of replication and a solution”. Em: Proceedings of the 1996 ACM SIGMOD international Conference on Management of Data. 1996, pp. 173– 182. [20] T. P. G. D. Group. PostgreSQL 12 Documentation: SQL Commands - VACUUM. 2019. url: https: //www.postgresql.org/docs/12/sql-vacuum.html (acedido em 04/07/2021). [21] R. Guerraoui e A. Schiper. “Software-based replication for fault tolerance”. Em: Computer 30.4 (1997), pp. 68–74. [22] Z. Guo, G. Fox e M. Zhou. “Investigation of Data Locality in MapReduce”. Em: 2012 12th IEEE/ACM International Symposium on Cluster, Cloud and Grid Computing (ccgrid 2012). 2012, pp. 419–426. [23] R. Kimball e J. Caserta. “Extracting Changed Data”. Em: Wiley, 2009, pp. 106–110. [24] B. Kolev, P. Valduriez, C. Bondiombouy, R. Jiménez-Peris, R. Pau e J. O. Pereira. “CloudMdsQL: Querying Heterogeneous Cloud Data Stores with a Common Language”. Em: Distributed and Parallel Databases 34.4 (2016), pp. 463–503. [25] Li Fan, Pei Cao, J. Almeida e A. Z. Broder. “Summary cache: a scalable wide-area Web cache sharing protocol”. Em: IEEE/ACM Transactions on Networking 8.3 (2000), pp. 281–293. [26] S. Melnik, A. Gubarev, J. J. Long, G. Romer, S. Shivakumar, M. Tolton e T. Vassilakis. “Dremel: Interactive Analysis of Web-Scale Datasets”. Em: Proc. of the 36th Int’l Conf on Very Large Data Bases. 2010, pp. 330–339. 64 BIBLIOGRAFIA [27] Y. Minsky, A. Trachtenberg e R. Zippel. “Set reconciliation with nearly optimal communication complexity”. Em: IEEE Transactions on Information Theory 49.9 (2003), pp. 2213–2218. [28] M. T. Özsu e P. Valduriez. Principles of distributed database systems. Springer, 1999. [29] M. R. Palattella, M. Dohler, A. Grieco, G. Rizzo, J. Torsner, T. Engel e L. Ladid. “Internet of Things in the 5G Era: Enablers, Architecture, and Business Models”. Em: IEEE Journal on Selected Areas in Communications 34.3 (2016), pp. 510–527. [30] PostgreSQL: Documentation: 12: 5.12. Foreign Data. 2021. url: https://www.postgresql. org/docs/12/ddl-foreign-data.html (acedido em 23/07/2021). [31] P. Potineni. Oracle Database Data Warehousing Guide, 21c. 2021. [32] Python TPCC. url: https://github.com/apavlo/py-tpcc/wiki (acedido em 24/11/2021). [33] H. Ramadhan, F. I. Indikawati, J. Kwon e B. Koo. “MusQ: A Multi-Store Query System for IoT Data Using a Datalog-Like Language”. Em: IEEE Access 8.0 (2020), pp. 58032–58056. [34] R. van Renesse e R. Guerraoui. “Replication Techniques for Availability”. Em: Springer Berlin Heidelberg, 2010, pp. 19–40. [35] I. Savnik. “Index Data Structure for Fast Subset and Superset Queries”. Em: Availability, Reliability, and Security in Information Systems and HCI. 2013, pp. 134–148. [36] R. Sethi, M. Traverso, D. Sundstrom, D. Phillips, W. Xie, Y. Sun, N. Yegitbasi, H. Jin, E. Hwang, N. Shingte e C. Berner. “Presto: SQL on Everything”. Em: 2019 IEEE 35th International Conference on Data Engineering (ICDE). 2019, pp. 1802–1813. [37] W. Shi e S. Dustdar. “The Promise of Edge Computing”. Em: Computer 49.5 (2016), pp. 78–81. [38] The mathematics of Minisketch sketches. url: https : / / github . com / sipa / minisketch / blob/master/doc/math.md (acedido em 03/12/2021). [39] A. Trachtenberg, D. Starobinski e S. Agarwal. “Fast PDA synchronization using characteristic polynomial interpolation”. Em: Twenty-First Annual Joint Conference of the IEEE Computer and Communications Societies. 2002, pp. 1510–1519. [40] P. Vassiliadis. “A survey of extract transform load technology”. Em: International Journal of Data Warehousing and Mining (IJDWM) 5.3 (2009), pp. 1–27. [41] M. Zaharia, M. Chowdhury, M. J. Franklin, S. Shenker, I. Stoica et al. “Spark: Cluster computing with working sets.” Em: HotCloud 10.10 (2010), p. 95. 65