Full text
Universidade do Minho Escola de Engenharia Nelson José Dias Teixeira HyLake: Atualização de lagos de dados com granularidade fina setembro de 2021 UMinho | 2021 Nelson José Dias Teixeira HyLake: Atualização de lagos de dados com granularidade fina
Nelson José Dias Teixeira HyLake: Atualização de lagos de dados com granularidade fina Dissertação de Mestrado Mestrado em Engenharia Informática Trabalho efetuado sob a orientação do: Professor Doutor José Orlando Roque Nascimento Pereira Doutor Fábio André Castanheira Luís Coelho Universidade do Minho Escola de Engenharia setembro de 2021
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 RepositoriUM da Universidade do Minho. Licença concedida aos utilizadores deste trabalho Creative Commons Atribuição-NãoComercial-SemDerivações 4.0 Internacional CC BY-NC-ND 4.0 ?iiTb,ff+`2iBp2+QKKQMbXQ`;fHB+2Mb2bf#v@M+@M/f9Xyf/22/XTi iv
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. , (Local) (Data) (Nelson José Dias Teixeira) v
“If the plan doesn’t work, change the plan, but never the goal.” vi
Agradecimentos Com a concretização do vigente documento, dou por terminado o respetivo tema de dissertação, concluindo assim mais uma etapa importante da minha vida académica. Esta fase encerra um ciclo de estudos árduo e desafiante, sobretudo em época de pandemia, pelo que não seria completado com sucesso sem a presença de muitas pessoas neste longo percurso de estudos. Como tal, quero agradecer a todas as entidades que fizeram parte desta caminhada durante os últimos cinco anos, tanto a nível institucional como a nível pessoal. Em primeiro lugar, quero expressar a minha gratidão para com o meu orientador, Professor Doutor José Orlando Roque Nascimento Pereira, que aceitou a minha candidatura para este tema de forma transparente e que se demonstrou sempre disponível no esclarecimento de dúvidas em todos os estágios deste projeto. Para além disso, tenho que destacar o tempo despendido nas várias sessões estipuladas ao longo do ano letivo e, como não poderia deixar de ser, as diversas discussões construtivas e esclarecedoras que tivemos acerca do tópico subjacente a este trabalho. Por fim, queria reconhecer a experiência e o conhecimento manifestados pelo professor, que acabaram por se tornar fulcrais na finalização do mesmo, sobretudo no que toca às decisões mais complexas a nível conceptual e de implementação. Quero também agradecer ao coorientador deste tema, Doutor Fábio André Castanheira Luís Coelho, pelas opiniões emitidas em relação à redação deste documento, tanto na perspetiva de ex-aluno da academia como de um profissional dotado neste ramo da informática. À instituição de investigação INESCTEC (Instituto de Engenharia de Sistemas e Computadores, Tecnologia e Ciência), em particular ao grupo HASLab ( High-Assurance Software Laboratory ), quero destacar a entreajuda e as ótimas condições exibidas durante o desenvolvimento deste trabalho. Também quero saudar o apoio exibido por parte desta instituição e da FCT (Fundação para a Ciência e a Tecnologia) através de um financiamento plurimensal sob a forma de uma bolsa de iniciação à investigação (9034/BII-E_B4/2021). Quero também agradecer à comunidade académica da Universidade do Minho pela minha formação e pela forma como fui recebido por todas as partes envolvidas. É mais do que evidente o ambiente saudável vivido nesta casa, desde o momento em que entrei na mesma até à data atual. A bagagem enriquecida que levo desta instituição permite-me enfrentar comodamente o mercado de trabalho, pelo que me resta agradecer novamente à Universidade do Minho por transformar a possibilidade de assegurar um futuro risonho na minha carreira profissional numa realidade. Não posso também deixar de salientar todos os docentes que integraram a Licenciatura em Ciências vii
da Computação (LCC) entre 2016 e 2019. Todos eles foram importantes no culminar deste percurso, uma vez que foi nessa fase que me foram dadas as bases de conhecimento expectáveis nas áreas da matemática e da computação. Aos restantes professores pertencentes ao Departamento de Informática da Universidade do Minho e que fazem parte do Mestrado em Engenharia Informática (MEI), obrigado por terem consolidado e expandido o meu conhecimento nos diferentes tópicos intrínsecos a este curso, nomeadamente na área de sistemas distribuídos e de ciência de dados. Aos meus colegas de curso obrigado pela boa disposição, pela camaradagem, pela amizade e pelas histórias e momentos passados que fizeram deste percurso uma aventura incrível e memorável. Por fim, quero agradecer a todos os elementos da minha família por me terem apoiado nesta jornada, em particular, aos meus pais que me incentivaram a integrar este curso e que me ofereceram uma oportunidade que nunca lhes tinha sido dada anteriormente. viii
Resumo HyLake: Atualização de lagos de dados com granularidade fina Os lagos de dados, também conhecidos por data lakes , suportam a recolha de grandes quantidades de informação em ficheiros imutáveis para processamento analítico. No entanto, tem surgido a necessidade de modificar e atualizar esta informação de forma fiável, seja porque os dados são recebidos de forma incremental (por exemplo, de sensores e outras fontes de eventos) ou para eliminar os mesmos (por exemplo, devido ao RGPD (Regulamento Geral sobre a Proteção de Dados)). As soluções atuais para o fazer não são no entanto ideais: o armazenamento em SGBD (Sistema de Gestão de Bases de Dados) NoSQL (Not only SQL) tem um grande impacto no desempenho analítico, enquanto que sistemas baseados em ficheiros, como o Delta Lake , permitem apenas atualizações de granularidade grossa. Neste trabalho aborda-se este problema propondo uma solução híbrida que combina o armazenamento de longo prazo em ficheiros com um armazenamento transitório num SGBD NoSQL de forma a obter as vantagens de ambos os sistemas. Para o efeito, é implementado uma prova de conceito usando Spark , com ficheiros Parquet ,e MongoDB . Assim, com a introdução deste sistema pretende-se possibilitar a execução de transações frequentes e de granularidade fina para suportar uma carga de trabalho OLTP (Online Transaction Processing) . Os resultados experimentais obtidos confirmam que esta proposta obtém desempenho analítico e transacional comparável a cada um dos sistemas isolados. Palavras-chave: Lagos de dados, Transações, Processamento híbrido transacional-analítico, Bases de dados, Sistemas distribuídos. ix
Índice de Tabelas 1 Caraterísticas de um data warehouse ......................... 12 2 Caraterísticas de um data lake ............................ 13 3 Esquema com o conteúdo relativo à ação metaData .................. 24 4 Esquema com o conteúdo relativo à ação add ..................... 25 5 Esquema com o conteúdo relativo à ação remove ................... 26 6 Esquema com o conteúdo relativo à ação txn ..................... 27 7 Esquema com o conteúdo relativo à ação protocol ................... 28 8 Cardinalidade de cada tabela da ferramenta TPC-H em função do fator a6 ....... 37 9 Métodos da interface do sistema HyLake ....................... 57 10 Variáveis de instância da classe do sistema HyLake .................. 58 11 Métodos secundários da classe do sistema HyLake .................. 59 12 Configuração global dos testes de desempenho .................... 62 13 Configuração dos testes de desempenho executados num ambiente computacional local 63 14 Configuração dos testes de desempenho executados num ambiente computacional remoto 67 xvi
Índice de Listagens 1 Ação correspondente à modificação dos metadados presentes numa tabela delta .. 25 2 Ação correspondente à adição de um ficheiro numa tabela delta .......... 26 3 Ação correspondente à remoção de um ficheiro numa tabela delta ......... 27 4 Ação relativa ao progresso da execução de uma transação numa tabela delta .... 28 5 Ação relativa à atualização do protocolo de uma tabela delta ............ 29 6 Ação associada à metainformação de uma operação executada numa tabela delta . 29 7 Acesso ao estado transacional de uma tabela delta ................ 40 8 Leitura de uma tabela delta ............................ 40 9 Modelo de armazenamento idealizado para um ficheiro JSON em MongoDB .... 41 10 Documento relativo à operação upsert do sistema HyLake ............. 48 11 Computação do conteúdo a ser atualizado ou inserido numa tabela HyLake .... 49 12 Documento relativo à remoção de dados presentes numa tabela HyLake ...... 50 13 Computação do conteúdo a ser removido numa tabela HyLake ........... 51 14 Filtragem de documentos de acordo com a versão atual da tabela HyLake ..... 53 15 Desconstrução das linhas persistidas num documento em MongoDB ........ 53 16 Adição e remoção de linhas entre dois dataframes distintos ............ 54 17 Modelo de armazenamento do sistema HyLake com pontos de verificação ..... 63 xvii
Glossário big data Conceito que descreve o grande volume de dados estruturados e não estruturados que são gerados a cada segundo. 1,3,5,6,7,8,11 cloud computing Modelo de desenvolvimento, baseado na nuvem, que disponibiliza uma quantidade flexível de capacidade de armazenamento e processamento para satisfazer as necessidades inerentes à implementação de um projeto computacional. 1 cluster Designação atribuída a um conjunto de computadores independentes que se encontram ligados a uma rede de modo a desempenhar operações complexas. Ao atuarem sobre as mesmas tarefas computacionais, estes equipamentos, também conhecidos por nós, formam um sistema unificado para dar resposta às necessidades dos seus utilizadores, acelerando substancialmente a capacidade de processamento de informação. 19,20,58,67 online Termo utilizado no âmbito da informática para tanto designar uma ligação ou conexão à rede com sucesso como para indicar que um determinado sistema se encontra operacional ou pronto a usar. 3,6 query Ação exercida sobre uma base de dados de forma a solicitar informações acerca de uma ou mais tabelas presentes na mesma. 9,10,14,15,21,37,38,39,49,73 xviii
Acrónimos e Siglas ACID Atomicidade, Consistência, Isolamento e Durabilidade 2,3,8,9,20,24 API Application Programming Interface 19,36,40,47 CPU Central Processing Unit 67 CSV Comma-Separated Values 15 DBMS Database Management System x ETL Extract , Transform , Load 12 GB Gigabyte 35,36,38,39,67 GDPR General Data Protection Regulation x HDFS Hadoop Distributed File System 61 HTAP Hybrid Transactional Analytical Processing x,xiv,8,10,11,18,72,73 JSON JavaScript Object Notation xvii,3,15,22,23,29,31,33,40,41,45,48,50,53,72 JVM Java Virtual Machine 36 NoSQL Not only SQL ix,x,15,36,39,44,52,72 xix
ACRÓNIMOS E SIGLAS NVMe Non-Volatile Memory Express 35 OLAP Online Analytical Processing 8,10,11,12,17,18,38,39,70,72,73 OLTP Online Transaction Processing ix,x,2,3,8,9,10,11,16,17,18,30,44,50,61,62,72,73 RAM Random Access Memory 35,67 RDBMS Relational Database Management System 8 RDD Resilient Distributed Dataset 19,20,36 RGPD Regulamento Geral sobre a Proteção de Dados ix SGBD Sistema de Gestão de Bases de Dados ix SQL Structured Query Language 8,19,20,43,54,55 SSD Solid State Drive 35 UDF User-Defined Function 38 XML Extensible Markup Language 15 xx
Capítulo 1 Introdução A utilização e vulgarização de aplicações digitais e o aumento do número de utilizadores que as frequentam diariamente acentua gravemente um problema já existente no contexto computacional, isto é, a enorme quantidade de dados por armazenar e processar. De facto, o volume de dados gerados atualmente encontra-se num nível sem precedentes, pelo que o seu tratamento ganha cada vez mais importância. Para demonstrar a influência destas plataformas nos dias de hoje e a quantidade de dados e utilizadores que lhes estão associadas, apresenta-se de seguida o exemplo da rede social Instagram 1, atualmente pertencente à empresa Facebook 2. Segundo os estudos estatísticos mais recentes [12], após uma semana a rede social Instagram ter sido fundada, cerca de cem mil utilizadores usufruíam esta plataforma para partilhar fotografias e vídeos. Em meados de fevereiro de 2011, aproximadamente dois meses após a sua criação, registou-se nesta aplicação um milhão de utilizadores ativos por mês, evidenciando-se um crescimento significativo. Apesar desta empresa ter sido comprada pela Facebook no início de 2012, hoje em dia entre quinhentos a seiscentos milhões de utilizadores usam diariamente a aplicação Instagram , sendo que, por mês, este valor já disparou para um marco histórico de um bilião de utilizadores. Tal como seria de esperar, este número de utilizadores traduz-se num volume considerável de informação consumida pelos mesmos que, no caso desta rede social, corresponde à partilha diária de noventa e cinco milhões de fotografias, perfazendo um total de quarenta biliões de fotografias partilhadas desde a sua criação [26]. Em resposta, os últimos anos testemunharam a consolidação do uso do paradigma de computação na nuvem, também conhecido por cloud computing . Este foi um dos maiores impulsionadores do conceito big data , permitindo não só a ingestão de grandes quantidades de informação como o desbloqueio dos casos de uso de analítica de dados. O armazenamento e processamento de informação em sistemas 1?iiTb,ff#QmiXBMbi;`KX+QK 2?iiTb,ff#QmiX7+2#QQFX+QK 1
CAPÍTULO 1. INTRODUÇÃO cloud tem vindo a ser uma solução bastante popular entre as diversas companhias do meio tecnológico, sendo que Google 3, Microsoft 4e Amazon 5constituem-se como três das maiores empresas a disponibilizarem este tipo de serviços. Google Cloud [11], Microsoft Azure [24]e Amazon Web Services [25] são alguns exemplos destes serviços e permitem aos seus clientes usufruírem de infraestruturas desejáveis para o desenvolvimento das suas soluções. Esta preferência deve-se sobretudo ao facto de os utilizadores poderem escalar significativamente os recursos de computação e de armazenamento separadamente. Para além disso, aspetos como a disponibilidade, a gestão e monitorização de infraestruturas, a simplicidade, a segurança e o custo reduzido na adesão a estes serviços acabam por tornar ainda mais atrativa a utilização destes sistemas num contexto empresarial. Assim, a virtualização e utilização do modelo de desenvolvimento baseado na nuvem são cada vez mais consideradas como a primeira escolha para o deployment das aplicações atuais. 1.1 Problema Todavia, a enorme quantidade de dados alojados neste tipo de repositórios e a sua desnormalização fazem com que se recorra frequentemente ao armazenamento por objetos ou chave-valor. Por conseguinte, a execução de transações torna-se cada vez mais delicada e complexa, quer a nível de desempenho quer a nível de consistência. Tanto os casos de uso analíticos atuais como os sistemas que os suportam possibilitam a ingestão de dados de várias fontes, criando pipelines com dados heterogéneos. Nestes ambientes, a oferta de propriedades transacionais é limitada ou até inexistente, sendo que neste último caso é ainda mais difícil considerarem-se várias fontes de dados em simultâneo. Como tal, é imprescindível a adoção de um sistema que proporciona a oferta de propriedades transacionais ACID (Atomicidade, Consistência, Isolamento e Durabilidade) sobre várias fontes de dados que, por sua vez, recorrem ao armazenamento baseado em objetos de dados. Desta forma, este sistema deve ser capaz de executar várias transações de granularidade fina num ambiente analítico. 1.2 Motivação As transações descritas anteriormente dizem respeito a um ambiente computacional que utiliza cargas de trabalho OLTP (Online Transaction Processing) . Os sistemas associados a este ambiente encarregamse de registar todas as transações contidas numa determinada operação organizacional. A carga de trabalho envolvida é identificada por uma base de dados que recebe solicitações e alterações de dados de vários utilizadores. Desta forma, evidencia-se um tempo de resposta curto na execução de transações de pequenas dimensões. Caraterísticas como a consulta rápida de dados e a preservação da integridade 3?iiTb,ff#QmiX;QQ;H2 4?iiTb,ffrrrXKB+`QbQ7iX+QKf#Qmi 5?iiTb,ffrrrX#QmiKxQMX+QK 2
1.3. OBJETIVOS E CONTRIBUIÇÕES dos mesmos fazem com que este tipo de carga de trabalho seja o mais indicado em bases de dados tradicionalmente relacionais [13,14,16–18]. Máquinas multibanco, sistemas de reservas em restaurantes e aplicações online de manutenção de registos na área da saúde são alguns exemplos de sistemas OLTP . Os ambientes analíticos modernos permitem conjugar várias fontes de dados em ambientes com alta concorrência de escritas e leituras. Tentar garantir coerência nos dados a partir de múltiplas origens é uma tarefa difícil, sobretudo num ambiente analítico. De forma a ultrapassar estes desafios computacionais e, consequentemente, garantir a consistência durante o processamento de transações, o repositório de dados considerado necessita de aderir obrigatoriamente às propriedades ACID. O sistema Delta Lake [8] oferece estas mesmas propriedades transacionais, permitindo operar nestas arquiteturas, potencialmente com muitas fontes de dados. Este sistema inclui uma camada de armazenamento que é aplicada sobre lagos de dados de modo a executar transações com as propriedades ACID, lidando com variações de esquemas de dados e a manipulação de metadados de uma forma escalável. Uma tabela delta corresponde a uma diretoria num sistema de armazenamento orientado a objetos ou num sistema de ficheiros. Nela existem objetos de dados, sob o formato Parquet [6], com o conteúdo da tabela, e um conjunto de registos, codificados na extensão JSON (JavaScript Object Notation) , com a respetiva metainformação. Com o uso da ferramenta Spark [35] é também possível executar interrogações analíticas, em particular aquelas que envolvem atualizações cirúrgicas aos dados. Para além disso, este utensílio é particularmente útil no processamento distribuído de big data . Apesar deste sistema permitir a execução de poucas transações de granularidade grossa para efetuar atualizações em lagos de dados, este apresenta-se como o candidato apropriado à elaboração de uma proposta que envolva muitas transações de pequenas proporções. 1.3 Objetivos e contribuições Atendendo às caraterísticas do sistema Delta Lake , este trabalho tem como objetivo o estudo e a expansão do limite da granularidade dos dados presentes na execução de transações relativas ao último sistema. Como tal, o foco deste tema de dissertação é possibilitar a execução de transações OLTP num ambiente analítico com as propriedades da arquitetura do sistema Delta Lake . Para isso, elaborou-se um novo sistema que permite a execução de várias transações de granularidade fina em lagos de dados. O surgimento desta solução, isto é, o sistema HyLake , tem por base a realização de diversas experiências que estudam a viabilidade do sistema Delta Lake na execução de cargas de trabalho OLTP . A implementação do sistema HyLake corresponde assim a um sistema capaz de suportar transações de atualização de dados utilizados em interrogações analíticas por sistemas distribuídos, como é o caso da ferramenta Spark . A última framework , ao proporcionar o processamento distribuído de grandes quantidades de informação, ajudará imenso na avaliação da usabilidade e escalabilidade do sistema produzido. Com a realização de leituras e escritas numa base de dados não relacional antecipam-se obter melhores resultados no sistema HyLake em relação ao que é observado no sistema Delta Lake , onde as operações em 3
CAPÍTULO 1. INTRODUÇÃO causa envolvem a manipulação de ficheiros. Assim, com a introdução do sistema HyLake pretende-se provar que é possível estender a granularidade dos dados na execução de transações analíticas no sistema Delta Lake . 1.4 Estrutura da dissertação O resto do documento está estruturado da seguinte forma: O Capítulo 2apresenta o estado de arte relativo ao sistema Delta Lake . Nele evidenciam-se alguns conceitos relevantes que complementam o sistema em causa. Posteriormente são expostos no Capítulo 3todos os testes intermédios realizados sobre o sistema Delta Lake . Nesse processo são indicadas as ferramentas computacionais utilizadas, fazendo-se referência à instalação e configuração das mesmas. Tendo em conta as experiências realizadas no capítulo anterior, procede-se, no Capítulo 4, à construção da arquitetura do sistema idealizado, isto é, o sistema HyLake . Ao longo da descrição desta nova proposta apresenta-se também a sua implementação, incluindo as operações de leitura e de escrita. Para interpretar corretamente a especificação do sistema HyLake são referenciados alguns operadores secundários que ajudam a elaborar a solução em questão. No Capítulo 5avalia-se o sistema HyLake , comparando o seu desempenho com os dos sistemas Delta Lake e MongoDB [28]. Por fim, no Capítulo 6conclui-se o vigente trabalho ao apontar o sistema que apresenta o melhor comportamento na execução de transações frequentes e de granularidade fina. Após essa tomada de decisão, é abordado o trabalho futuro deste tema de dissertação. 4
Capítulo 2 Estado de arte Esta secção do documento apresenta diversas definições e detalhes computacionais de forma a introduzir o sistema Delta Lake . Este sistema é o ponto essencial deste trabalho pelo que será explorado a nível de estruturação e de funcionalidades. 2.1 Big data Big data é a área do conhecimento que estuda como tratar, analisar e obter informações a partir de grandes conjuntos de dados. O termo em questão surgiu em 1997 e foi utilizado para referenciar a manipulação de grandes quantidades de dados não estruturados. Os principais aspetos de big data podem ser definidos à custa da regra dos três V’s : •Volume: designação relacionada com a grande quantidade de dados gerada por unidade de tempo; •Velocidade: nome que traduz o ritmo com que os dados são produzidos e manipulados; •Variedade: palavra que expressa quer a diversidade das fontes de informação quer a multiplicidade dos formatos de dados. Estas duas caraterísticas aumentam, claro está, a complexidade das análises efetuadas. A quantidade de dados gerada atualmente tem crescido de forma exponencial. Na Figura 1salientase o crescimento contínuo e evolutivo de dados ao longo dos últimos anos. De notar que os dados apontados entre 2018 e 2025 são valores baseados em predições. Da análise, a quantidade total de dados criados, capturados, copiados e consumidos globalmente aumentou rapidamente desde 2010. Por conseguinte, prevê-se que a capacidade dos sistemas de armazenamento em questão acompanhe este 5
CAPÍTULO 2. ESTADO DE ARTE organizações que adotam este modelo conseguem armazenar consistentemente a informação proveniente das suas atividades, tomando melhores decisões sobre os seus negócios. Para dar suporte a estas decisões é comum realizarem-se análises sobre grandes conjuntos de dados com recurso a cargas de trabalho OLAP . A Tabela 1e a Figura 5sumariam as caraterísticas e a arquitetura deste modelo, respetivamente. Caraterística Data warehouse Tipo de dados Relacionais, provenientes de sistemas transacionais. Armazenamento Captura de informação previamente formatada. Processamento Requer maior manutenção e utiliza o processo ETL (Extract, Transform, Load) . Esquema Concebido antes da implementação e execução de tarefas. Desempenho e custo Resultados rápidos de consulta usando armazenamento de alto custo. Qualidade de dados Elevada. Utilizadores Analistas de dados e de negócio. Finalidade Decisões analíticas. Tabela 1: Caraterísticas de um data warehouse Figura 5: Arquitetura de um data warehouse 12
2.3. MODELOS DE ARMAZENAMENTO DE DADOS 2.3.2 Data lake Um lago de dados, também conhecido por data lake , define uma coleção centralizada com dados diversificados. Estes dados possuem diferentes formatos a nível de estruturação, pelo que os mesmos se podem constituir como estruturados, semi-estruturados ou não estruturados. De forma mais informal, um data lake pode ser visto como um contentor de informação de grandes proporções onde se armazenam dados nativamente, independentemente das suas dimensões ou origens. O esquema de dados deste repositório não é definido quando os mesmos são capturados. Estes são guardados até que sejam estritamente necessários para a seleção, organização e execução de determinadas tarefas computacionais. Tal como o termo em inglês transparece, este modelo de persistência permite uma acumulação significativa de informação, sem haver qualquer tipo de tratamento prévio sobre os mesmos. Como tal, é possível produzir cargas de trabalho mistas e, em último caso, aumentar o desempenho analítico sobre este volume de informação. A Tabela 2sintetiza as caraterísticas deste modelo. Caraterística Data lake Tipo de dados Estruturados, semi-estruturados e não estruturados. Armazenamento Informação persistida no seu formato original, independentemente da sua origem. Processamento Execução moderada consoante o tipo de dados. Esquema Definido após o armazenamento de dados e identificado no momento de análise. Desempenho e custo Resultados moderados de consulta usando armazenamento de baixo custo. Qualidade de dados Diversificada, dada a natureza dos dados. Utilizadores Engenheiros e cientistas de dados. Finalidade Mineração e análise preditiva de dados. Tabela 2: Caraterísticas de um data lake Quanto à arquitetura de um data lake , pode-se inferir a sua simplicidade dado os diferentes tipos de estruturação que este modelo suporta. Por conseguinte, é possível atingir índices consideráveis de escalabilidade, algo que não acontece, por exemplo, em sistemas de armazenamento tradicionais. 2.3.3 Estruturação e formatação de dados Tal como foi possível constatar anteriormente, existem múltiplas estruturas de dados para diferentes modelos de armazenamento. Esta secção explora não só os detalhes inerentes a estes formatos como 13
CAPÍTULO 2. ESTADO DE ARTE também as vantagens e desvantagens que lhes estão associadas. Fatores como a organização, a flexibilidade e escalabilidade são realçados no momento da sua caraterização. Ao longo destas descrições são também dados alguns exemplos de formatos de dados para sustentar a definição da sua estruturação. A Figura 6apresenta alguns exemplos relativos aos três tipos de estruturação de dados presentes nos modelos de armazenamento atuais. Figura 6: Dados estruturados (1), semi-estruturados (2) e não estruturados (3) 2.3.3.1 Dados estruturados Os dados estruturados disponibilizam tanto um esquema de dados fixo e bem definido como a existência de registos orientados à linha e à coluna. Esta estruturação apresenta velocidades superiores na execução de queries e garante uma persistência de dados eficiente. O último facto transparece, de certa forma, o requerimento de uma menor capacidade de armazenamento. Assim sendo, estes podem ser vistos como registos que podem ser guardados em bases de dados relacionais. Hoje, os dados estruturados são a maneira mais comum e simples de gerir informações. No entanto, estes representam apenas cinco a dez porcento de todos os dados gerados atualmente. São casos desta estruturação de dados os formatos Parquet [6], Avro [33]e ORC [5]. Tal como a Figura 6clarifica, os dados estruturados são pouco flexíveis, possuem o maior nível organizacional de todas estruturações de dados referidas e são persistidos num formato bem definido. Nesta estruturação é nativamente disponibilizado pelo gestor de uma base de dados relacional os processos relativos ao controlo de concorrência e à gestão de transações. 14
2.3. MODELOS DE ARMAZENAMENTO DE DADOS 2.3.3.2 Dados semi-estruturados Ficheiros com extensão JSON , CSV (Comma-Separated Values) ou XML (Extensible Markup Language) adotam a semi-estruturação de dados. Este tipo de estrutura representa entre a cinco a dez porcento dos dados gerados atualmente. Apesar destes não possuírem um esquema de dados por omissão, é possível deduzir o seu conteúdo pela existência de registos delimitados por algum tipo de separador (vírgula, dois pontos, entre outros). Assim, quando comparado com o caso anterior (Secção 2.3.3.3), esta estruturação, ainda que não seja rígida, garante uma melhoria substancial na rapidez do respetivo processamento de dados e, consequentemente, obtêm-se execuções de queries mais velozes. Por conseguinte, adquirese um armazenamento eficiente e com melhor desempenho. Bases de dados não relacionais, também conhecidas pelo termo NoSQL (Not only SQL) , fazem parte deste tipo de estruturação. De salientar que este género de dados pode ser facilmente convertido para uma estrutura fixa, como é o caso dos dados estruturados (Secção 2.3.3.1). Dito isto, pode-se concluir que os dados semi-estruturados têm um grau organizacional intermédio, dispõem uma flexibilidade superior aos dados não estruturados e são mais simples de escalar. Todavia, não integram o controlo de concorrência e incorporam uma gestão de transações adaptada. 2.3.3.3 Dados não estruturados Os dados não estruturados correspondem a cerca de oitenta porcento dos dados produzidos atualmente. Estes dados possuem uma formatação heterogénea e contêm alguma redundância e ambiguidade, pelo que o seu processamento é difícil. Nesta estrutura de dados, não existe a noção de registo, isto é, não há um esquema de dados predefinido. Consequentemente, não é possível exibir este tipo de dados em linhas, colunas ou em bases de dados relacionais. Apesar dos próprios possuírem uma estrutura interna, o conteúdo dos mesmos não é incorporável numa base de dados sem a ocorrência de processamento. Desta maneira, este tipo de dados é extraído e retido em data lakes para desempenhar análises preditivas. Ficheiros com extensão TXT , isto é, ficheiros com conteúdo textual, são exemplos deste formato. Tendo em consideração a volumetria deste tipo de dados, é possível inferir que a ausência de uma estrutura fixa prejudica o processamento dos mesmos, ainda que exista flexibilidade. Em síntese, os dados não estruturados fornecem principalmente informações qualitativas, pelo que não podem ser mapeados num modelo de dados predefinido. Estes constituem-se como a estrutura de dados com maior flexibilidade e escalabilidade, uma vez que não adotam um esquema de dados concreto. Por fim, nos dados não estruturados não se verifica qualquer tipo de controlo de concorrência, pelo que também não há nenhuma gestão a nível de execução de transações. 2.3.4 Modelos lógico e físico de armazenamento de dados Nesta secção exibem-se os modelos lógico e físico de armazenamento de dados. Neles é possível ter um visão externa de como se procede à persistência de informação num sistema de armazenamento 15
CAPÍTULO 2. ESTADO DE ARTE local ou até remoto. No que toca ao plano lógico, é possível observar o conteúdo original de um conjunto de dados e inferir a sua dimensão e esquematização. Quanto ao plano físico, transparece-se a orientação dos dados persistidos num determinado local de armazenamento. Na Figura 7é possível identificar três tipos de orientações no armazenamento de dados relativamente ao plano físico: orientação à linha, à coluna e, ainda, uma estratégia híbrida. Na primeira dá-se a partição horizontal da tabela, ou seja, cada registo presente numa determinada linha é guardado segundo essa orientação, preservando a sua ordem. Na segunda ocorre a partição vertical da tabela, isto é, cada registo presente numa determinada coluna é armazenado de acordo com essa orientação, mantendo a respetiva ordem. Por fim, na terceira verifica-se uma combinação das estratégias anteriores, percorrendose e extraindo-se porções de registos presentes quer numa determinada linha quer numa coluna em específico. Figura 7: Esquematização de modelos físicos na persistência de dados 2.3.4.1 Orientação à linha Este tipo de orientação no armazenamento de dados traduz-se na partição horizontal de um conjunto de dados. A Figura 8ilustra a orientação à linha na persistência de dados. Tendo em consideração os aspetos mencionados na Secção 2.2.1, é simples concluir que esta estratégia de armazenamento adequa-se perfeitamente a uma carga de trabalho OLTP . Alterações sobre uma determinada tabela, isto é, operações como a inserção, atualização e remoção de registos afetam as linhas que contêm a informação. Por outro lado, pode-se também inferir que esta orientação não é apropriada para um ambiente onde hajam poucas 16
2.3. MODELOS DE ARMAZENAMENTO DE DADOS operações de grandes dimensões sobre um subconjunto de todas as colunas de um repositório de dados, como é o caso de uma carga de trabalho OLAP . Figura 8: Modelo físico de armazenamento de dados orientado à linha 2.3.4.2 Orientação à coluna Tal como a Figura 9demonstra, este tipo de orientação exprime-se, de forma simplificada, na partição vertical de um conjunto de dados. Recuando ao que foi exposto na Secção 2.2.2, pode-se afirmar que esta estratégia alinha-se com as caraterísticas de uma carga de trabalho OLAP . Uma vez que o interesse da última passa pela obtenção de uma porção de todas as colunas de um conjunto de dados, a execução de projeções torna-se bastante eficiente devido à contiguidade dos últimos. Dado que os registos de uma determinada coluna são sequencialmente persistidos, existe a possibilidade de efetuar uma compressão dos mesmos. Em contrapartida, esta orientação não se adequa a uma carga de trabalho OLTP porque, para a realização de operações que envolvam a inserção de registos, é necessário a adição dos mesmos em diferentes locais do armazenamento. Figura 9: Modelo físico de armazenamento de dados orientado à coluna 17
CAPÍTULO 2. ESTADO DE ARTE 2.3.4.3 Orientação híbrida Retomando o seu conceito, este tipo de orientação no armazenamento de dados representa uma junção das caraterísticas inerentes à orientação à linha (Secção 2.3.4.1) e à coluna (Secção 2.3.4.2). Desta forma, dá-se tanto a partição horizontal como a partição vertical da tabela. No caso da Figura 10, a partição horizontal é invocada a cada três linhas, sendo-lhe posteriormente aplicada a partição vertical. Este procedimento de armazenamento de dados é usado em alguns formatos, como Parquet e ORC ,e oferece tanto as propriedades de localização de registos relativas a uma carga de trabalho OLTP como a contiguidade dos últimos num ambiente OLAP . Assim, a orientação híbrida é apropriada a uma carga de trabalho HTAP . Figura 10: Modelo físico de armazenamento de dados com orientação híbrida 2.4 Spark O paradigma MapReduce [10], introduzido pela Google em 2004, é um modelo que visa a computação paralela de grandes volumes de dados segundo a estratégia divide and conquer . O surgimento deste paradigma permitiu resolver alguns problemas computacionais, como por exemplo o balanceamento de cargas de trabalho e a movimentação de dados observada durante a fase de processamento. A estratégia adotada no processamento de um conjunto de dados arbitrário pode ser decomposta em quatro etapas: divisão, mapeamento, agrupamento e ordenação e, por fim, redução. Inicialmente é realizada a divisão do conjunto de dados considerado em diversos segmentos com aproximadamente a mesma dimensão. De seguida, cada bloco de dados é mapeado num ou mais pares chave-valor, sendo os mesmos agrupados e ordenados segundo a primeira componente. Por fim, para cada chave de cada par é aplicada a redução dos respetivos valores. A Figura 11 ilustra cada uma das etapas descritas. Apesar deste modelo ser caraterizado pela sua escalabilidade, simplicidade, flexibilidade e, ainda, pela sua tolerância a faltas, é possível apontar alguns defeitos. Tal como a Figura 11 transparece, existem demasiadas escritas em disco ou em memória entre os vários estágios de processamento. Para além 18
2.4. SPARK disso, a composição destas etapas não é feita de uma forma livre. Com a ausência de abstrações de memória distribuída, observa-se ineficiência na execução de tarefas que necessitam de reutilizar cálculos intermédios em diferentes etapas de computação. Figura 11: Paradigma MapReduce De forma a ultrapassar estas adversidades, foi proposta uma nova solução, denominada por Spark [35], que introduz o conceito de RDD (Resilient Distributed Dataset) . Esta abstração de memória distribuída corresponde a uma coleção de objetos particionados que podem ser reconstruídos caso uma das partições seja perdida. Esta coleção é computada e armazenada pelos nós que constituem um determinado cluster . Com estas caraterísticas, os utilizadores desta ferramenta podem preservar RDDs em memória pelos diversos nós de um cluster e, desta forma, reutilizá-los em fases posteriores de processamento. Spark é por isso um motor de processamento distribuído que possibilita a análise de dados em grande escala. Esta ferramenta fornece APIs (Application Programming Interfaces) em diversas linguagens de programação, nomeadamente em Python [34]e Scala [29], para garantir o paralelismo de dados e a tolerância a faltas em cluster sde variadas dimensões. Ainda assim, Spark não oferece, por exemplo, um sistema de gestão de ficheiros. A quantidade de memória usufruída na execução de tarefas e a ausência de suporte no processamento de dados em tempo real são duas das suas principais limitações. Uma vez discutida a sua natureza, expõe-se de seguida os diferentes módulos disponibilizados pela ferramenta Spark : • Spark SQL : biblioteca usada no processamento de dados estruturados; • MLlib : módulo utilizado em aprendizagem automática; • GraphX : biblioteca usada no processamento de grafos; • Structured streaming : módulo utilizado no processamento de fluxo ( stream ) e na computação incremental. 19
CAPÍTULO 2. ESTADO DE ARTE Dos quatro módulos apontados nesta enumeração, apenas a biblioteca Spark SQL [7] é utilizada na construção da solução deste trabalho, dada a sua utilidade no armazenamento e na manipulação de dados estruturados. No caso do sistema elaborado, estes dados são usados sob a forma de um dataframe . Esta estrutura de dados pode ser vista como uma tabela presente numa base de dados relacional, onde cada coluna possui um único nome e tipo de dados. Na perspetiva da ferramenta Spark , um dataframe representa uma coleção de dados distribuída que inclui diversas otimizações para conseguir atingir um melhor desempenho comparativamente a um RDD . Durante a distribuição de dados pelos nós de um cluster , a ausência de técnicas otimizadas na serialização de dados traduz um dos contratempos verificados em aplicações Spark que recorrem ao uso de RDDs . Uma vez que a informação contida num dataframe é guardada sob o formato binário, não existe a necessidade de executar a serialização dos dados envolvidos, pelo que é recomendada a sua utilização. 2.5 Delta Lake O modelo de armazenamento de dados subjacente a um lago de dados é atualmente utilizado para a persistência de informação com diferentes estruturas. Como tal, é comum verificar-se a coleção da mesma, numa primeira fase, e posteriormente a sua persistência. Com estes dados são subsequentemente exercidas análises em contextos de ciência de dados ou até de aprendizagem automática. Todavia, a qualidade inicial dos dados coletados é bastante reduzida, sobretudo pela sua heterogeneidade. Desta forma, não é possível assegurar a sua qualidade, pelo que este processo acaba por se tornar inútil para o efeito pretendido, ou seja, extrair valor dos dados previamente armazenados. A esta contrariedade juntam-se também alguns aspetos negativos, como por exemplo o facto de num lago de dados não existir atomicidade e, ainda, consistência ou isolamento dos dados. Na eventualidade de ocorrerem falhas durante a execução de uma tarefa computacional, os dados poderão terminar numa vista incoerente. Com a introdução do sistema Delta Lake [8] pretende-se combater todos os aspetos negativos que foram expostos anteriormente em relação a um lago de dados. A principal caraterística que distingue este sistema de um lago de dados reside na inclusão de uma camada de armazenamento onde é possível executar transações que integram as propriedades ACID (Secção 2.2.1.1). Assim, conseguem-se obter ambientes transacionais mais coerentes num âmbito de execução analítica distribuída. Para além disso, este utiliza a ferramenta Spark como o motor de computação distribuída. Com a integração da mesma, dá-se a unificação de tarefas que envolvem o processamento de dados em tempo real ( streaming )e em lotes ( batch ). Ainda acerca da incorporação da ferramenta referida, o sistema Delta Lake utiliza-a para proceder a uma manipulação escalável dos metadados subjacentes. Como tal, realizam-se escritas pontuais derivadas de operações analíticas de uma forma coerente, algo que não acontece naturalmente num lago de dados. Outra funcionalidade importante é a capacidade que o mesmo tem em lidar com variações de esquema de dados. Este último tópico é tratado automaticamente pelo sistema Delta Lake para evitar a inserção de registos inválidos durante o processamento de uma tarefa. Por fim, o sistema 20
2.5. DELTA LAKE em questão disponibiliza aos seus utilizadores o acesso a todas as versões dos dados persistidos, sendo possível observar o histórico de todas as operações efetuadas e, ainda, executar queries sobre uma versão em específico. 2.5.1 Arquitetura A Figura 12 ilustra a categorização de dados segundo a sua qualidade. Esta divisão é efetuada por parte do sistema Delta Lake para incrementalmente melhorar a qualidade dos dados até que estes estejam disponíveis para serem consumidos. Nesta segmentação evidenciam-se três classes de dados distintas. Na primeira, a categoria bronze, apresentam-se todos os dados que ainda possuem a sua estrutura original, ou seja, que permaneceram inalterados. Com o armazenamento destes, ainda que não tenham a estrutura desejada, consegue-se preservar toda a informação desde o momento da sua génese. Para além disso, ao proceder-se à persistência inicial destes dados, evita-se prematuramente o pré-processamento dos mesmos. Quanto aos dados pertencentes à categoria prata, estes, à semelhança da classe anterior, também não se encontram prontos para serem utilizados. Contudo, já contêm algum tipo de tratamento, como por exemplo filtragens, incorporação de esquemas de dados ou até junções com outros conjuntos do mesmo tipo. Desta forma, este grupo já possui uma ligeira melhoria a nível de estruturação face ao que se observou na categoria bronze. Esta categoria também pode ser visualizada como um classe intermédia onde se podem realizar processos de verificação sobre os dados de forma a detetar algumas anomalias. Por fim, na classe ouro, existe uma coletânea de dados totalmente processados para posteriormente serem consumidos em operações de análise. Relativamente à execução de operações do tipo streaming , estas movimentam os dados através do sistema Delta Lake , percorrendo as diversas classes mencionadas anteriormente. Com este comportamento é possível eliminar a gestão do escalonamento das respetivas tarefas computacionais. Figura 12: Processo de categorização de dados no sistema Delta Lake 21
CAPÍTULO 2. ESTADO DE ARTE tarefa. Os identificadores em causa são armazenados em pares, onde o primeiro componente, appId ,é o identificador único do processo que modifica a tabela e o segundo componente, version , transparece o progresso obtido por uma determinada aplicação. O registo atómico destas informações junto com as modificações na tabela permite que os sistemas externos em questão façam as suas alterações numa tabela delta idempotente. Para assimilar melhor esta estruturação de dados, apresenta-se na Listagem 4 um exemplo desta ação. 1& 2]itM],& 3]TTA/],]j#Rj3dk@k/9d@92Rd@3ey@kR7/kkkjN8]- 4]p2`bBQM],je99d8 5' 6' Listagem 4: Ação relativa ao progresso da execução de uma transação numa tabela delta ⇒ Change Protocol :esta ação é usada para atualizar a versão do protocolo transacional do sistema Delta Lake na execução de leituras e escritas. Por conseguinte, é possível interpretar os metadados da tabela e os respetivos registos transacionais em diferentes versões do protocolo delta . Assim, os utilizadores que possuem versões antigas ou mais recentes deste sistema podem executar normalmente as operações de leitura ou de escrita numa determinada tabela. De notar que a versão do protocolo será aumentada sempre que alterações não compatíveis com versões futuras forem feitas nesta especificação. Uma vez que mudanças significativas sobre uma tabela delta devem ser acompanhadas por um aumento na versão do protocolo, os utilizadores podem assumir que as propriedades ou ações não reconhecidas nunca são necessárias para interpretar corretamente os registos transacionais da mesma. A Tabela 7apresenta o esquema de dados relativo a esta ação. Nome do campo Descrição KBM_2/2`o2`bBQM Versão mínima do protocolo de leitura de uma tabela delta que um utilizador deve adotar de modo a ler corretamente a tabela em causa. KBMq`Bi2`o2`bBQM Versão mínima do protocolo de escrita de uma tabela delta que um utilizador deve seguir de maneira a escrever corretamente na tabela em questão. Tabela 7: Esquema com o conteúdo relativo à ação protocol Para interpretar na totalidade o conteúdo da Tabela 7, evidencia-se na Listagem 5um exemplo da ação protocol num registo transacional. 28
2.5. DELTA LAKE 1& 2]T`QiQ+QH],& 3]KBM_2/2`o2`bBQM],R4]KBMq`Bi2`o2`bBQM],k 5' 6' Listagem 5: Ação relativa à atualização do protocolo de uma tabela delta ⇒ Commit Info :esta ação indica alguns detalhes que dizem respeito à alteração efetuada. Aspetos como a metainformação, a versão da tabela e as listas de ficheiros e de transações são devidamente apresentados. Neste tipo de ação, qualquer formato JSON válido pode ser persistido, pelo que não existe um esquema de dados predefinido. Tal como se verificou para as ações anteriores, expõe-se na Listagem 6um exemplo para a ação commitInfo . 1& 2]+QKKBiAM7Q],& 3]iBK2biKT],ReRk888deR3y84]mb2`A/],]3]- 5]mb2`LK2],]!#X+QK]- 6]QT2`iBQM],]q_Ah1]- 7]QT2`iBQMS`K2i2`b],& 8]KQ/2],]1``Q`A71tBbib]- 9]T`iBiBQM"v],]()] 10 '- 11 ]Bb"HBM/TT2M/],i`m212 ]QT2`iBQMJ2i`B+b],& 13 ]MmK6BH2b],]9e]- 14 ]MmKPmiTmi"vi2b],]9k39NNyNy]- 15 ]MmKPmiTmi_Qrb],]RRNNdNNe] 16 '- 17 ]MQi2#QQF],& 18 ]MQi2#QQFA/],]999jykN]- 19 ]MQi2#QQFSi?],]lb2`bf!#X+QKf+iBQMb]- 20 ]+Hmbi2`A/],]Rykd@kyk9ye@yyyyNNR] 21 ' 22 ' 23 ' Listagem 6: Ação associada à metainformação de uma operação executada numa tabela delta Ao considerar que um utilizador pretende criar uma transação que envolva a adição de uma nova coluna, com novos dados, a uma tabela delta , o sistema Delta Lake vai separar a transação em causa em dois segmentos distintos e, assim que a mesma termine, acrescenta as seguintes ações aos registos transacionais: 29
CAPÍTULO 2. ESTADO DE ARTE 1. Change metadata : alteração do esquema da tabela para incluir a nova coluna; 2. Add file : ação invocada por cada ficheiro adicionado. Tal como se pode observar no exemplo acima, as ações são ordenadas atomicamente segundo a transação requisitada por parte do utilizador. 2.5.2.3 Pontos de verificação O sistema Delta Lake , para obter um bom desempenho na execução da operação de leitura, procede à compactação periódica dos registos transacionais em pontos de verificação. Os pontos de verificação armazenam todas as ações não redundantes no estado da tabela até um determinado identificador de registo, no formato Parquet . Alguns conjuntos de ações são redundantes e podem ser removidos. A Figura 15 apresenta esses casos. 𝑎𝑑𝑑 (𝑥)𝑟𝑒𝑚𝑜𝑣𝑒(𝑥)𝑎𝑑𝑑 (𝑦)𝑟𝑒𝑚𝑜𝑣𝑒(𝑦)… 𝑎𝑑𝑑 (𝑥)𝑎𝑑𝑑 (𝑥)… Figura 15: Ações redundantes no sistema Delta Lake O resultado final do processo de verificação corresponde a um ficheiro Parquet que contém um registo de adição para cada objeto ainda presente na tabela. Mais, o mesmo também possui os registos de objetos que foram excluídos, cujo período de retenção ainda não expirou, e um pequeno número de registos, como txn , protocol e changeMetadata . Tal como a Figura 14 transparece, o ficheiro 000002.parquet representa um ponto de verificação com todos os registos anteriores, ou seja, deste o registo 000000.json até ao 000002.json , inclusive. Por omissão, o sistema Delta Lake realiza pontos de verificação a cada dez transações. Os ficheiros referidos com extensão Parquet encontram-se num formato ideal para consultar metadados sobre a tabela e para descobrir quais os objetos que podem conter dados relevantes para uma consulta seletiva com base nas suas estatísticas de dados. Este facto potencia o propósito desta dissertação, isto é, a adaptação do sistema Delta Lake na execução de transações frequentes e de granularidade fina para suportar cargas de trabalho OLTP . De modo a aceder eficientemente ao último ponto de verificação da tabela armazenada no sistema Delta Lake , os utilizadores podem encontrar essa informação no ficheiro _last_checkpoint , tal como a Figura 16 indica. 30
2.5. DELTA LAKE Figura 16: Último ponto de verificação de uma tabela delta 2.5.3 Protocolos de acesso Os protocolos de acesso do sistema Delta Lake foram concebidos de maneira a obter execuções de transações serializáveis sobre uma determinada tabela. Esta secção descreve as operações de leitura e de escrita presentes na realização de transações. 2.5.3.1 Leitura de uma tabela delta De modo a executar transações que apenas contêm operações de leitura de uma tabela pertencente ao sistema Delta Lake é necessário seguir diversos passos: 1. Ler o ficheiro _last_checkpoint para descobrir o último ponto de verificação concretizado (caso exista); 2. Listar todos os registos transacionais desde o último ponto de verificação (caso exista, se não existir considera-se o registo transacional com a versão inicial da tabela, isto é, a 0) até ao ficheiro JSON ou Parquet mais recente (isto é, o correspondente à última versão da tabela delta ). Esta operação é usada através do método listFrom ; 3. Usar os ficheiros do passo anterior para reconstruir o estado da tabela, nomeadamente aqueles que possuem a ação add sem haver uma ação remove associada; 4. Utilizar as estatísticas das ações presentes nos registos transacionais para identificar quais os objetos de dados a serem usados; 5. Executar a operação de leitura sobre os objetos de dados do passo anterior. 31
CAPÍTULO 2. ESTADO DE ARTE Caso um utilizador leia uma versão antiga do ponto de verificação presente no ficheiro _last_checkpoint , este consegue encontrar na mesma os registos transacionais na futura operação de listagem e, consequentemente, reconstruir o estado da tabela delta . Desta forma, um utilizador consegue tolerar inconsistência na listagem dos registos mais recentes ou na leitura dos objetos de dados referenciados nos registos transacionais. Tal como foi referido na Secção 2.5.2.3, o ficheiro _last_checkpoint ajuda apenas a reduzir o tempo de execução da operação listFrom ao disponibilizar o identificador do ponto de verificação mais recente. Assim, na existência de um número assinalável de registos, é possível aceder eficientemente ao último ponto de verificação de uma tabela delta . De modo a simplificar a interpretação do processo relativo à leitura de uma tabela delta , procedeu-se à esquematização do mesmo. A Figura 17 indica os passos necessários a uma realização válida deste tipo de operação na execução de uma determinada transação. Figura 17: Leitura de uma tabela delta 2.5.3.2 Escrita de um novo estado transacional numa tabela delta De forma a executar transações que incluem operações de escrita numa tabela delta , devem-se tomar os seguintes passos: 32
2.5. DELTA LAKE 1. Identificar o registo transacional mais recente, versão 𝑖(última versão da tabela delta ), usando os passos 1 e 2 referidos no processo de leitura. Após a transação efetuar a leitura dos registos transacionais procede-se à escrita do registo 𝑖+1; 2. Leitura dos dados da tabela (idêntico ao processo de leitura); 3. Escrita dos objetos de dados ( Parquet ) da transação em causa para a diretoria em questão. Uma vez terminada esta ação, os últimos estão prontos para serem referenciados no registo transacional que se pretende adicionar; 4. Escrita atómica do ficheiro JSON com a versão 𝑖+1no estado transacional da tabela. À semelhança da leitura de uma tabela delta , elaborou-se um esquema que sumaria o processo relativo à escrita de objetos de dados, em formato Parquet , e os respetivos registos transacionais, com extensão JSON . A Figura 18 expõe esse mesmo esquema. Figura 18: Escrita de um novo estado transacional numa tabela delta Para além dos quatro passos mencionados em cima, é também possível executar outro opcionalmente. Após a escrita do ficheiro JSON com o conteúdo da operação de escrita efetuada, existe a possibilidade de também adicionar nos registos transacionais o ponto de verificação, com formato Parquet , para 33
CAPÍTULO 2. ESTADO DE ARTE essa versão da tabela delta . Tal como foi referido na Secção 2.5.2.3, o sistema Delta Lake , por omissão, realiza pontos de verificação a cada dez transações. Contudo, este último valor pode ser configurado por parte do utilizador. Assim que a escrita do ficheiro em questão esteja concluída, é devidamente atualizado o ficheiro _last_checkpoint com o ponto de verificação relativo à versão 𝑖+1. De notar que este último passo apenas afeta o desempenho do sistema em causa e, na eventualidade da ocorrência de uma falha, os dados correspondentes não são corrompidos. Tal como seria de esperar, este passo opcional só é completado se o anterior for efetuado com sucesso. 34
Capítulo 3 Testes intermédios De forma a compreender na totalidade as particularidades do sistema Delta Lake [8], foram realizados alguns testes intermédios para ganhar uma maior intuição acerca do seu funcionamento. As experiências que se seguem foram implementadas tanto em Python [34] como na linguagem de programação Scala [29]. A primeira linguagem referida foi selecionada em detrimento doutras para a implementação do teste intermédio da Secção 3.3 devido à sua simplicidade, manutenção e legibilidade. Neste teste efetuouse a análise do tempo de execução das interrogações do benchmark TPC-H [32] sobre tabelas delta de diversas dimensões. A concretização deste teste intermédio permitiu conhecer o nível de desempenho do sistema Delta Lake na execução de interrogações com diferentes caraterísticas. Para tornar a codificação deste teste ainda mais percetível, foi utilizada a ferramenta JupyterLab [20] para, por exemplo, repetir a execução de blocos de código relevantes. Dado que o sistema Delta Lake é nativamente codificado em Scala , esta foi usada na implementação do testes intermédio que estuda a possibilidade da integração do estado transacional de uma tabela delta numa base de dados não relacional e no teste de desempenho que estabelece comparações entre o sistema Delta Lake e uma versão adaptada do último em MongoDB [28]. Deste modo, os testes intermédios presentes nas Secções 3.4 e3.5 foram especificados nesta linguagem de programação. Esta decisão foi tomada desta forma uma vez que para aceder ao estado transacional de uma tabela delta é necessário invocar um conjunto de instruções muito específicas que só estão disponíveis em Scala . Apesar desta realidade, também é possível reproduzir os mesmos resultados em Python utilizando instruções semelhantes às anteriores. Quanto ao ambiente computacional, os testes em questão são conduzidos localmente num computador com seis núcleos de processamento e 16 GB (Gigabyte) de memória RAM (Random Access Memory) . Em relação ao espaço de armazenamento, existe um disco SSD (Solid State Drive) NVMe (Non-Volatile 35
CAPÍTULO 3. TESTES INTERMÉDIOS Memory Express) com a capacidade de 512 GB . Dito isto, com a realização de múltiplas experiências sobre o sistema Delta Lake e com a aglomeração dos resultados obtidos, pretende-se estudar e validar a solução que será posteriormente proposta no Capítulo 4, onde é exposta a sua arquitetura e implementação. 3.1 Instalação Recorreu-se às ferramentas Spark (versão 3.1.2) [6], framework útil para o processamento distribuído, JupyterLab (versão 3.1.0), onde é implementado o primeiro teste intermédio, e aos sistemas Delta Lake (versão 1.0.0) e MongoDB (versão 5.0.1), que é importante para a persistência do estado transacional de uma tabela delta numa base de dados não relacional ( NoSQL ). Considerou-se ainda as linguagens Python e Scala , com as versões 3.9.6 e 2.12.14, respetivamente. Relativamente à primeira linguagem de programação, foram também instalados três pacotes distintos: PySpark [22] (versão 3.1.2), FindSpark [15] (versão 1.4.2) e PyMongo [21] (versão 3.12.0). Dado que Spark é nativamente implementado em Scala , existe a necessidade de usar um utensílio que estabeleça a ponte entre a própria e a linguagem de programação Python . Como tal, PySpark surgiu para tratar esta interação e representa, de forma sucinta, uma API , em Python , para Spark . Além disso, o último ajuda a fazer a interface com conjuntos de dados distribuídos e resilientes, também conhecidos por RDD , através da biblioteca Py4j . Esta biblioteca é bastante popular e também permite estabelecer uma interface dinâmica com objetos JVM (Java Virtual Machine) . A resiliência dos dados apontados devese ao facto dos mesmos serem regenerados caso sejam perdidos. Quanto à sua distribuição, estes são computados e guardados em diversos nós de armazenamento. Em suma, este pacote fornece uma API simples e abrangente com várias opções de visualização de dados, o que não se verifica nas linguagens Java ou Scala . Em relação ao pacote FindSpark este foi usado para localizar convenientemente a pasta de instalação da ferramenta Spark . Isto é bastante útil sobretudo no que diz respeito à configuração inicial de um notebook em JupyterLab onde se têm de especificar quais as ferramentas a serem utilizadas. Quanto ao pacote PyMongo , este proporciona a conexão entre um cliente Python e uma base de dados em MongoDB . No caso do teste relativo à Secção 3.3, este pacote foi também utilizado para lidar com a criação de uma base de dados concreta que, por sua vez, possui diversas coleções. No que toca à linguagem Scala , foi incluído um utensílio, equivalente ao que foi mencionado previamente, que estabelece uma conexão entre a ferramenta Spark e uma base de dados MongoDB [27]. Para além disso, foi ainda incorporado o pacote ScalaTest [23] que possibilita a especificação de testes unitários de modo a assegurar a funcionalidade e a qualidade do código desenvolvido. De modo a inspecionar e a visualizar o conteúdo dos objetos de dados, sob o formato Parquet [6], que se encontram disponíveis no estado transacional de uma tabela delta , foi também usado o utensílio parquet tools , na versão 1.12.0, para o efeito. 36
3.2. CONFIGURAÇÃO Por fim, foi usada uma outra ferramenta dedicada à realização de benchmarks , isto é, o TPC-H . O TPC-H é um benchmark desenhado para avaliar sistemas analíticos, isto é, sistemas que realizam operações de leitura sobre dados de variadas proporções que, por sua vez, pertencem a múltiplas bases de dados. Para o contexto deste tema de dissertação, foram utilizados tanto os conjuntos de dados disponibilizados pelo mesmo como as respetivas queries . Estes dados estão estruturados em diversas tabelas, com múltiplas relações e diversos graus de complexidade. O esquema do conjunto de dados em causa pode ser visualizado com mais detalhe no Apêndice Ado vigente documento. 3.2 Configuração A utilização do benchmark TPC-H pressupõe a escolha de um factor de escalonamento, isto é, a quantidade de dados gerada e que por sua vez será carregada no sistema de dados. Este fator, denominado por a6 , é aplicado à maioria das tabelas em questão para ir ao encontro com o valor representativo da dimensão do conjunto de dados especificado pelo utilizador. Durante a execução do benchmark , as interrogações incidem sobre os dados carregados, sendo que os resultados estão relacionados com a dimensão das relações. A Tabela 8expõe a cardinalidade de cada tabela presente no utensílio TPC-H em função de a6 , que traduz o fator mencionado. Nome da tabela Cardinalidade GAL1Ah1J a6 ×6000000 P_.1_a a6 ×1500000 S_halSS a6 ×800000 S_h a6 ×200000 *lahPJ1_ a6 ×150000 alSSGA1_ a6 ×10000 LhAPL 25 _1:APL 5 Tabela 8: Cardinalidade de cada tabela da ferramenta TPC-H em função do fator a6 3.3 Análise do tempo de execução das interrogações da ferramenta TPC-H No primeiro teste intermédio planeou-se o estudo e medição do tempo de execução de queries para diferentes dimensões de um conjunto de dados persistido em diversas tabelas delta . Para averiguar 37
CAPÍTULO 3. TESTES INTERMÉDIOS 3.6 Observações Atendendo aos resultados obtidos nos testes intermédios deste capítulo, é inequívoco que o sistema Delta Lake apresenta alguns compromissos na execução de operações que transformam o seu conteúdo. Com a realização do último teste (Secção 3.5), confirma-se que a adaptação em MongoDB é possível ser alcançada, ainda que a mesma tenha sido desenvolvida num ambiente experimental. Olhando exclusivamente para as Figuras 20 e21 pode-se afirmar que, em média, a versão adaptada em MongoDB é 90% mais rápida na inserção de uma linha e 84% mais rápida na remoção de uma linha em relação ao sistema Delta Lake , independentemente da dimensão da tabela considerada. Como tal, a integração do estado transacional de uma tabela delta numa base de dados não relacional ( NoSQL ), como é o caso do MongoDB , constitui-se como uma boa solução para a realização de transações que afetam as linhas de uma tabela deste tipo. Assim, este sistema seria adequado num ambiente computacional onde se realizam várias atualizações de granularidade fina, ou seja, onde ocorrem cargas de trabalho OLTP . 44
Capítulo 4 HyLake De modo a adaptar o sistema Delta Lake [8] à execução de transações frequentes e de granularidade fina em lagos de dados, desenhou-se um novo sistema, intitulado por HyLake , que assenta quer nas componentes do sistema Delta Lake quer nas caraterísticas de uma base de dados MongoDB [28]. Este novo sistema pode ser visto como uma solução híbrida onde os dados correspondentes à versão inicial da tabela considerada são armazenados em Delta Lake e as transformações são guardadas em MongoDB . Os dados com maiores proporções são armazenados na infraestrutura clássica do sistema Delta Lake enquanto que os dados mais recentes e de menores dimensões são guardados em MongoDB . A integração da base de dados MongoDB neste sistema deve-se não só aos resultados obtidos no teste intermédio da Secção 3.5 como também a outros fatores relevantes, nomeadamente a obtenção de uma escalabilidade horizontal, a inclusão de esquemas de dados flexíveis, o acesso a dados através de qualquer linguagem e a elevada capacidade de interrogar e analisar dados a diferentes níveis de complexidade. Com a leitura e escrita dos vários objetos delta numa base de dados não relacional procura-se a obtenção de um desempenho superior, permitindo mais transações de pequenas dimensões. Atendendo à primeira propriedade e ao facto do MongoDB ter sido concebido inicialmente como uma base de dados distribuída, esta adequa-se perfeitamente ao objetivo deste trabalho. A Figura 22 evidencia a arquitetura do sistema HyLake . Os objetos de dados representados com as cores azul e branca dizem respeito aos ficheiros codificados no formato Parquet [6]. Estes exprimem a primeira versão da tabela delta e, por sua vez, da tabela HyLake , pelo que permanecem armazenados no sistema Delta Lake . À semelhança do que acontece para os ficheiros Parquet , o ficheiro JSON que possui todos os registos transacionais relativos à criação da tabela também é persistido no sistema Delta Lake . Os ficheiros com extensão JSON são assinalados na Figura 22 com as cores verde e branca. Relativamente às transformações, estas são alojadas numa coleção de uma base de dados MongoDB com o mesmo nome da tabela delta . Ao contrário do que acontece no sistema Delta Lake , as alterações efetuadas 45
CAPÍTULO 4. HYLAKE sobre uma tabela HyLake dão origem a novos documentos com versões sucessivamente crescentes. Aos dois componentes presentes na arquitetura do sistema HyLake junta-se ainda a ferramenta Spark [6] que possibilita, por exemplo, a leitura de uma tabela delta . Esta permite o estabelecimento de uma conexão a uma base de dados MongoDB , promovendo a leitura e a escrita de dados. Como tal, considerou-se para o efeito o conector oficial do MongoDB [27]. No que toca à escrita de dados, este utensílio também é útil na criação e no armazenamento de dataframes numa base de dados orientada a documentos. Figura 22: Arquitetura do sistema HyLake Quanto aos pontos de verificação presentes no estado transacional de uma tabela delta , estes não são incorporados na base de dados em causa, uma vez que o objetivo principal deste sistema é obter um desempenho positivo no momento em que ocorrem múltiplas escritas de dados. O sistema Delta Lake adota este mecanismo para simplesmente reduzir o tempo da leitura do estado atual de uma tabela delta , pelo que esta estratégia não será considerada no sistema final. Contudo, na Secção 5.2, onde se aborda a avaliação da solução construída, será comparado o sistema HyLake original com uma versão do último onde existem pontos de verificação para provar que a primeira proposta é realmente aquela que se pretende utilizar. Na receção de um pedido de um utilizador, o sistema HyLake apresenta um comportamento distinto consoante a operação invocada. Como tal, é indispensável discutir nesta fase todos os detalhes 46
4.1. CRIAÇÃO DA TABELA relacionados com a criação, a atualização, a inserção, a remoção e, por último, a leitura de uma tabela HyLake . 4.1 Criação da tabela A criação de uma tabela no sistema HyLake reduz-se globalmente à instanciação de uma tabela delta . Em relação à definição da coleção em MongoDB que diz respeito à tabela HyLake , esta tarefa é tratada implicitamente pelo respetivo conector [27] e é apenas efetuada no momento em que ocorre a primeira transformação sobre o seu estado inicial. No que toca à sua implementação, a criação de uma tabela HyLake é trivial, bastando invocar o método disponibilizado pela API do sistema Delta Lake . A Figura 23 apresenta o fio de execução da criação de uma tabela HyLake . Figura 23: Fio de execução relativo à criação de uma tabela HyLake No início da génese de uma tabela HyLake são removidos todos os documentos presentes na respetiva coleção da base de dados MongoDB (1). A eliminação destes documentos assegura que não ocorra nenhuma sobreposição de dados entre o antigo e o novo estado transacional de uma tabela HyLake . Desta forma, quando a tabela em questão sofrer alterações, os documentos em causa são armazenados sequencialmente, ou seja, de acordo com a ordem com que as transformações são exercidas. Uma vez concluído este passo, procede-se à criação propriamente dita da tabela (2). Nesta etapa é gerado aleatoriamente o seu conteúdo de acordo com a dimensão especificada pelo utilizador. No fim deste 47
CAPÍTULO 4. HYLAKE processo, os dados produzidos são inseridos na tabela delta de modo a expressar a versão inicial da tabela HyLake . 4.2 Atualização e inserção de dados No que toca à atualização e inserção de dados, optou-se por aglomerar o comportamento das duas operações numa só, isto é, a operação upsert . De forma sucinta, no instante em que é invocada esta operação, é implicitamente verificado se o conteúdo que é passado como argumento existe ou não no estado atual da tabela HyLake . Se existir, procede-se naturalmente à atualização do seu conteúdo, caso contrário é executada a inserção desses mesmos dados. Como tal, surge uma nova abordagem na execução desta operação no sistema HyLake face ao que o sistema Delta Lake oferece. Esta modificação reside no facto do estado transacional de uma tabela delta ser agora persistido numa base de dados não relacional em MongoDB . Aquando da sua invocação, a operação upsert regista todas as alterações efetuadas no repositório de dados em questão, quer se trate de uma atualização ou de uma inserção de dados. De notar que apenas a informação dos ficheiros que possuem o formato JSON é guardada em MongoDB . Dado que se deseja transferir o conteúdo destes ficheiros para documentos, é indispensável estabelecer um modelo de armazenamento para cada documento alojado em MongoDB . A Listagem 10 evidencia o esquema de dados adotado para a operação upsert . Para simplificar a sua interpretação, exibe-se nesta listagem um exemplo de uma inserção de dados. 1& 2]p2`bBQM],R3]iBK2biKT],]6`B CM R Rk,yy,yy q1ah kykR]- 4]QT2`iBQM],]lSa1_h]- 5]QT2`iBQMS`K2i2`b],& 6]T`2/B+i2],]B/ AL UR-kV] 7'- 8]QT2`iBQMJ2i`B+b],& 9]MmKlT/i2/P`AMb2`i2/_Qrb],k 10 '- 11 ]//],( 12 &]B/],R-]+QHR],R'- 13 &]B/],k-]+QHR],k' 14 )- 15 ]`2KQp2], () 16 ' Listagem 10: Documento relativo à operação upsert do sistema HyLake Dos campos presentes nesta estrutura, destacam-se os seguintes três: version , add e remove . Tal como o próprio nome indica, o campo version indica o número da versão relativo a uma determinada 48
4.2. ATUALIZAÇÃO E INSERÇÃO DE DADOS tabela HyLake . Quanto aos campos add e remove , estes evidenciam as transformações exercidas sobre a mesma. Caso se trate de uma inserção de dados, o campo add é devidamente preenchido para refletir a nova modificação do estado da tabela. Na eventualidade de ocorrer uma atualização, são adicionados tanto ao campo add como ao campo remove as informações relativas a essa mesma operação. Com esta estratégia é possível realizar uma leitura válida do estado transacional de uma tabela HyLake em qualquer estágio, desde a sua génese. Relativamente aos restantes campos, estes possuem apenas metadados que complementam a informação da operação upsert . A Figura 24 resume o novo comportamento desta operação. Figura 24: Fio de execução relativo à operação upsert do sistema HyLake Na execução de uma query que envolva a inserção ou a atualização de dados, é invocada a operação upsert sobre o sistema HyLake . Inicialmente é determinado o estado atual da tabela HyLake para posteriormente colocar nos campos add e remove os dados referentes às linhas afetadas. Como tal, procede-se à leitura da tabela (1), sendo filtrados todos os registos cujas chaves encontram-se presentes no argumento passado ao método desta operação. Uma vez computadas as linhas da tabela que possuem estes identificadores, a estrutura de dados resultante é assinalada para remoção (2). Imediatamente a seguir a este cálculo é gerado deterministicamente um dataframe com as chaves referidas para serem devidamente anotadas para adição (2). A Listagem 11 demonstra este processo. 1pH `2KQp2.6 4?vGF2.6X7BHi2`Ub]B/ AL 0&B/bXKFai`BM;U]U]-]-]-]V]V']V 2pH //.6 4+`2i2_M/QK.i7`K2UB/b- ?vGF2.6X+QHmKMbXH2M;i?- BblTb2`i 4i`m2V Listagem 11: Computação do conteúdo a ser atualizado ou inserido numa tabela HyLake 49
CAPÍTULO 4. HYLAKE Após a criação dos dataframes referentes às linhas adicionadas e removidas, é construído um novo documento (3) segundo o esquema de dados apontado na Listagem 10 e os dados produzidos nos passos computacionais da Listagem 11. De modo a encerrar a execução desta operação, é persistido o documento em causa na respetiva coleção da base de dados MongoDB (4), incrementando não só a variável version , que diz respeito à versão atual da tabela, como a variável upserts , que espelha o número de execuções da operação upsert . Com estas alterações, antecipa-se uma melhoria substancial, em termos de desempenho, durante a execução de transações com as caraterísticas de um ambiente OLTP quando comparado com o sistema Delta Lake . O tratamento de ficheiros presentes no estado transacional de uma tabela delta é substituído pelo tratamento de documentos alojados numa coleção de uma base de dados em MongoDB . Com o aumento do número de transações que envolvem inserções ou atualizações, são muitos os ficheiros delta , Parquet e JSON , que iriam ser instanciados no sistema Delta Lake . Por conseguinte, observaria-se neste caso um pobre desempenho desta operação. 4.3 Remoção de dados Tratando-se de uma operação que modifica o estado atual de uma tabela HyLake , é novamente necessário especificar uma estrutura para o armazenamento de documentos relativos à operação de remoção de dados. A Listagem 12 aponta um exemplo concreto do modelo subjacente a esta operação. 1& 2]p2`bBQM],R3]iBK2biKT],]6`B CM R Rk,yy,yy q1ah kykR]- 4]QT2`iBQM],].1G1h1]- 5]QT2`iBQMS`K2i2`b],& 6]T`2/B+i2],]B/ AL UR-kV] 7'- 8]QT2`iBQMJ2i`B+b],& 9]MmK.2H2i2/_Qrb],k 10 '- 11 ]//], ()- 12 ]`2KQp2],( 13 &]B/],R-]+QHR],R'- 14 &]B/],k-]+QHR],k' 15 ) 16 ' Listagem 12: Documento relativo à remoção de dados presentes numa tabela HyLake Tal como se verificou na Listagem 10, o esquema de dados desta operação é muito semelhante ao da operação upsert . Todavia, existem algumas diferenças. A mais notória é sem dúvida ao campo add estar agora associada uma lista vazia para exprimir a essência da operação em causa. Quanto às 50
4.3. REMOÇÃO DE DADOS restantes modificações, estas residem na indicação da metainformação relativa à operação de remoção, incluindo, por exemplo, o número total de registos eliminados. Em relação ao campo remove é colocada a informação que diz respeito às linhas da tabela que são removidas para posteriormente realizar uma reconstituição correta do estado atual da mesma. Aquando da invocação desta operação, é adicionado um novo documento com a respetiva metainformação e com os dados a serem eliminados da tabela. A Figura 25 exibe o fio de execução da operação discutida. Figura 25: Fio de execução relativo à operação de remoção do sistema HyLake O caminho observado desde que a operação é iniciada até ao instante em que a última é finalizada não modifica os ficheiros Parquet presentes no sistema Delta Lake . Por conseguinte, para não comprometer o propósito desta operação, dá-se a preservação do estado inicial de uma tabela HyLake . Considerando o comportamento da operação de remoção do sistema HyLake , é crucial calcular o estado atual da tabela para indicar corretamente os dados referentes às linhas afetadas. Como tal, inicia-se a leitura da tabela (1), recolhendo-se no fim todas as linhas cujos identificadores encontram-se presentes no argumento passado à função em questão. A Listagem 13 exibe a filtragem presente no cálculo anterior. 1pH `2KQp2.6 4?vGF2.6X7BHi2`Ub]B/ AL 0&B/bXKFai`BM;U]U]-]-]-]V]V']V Listagem 13: Computação do conteúdo a ser removido numa tabela HyLake Assim que estes registos sejam totalmente determinados, a estrutura de dados resultante é assinalada para remoção (2). Após este procedimento, é elaborado o documento alusivo a esta operação (3) de acordo com o modelo de dados referido na Listagem 12 e os dados produzidos no passo computacional da Listagem 13. De maneira a concluir a execução desta operação, é armazenado o documento em causa 51
CAPÍTULO 4. HYLAKE na respetiva coleção da base de dados MongoDB (4), atualizando a versão atual da tabela em questão com mais uma unidade. 4.4 Leitura da tabela Atendendo à nova arquitetura proposta para as operações que transformam uma tabela HyLake , a leitura do estado transacional da mesma passa a ser totalmente distinta daquela que é definida por omissão no sistema Delta Lake . Com a escrita de registos transacionais numa coleção de uma base de dados não relacional ( NoSQL ), a operação de leitura requer uma adaptação de modo a acompanhar as modificações efetuadas. Este ajuste resume-se não só na leitura dos ficheiros Parquet presentes no estado transacional de uma tabela delta como também na leitura das alterações guardadas em MongoDB . Sendo uma operação de consulta, é requerido saber de antemão o formato com que os registos transacionais são persistidos na base de dados não relacional. Neste esquema de armazenamento, que é indicado explicitamente nas Listagens 10 e12, os campos version , add e remove são os únicos utilizados para realizar a leitura da tabela. Para cada versão desta estrutura de dados, é extraída a informação em causa, aplicando-a aos ficheiros Parquet que se encontram persistidos no sistema Delta Lake . Tendo em consideração não só o último esquema de dados como também o facto dos nomes das coleções presentes na base de dados MongoDB corresponderem aos respetivos nomes das tabelas delta , a operação de leitura é intuitiva a nível arquitetural. A Figura 26 sintetiza a forma como a operação de leitura é executada. Figura 26: Fio de execução relativo à leitura de uma tabela HyLake Com esta operação é possível aceder ao estado atual de uma tabela HyLake , independentemente 52
4.4. LEITURA DA TABELA da sequência de operações exercida sobre a última. Tal como se pode verificar nas Figuras 24 e25,a leitura de uma tabela HyLake está presente tanto na operação upsert como como no instante em que se dá a remoção de dados. Atendendo à importância desta operação no funcionamento do sistema HyLake , foram elaborados alguns testes unitários adicionais em Scala [29] para reforçar a qualidade da sua implementação. Aquando da invocação da operação de leitura, esta é automaticamente redirecionada tanto para o sistema Delta Lake como para a respetiva coleção de documentos em MongoDB . Numa primeira fase é lida a tabela delta em questão a partir dos seus ficheiros Parquet e do único documento JSON (1). A leitura desta tabela é obrigatória uma vez que a mesma traduz a versão inicial da tabela HyLake considerada. Caso não exista nenhuma transformação efetuada, o estado atual da tabela HyLake é exatamente igual ao da tabela delta , pelo que é devolvido o conteúdo dessa estrutura de dados sob a forma de um dataframe . Na eventualidade de haver alguma modificação persistida em MongoDB , são então acedidas e extraídas as mesmas com recurso a quatro operadores fundamentais: filter 1, explode 2, unionAll 3e, por fim, exceptAll 3. Tal como o próprio nome indica, o primeiro operador possibilita a filtragem dos documentos armazenados em MongoDB até uma certa versão da tabela HyLake . Assim, com o uso deste operador, é possível aceder os documentos pretendidos para posteriormente proceder-se à extração dos dados presentes nos campos add e remove (2). A Listagem 14 apresenta a filtragem mencionada. 1`//XiQ.6UVX7BHi2`U0]p2`bBQM] I4 p2`bBQMX;2iUVVX++?2UV Listagem 14: Filtragem de documentos de acordo com a versão atual da tabela HyLake No que toca à operação explode , esta permite expandir uma lista nos seus diversos elementos. No caso do sistema HyLake , esta operação é utilizada para expandir as linhas do dataframe persistido nos campos add e remove num determinado documento em MongoDB . A Listagem 15 exibe a aplicação desta operação. 1KQM;Q.6Xb2H2+iU2tTHQ/2U0]//]VXbU]//]VV Listagem 15: Desconstrução das linhas persistidas num documento em MongoDB Para tornar a interpretação da operação explode ainda mais clara, a Figura 27 ilustra um exemplo elucidativo. Tal como se pode observar nesta figura, o campo add é decomposto em dois segmentos 1?iiTb,ffbT`FXT+?2XQ`;f/Q+bfHi2bifTBfb[HfO7BHi2` 2?iiTb,ffbT`FXT+?2XQ`;f/Q+bfHi2bifTBfb[HfO2tTHQ/2 3?iiTb,ffbT`FXT+?2XQ`;f/Q+bfHi2bifb[H@`27@bvMit@[`v@b2H2+i@b2iQTbX?iKH 53
Capítulo 5 Avaliação De modo a avaliar o sistema HyLake procede-se à execução de diversos testes de desempenho com o intuito de o comparar não só com o sistema Delta Lake [8] como com a adaptação do último em MongoDB [28]. Em relação ao último sistema, este preserva exclusivamente o conteúdo das transações numa base de dados não relacional. Para atingir este propósito é utilizado o mesmo conector [27] que foi usado no sistema HyLake . O modelo de armazenamento subjacente a este conector corresponde precisamente ao conteúdo inerente às linhas que atualmente pertencem à tabela em causa. Esta informação equipara-se, por exemplo, aos dados alojados nos campos add e remove da Figura 10. Contudo, como se pode ver na Figura 29, cada linha armazenada na tabela dá origem a um documento diferente. Figura 29: Modelo de armazenamento do conector MongoDB - Spark 60
5.1. BENCHMARK Desta maneira, aquando da invocação de uma atualização, inserção ou remoção de linhas, o sistema MongoDB manipula apenas as linhas afetadas pela operação desencadeada. Ao contrário dos restantes sistemas, esta solução não guarda o estado transacional da tabela de uma forma iterativa, isto é, versão após versão. Em vez disso, o sistema MongoDB preserva apenas o estado atual da tabela em causa. Com a introdução desta nova solução é possível determinar qual destes três sistemas é o mais indicado à execução de transações frequentes e de granularidade fina para suportar uma carga de trabalho OLTP . 5.1 Benchmark Para obter uma comparação válida entre estes sistemas, foi construído um benchmark minimalista, em Scala [29], para os diferentes testes de desempenho. O surgimento deste benchmark deve-se ao facto da ferramenta TPC-H [32] apenas avaliar sistemas analíticos, ou seja, sistemas que realizam operações de leitura sobre dados de múltiplas dimensões e origens. Como para estudar o sistema HyLake é necessário medir o desempenho das suas escritas, o uso do benchmark TPC-H torna-se prescindível. Assim, optou-se pela elaboração de um novo benchmark que permita a leitura e a escrita de dados numa tabela HyLake . Para além disso, os resultados dos testes de desempenho deste utensílio podem ser guardados tanto num sistema de ficheiros tradicional como num sistema de ficheiros distribuído, como é o caso do HDFS (Hadoop Distributed File System) [2]. Por forma a encapsular e a abstrair o seu comportamento, foram instanciadas uma interface e uma classe para cada um dos três sistemas considerados. Em cada classe são especificados os testes de desempenho das operações de escrita e de leitura. Em ambas as operações dá-se a produção de dados de forma aleatória para injetar nas tabelas criadas. Esta geração de dados é feita através da ferramenta Spark [35] e obedece às restrições impostas pela configuração adotada no sistema HyLake , em particular a dimensão da tabela utilizada. Na medição do tempo de execução das operações de leitura e de escrita do sistema HyLake são ainda produzidos conjuntos aleatórios de identificadores numéricos. No caso da primeira operação, estes valores são passados como argumento tanto à operação upsert como à operação de remoção de dados de maneira a assegurar uma quantidade fixa de transformações numa determinada tabela HyLake . Quanto à segunda operação, o conjunto de identificadores gerado é passado como argumento à operação upsert para possibilitar a avaliação do desempenho das escritas do sistema HyLake . De modo a obter resultados fidedignos, no fim de cada teste realiza-se a média dos tempos de execução da operação pretendida. Quanto ao ambiente computacional, este benchmark executa os testes de desempenho em contextos distintos. Os testes executados localmente inserem-se no mesmo ambiente que foi descrito no início do Capítulo 3. No que toca aos testes remotos, estes enquadram-se no ambiente que é caraterizado na Secção 5.3. A fim de avaliar corretamente o comportamento de tais soluções, as configurações dos testes de desempenho são iguais entre os três sistemas. Nela é especificada, por exemplo, o número de tentativas 61
CAPÍTULO 5. AVALIAÇÃO de execução por teste. Quanto aos restantes parâmetros, estes variam de acordo com o ambiente computacional usado, isto é, local ou remoto. A Tabela 12 apresenta a única configuração comum entre os testes de desempenho concebidos. Variável Valor Número de tentativas de execução por teste 20 Tabela 12: Configuração global dos testes de desempenho 5.2 HyLake O primeiro teste de desempenho procura confirmar se de facto a inclusão de pontos de verificação no sistema HyLake é ou não benéfica na execução de transações em ambientes que utilizam cargas de trabalho OLTP . Tal como foi mencionado na Secção 4.4, foram elaboradas quatro alternativas distintas para a operação de leitura de uma tabela relativa ao sistema HyLake . Consequentemente, produzem-se quatro formas diferentes de executar a reconstrução do estado transacional da mesma. Para facilitar a identificação destas propostas na legenda dos respetivos gráficos, estas são numeradas entre um a quatro da seguinte forma: • Alternativa 1: alternativa mencionada na Secção 4.4.1 onde ocorre a aplicação simultânea de transformações usando a linguagem de programação Scala e a ferramenta Spark ; • Alternativa 2: proposta referida na Secção 4.4.2 onde se dá a aplicação sucessiva de transformações utilizando a linguagem de programação Scala e e a ferramenta Spark ; • Alternativa 3: solução exposta na Secção 4.4.3 onde se realiza a aplicação simultânea de transformações usando a notação MongoDB aggregation pipeline ; • Alternativa 4: alternativa apontada na Secção 4.4.4 onde acontece a aplicação sucessiva de transformações utilizando a notação MongoDB aggregation pipeline . Para realizar esta experiência, desenhou-se uma nova vertente do sistema HyLake que inclui pontos de verificação no seu estado transacional. Nela, à semelhança do que se sucede no sistema Delta Lake , é assinalada, de dez em dez transações, a existência de pontos de verificação que, por sua vez, resumem o estado transacional até uma determinada versão. De modo a identificar a ocorrência de um ponto de verificação no sistema HyLake , basta observar se a versão correspondente a um determinado documento guardado em MongoDB é ou não múltiplo de dez. Para tornar ainda mais evidente o seu reconhecimento, 62
5.2. HYLAKE é acrescentado um novo campo, com o nome checkpoint , para informar a presença de um ponto de verificação. A Listagem 17 exibe um exemplo do novo modelo de armazenamento adotado na execução da operação upsert do sistema HyLake que inclui pontos de verificação. Nesta listagem destaca-se a existência do campo checkpoint e o facto do número da versão do documento alusivo a esta transformação ser múltiplo de dez. 1& 2]+?2+FTQBMi],i`m23]p2`bBQM],Ry4]iBK2biKT],]6`B CM R Rk,yy,yy q1ah kykR]- 5]QT2`iBQM],]lSa1_h]- 6]QT2`iBQMS`K2i2`b],& 7]T`2/B+i2],]B/ AL UR-kV] 8'- 9]QT2`iBQMJ2i`B+b],& 10 ]MmKlT/i2/P`AMb2`i2/_Qrb],k 11 '- 12 ]//],( 13 &]B/],R-]+QHR],R'- 14 &]B/],k-]+QHR],k' 15 )- 16 ]`2KQp2], () 17 ' Listagem 17: Modelo de armazenamento do sistema HyLake com pontos de verificação Em relação às configurações inicialmente apontadas na Tabela 12 foram adicionados dois novos parâmetros. O primeiro define o número total de transformações efetuadas sobre a tabela HyLake . A esta configuração junta-se também a especificação de uma dimensão fixa para a tabela gerada de modo a permitir uma análise coerente do comportamento do sistema HyLake . A Tabela 13 indica as configurações referidas. Variável Valor Número de transformações 1, 5, 10, 25, 50 Dimensão da tabela 100 linhas ×10 colunas Tabela 13: Configuração dos testes de desempenho executados num ambiente computacional local 5.2.1 Operação de leitura A Figura 30 evidencia o comportamento das diferentes alternativas do sistema HyLake sem pontos de verificação ao executar a operação de leitura. Olhando atentamente para esta figura, verifica-se um 63
CAPÍTULO 5. AVALIAÇÃO comportamento divergente entre os dois pares de alternativas deste sistema. Relativamente à primeira e terceira alternativas, é clara a sua superioridade em termos de desempenho quando comparado com as restantes. Tal como foi referido nas Secções 4.4.1 e4.4.3, estas duas alternativas aplicam de uma só vez as transformações efetuadas sobre uma tabela HyLake aos ficheiros Parquet [6] que exprimem a sua primeira versão, de forma a reconstruir corretamente o seu estado transacional. A eficiência da última tarefa nestas alternativas é nitidamente evidenciada na Figura 30. O comportamento de ambas é praticamente idêntico, isto é, para um número de transformações sucessivamente maior, o tempo de execução mantém-se reduzido e constante. Já para a segunda e quarta alternativas isto já não acontece, ou seja, o tempo de execução cresce à medida que a quantidade de modificações aumenta. Este desempenho deve-se ao facto da aplicação referida anteriormente ser feita alteração após alteração. 1510 25 50 0 1 2 3 4 5 6 7 Número de transformações Tempo de execução (em segundos) HyLake - Alternativa 1 HyLake - Alternativa 2 HyLake - Alternativa 3 HyLake - Alternativa 4 Figura 30: Teste de desempenho relativo à leitura de uma tabela HyLake sem pontos de verificação Com esta experiência, é irrefutável a escolha da primeira e terceira alternativas em detrimento da segunda e quarta alternativas. Desta forma, a partir da Secção 5.3 apenas são consideradas estas alternativas para se realizar a comparação com os diferentes sistemas. Passando para o sistema HyLake com pontos de verificação, é possível averiguar na Figura 31 praticamente o mesmo comportamento para as alternativas consideradas previamente. No entanto, para a segunda e quarta alternativas observa-se uma oscilação em termos de desempenho consoante o número de transformações. Isto ocorre tanto pelos mesmos motivos do sistema anterior como pela ocorrência de pontos de verificação de dez em dez modificações. Por conseguinte, a leitura do estado transacional de uma tabela HyLake nesse instante é um pouco mais rápida. Para os restantes casos, o tempo de 64
5.2. HYLAKE execução é superior uma vez que é necessário executar não só a leitura desse mesmo estado, que é potencialmente cada vez maior, como as modificações que foram exercidas sobre a tabela. 1510 25 50 0 10 20 30 40 50 60 Número de transformações Tempo de execução (em segundos) HyLake - Alternativa 1 HyLake - Alternativa 2 HyLake - Alternativa 3 HyLake - Alternativa 4 Figura 31: Teste de desempenho relativo à leitura de uma tabela HyLake com pontos de verificação Analisando o desempenho das quatro alternativas da leitura de uma tabela HyLake com e sem pontos de verificação, confirma-se que a aplicação sucessiva de transformações tem um impacto negativo na reconstrução do seu estado transacional. Por este motivo, a segunda e a quarta alternativas desta operação não são ponderadas no sistema HyLake final. Quanto à primeira e terceira alternativas, estas apresentam um comportamento similar nas duas propostas do último sistema. Dado que a atualização, a inserção e a remoção de dados dependem diretamente da operação de leitura e que, no sistema HyLake com pontos de verificação, é escrito de dez em dez transações o estado completo da tabela, a adoção da última proposta prejudicaria a execução de escritas. Assim, na avaliação final deste projeto escolheu-se apenas a primeira e a terceira alternativas do sistema HyLake que não inclui pontos de verificação. 5.2.2 Operação de escrita A Figura 32 apresenta o comportamento das diferentes alternativas do sistema HyLake sem pontos de verificação ao executar a operação de escrita. Nela ainda são exibidas a segunda e quarta alternativas para comprovar o péssimo desempenho que estas têm quando comparadas com as restantes. Basta por exemplo visualizar que o tempo dispendido na execução de dez transformações é idêntico aos resultados 65
CAPÍTULO 5. AVALIAÇÃO da primeira e terceira alternativas para cinquenta modificações. Às quatro alternativas mencionadas juntase também a primeira versão do sistema HyLake com pontos de verificação, versão essa que se encontra representada com uma linha a tracejado. Tal como se pode observar na Figura 32, esta possui um pior comportamento do que a mesma versão no sistema HyLake sem pontos de verificação, validando o que foi dito no fim da Secção 5.2.1. 1510 25 50 0 10 20 30 40 50 60 70 Número de transformações Tempo de execução (em segundos) HyLake - Alternativa 1 sem pontos de verificação HyLake - Alternativa 1 com pontos de verificação HyLake - Alternativa 2 sem pontos de verificação HyLake - Alternativa 3 sem pontos de verificação HyLake - Alternativa 4 sem pontos de verificação Figura 32: Teste de desempenho relativo à operação de escrita do sistema HyLake Olhando exclusivamente para a primeira e terceira alternativas, cujas linhas estão coloridas a roxo e a vermelho respetivamente, pode-se afirmar que a primeira é aquela que exibe a maior eficiência na execução de escritas. A única diferença que separa estas implementações é a linguagem utilizada, ou seja, a primeira usa a linguagem Scala , com o auxílio da ferramenta Spark , enquanto que a terceira usufrui a notação MongoDB aggregation pipeline . Dada a proximidade entre estas últimas duas versões em termos de desempenho, não é totalmente claro qual delas é que deve ser utilizada na comparação efetuada na Secção 5.3. Com a existência desta incerteza, são ainda consideradas nesse segmento ambas as versões. 5.3 Delta Lake vs HyLake vs MongoDB Com a escolha definitiva da versão final do sistema HyLake , avança-se para a realização dos testes de desempenho das operações de leitura e escrita para os três sistemas referidos neste capítulo. À semelhança do que se sucedeu na Secção 5.2, os sistemas Delta Lake , HyLake e MongoDB são sujeitos 66
5.3. DELTA LAKE VS HYLAKE VS MONGODB aos mesmos testes. Estes são executados tanto localmente como remotamente, pelo que são tomadas diferentes configurações. Quanto aos testes locais, estes adotam as mesmas configurações que foram apresentadas nas Tabelas 12 e13. Dado que na Secção 5.2 já foram determinadas as versões mais eficientes do sistema HyLake , estas são unicamente utilizadas neste segmento para efeitos de comparação, restando apenas medir o desempenho dos sistemas Delta Lake e MongoDB . Os resultados destes testes são expostos nas Figuras 33 e35. No que toca à execução dos testes de desempenho remotos, estes são conduzidos num cluster pertencente à Google Cloud Platform [11]. Ao contrário do que se verificou nos testes locais, neste ambiente computacional são avaliados todos os sistemas em causa. O cluster mencionado é especificado através do serviço DataProc [9] e é constituído por um nó master e dois nós workers para permitir uma taxa significativa de execuções transacionais. Estes três nós correspondem a máquinas da série N1 (n1-standard-2) com dois núcleos de processamento e 7,5 GB de memória RAM . Em relação ao espaço de armazenamento, cada uma das instâncias possui cerca de 50 GB . Para além destas três instâncias, é adicionada uma outra máquina que contém uma base de dados não relacional, isto é, MongoDB . Esta máquina partilha a mesma série das instâncias mencionadas previamente, sendo que usufrui apenas um CPU (Central Processing Unit) e 3,75 GB de memória RAM (n1-standard-1). Ao contrário do que se constatou nos três nós do cluster , esta instância requer uma quantidade reduzida de armazenamento persistente, isto é, cerca de 10 GB . Após a realização de múltiplas experiências, sintetizam-se nas Figuras 34 e36 os resultados obtidos no presente contexto. Por fim, expõem-se ainda na Tabela 14 as configurações exclusivas a este ambiente. Variável Valor Número de linhas afetadas 10 Número de linhas da tabela 1000, 10000, 100000, 1000000 Número de colunas da tabela 10 Tabela 14: Configuração dos testes de desempenho executados num ambiente computacional remoto Em relação ao parâmetro relativo ao número de linhas afetadas, este é fixado num único valor para facilitar a interpretação dos resultados obtidos. Em suma, aquando da invocação de uma transformação, é executada uma atualização, inserção ou remoção de dados que afeta a quantidade de linhas especificada. Isto é conseguido através da identificação das chaves que determinam univocamente cada linha presente numa determinada tabela. 5.3.1 Operação de leitura Para iniciar a execução dos testes que dizem respeito à operação de leitura, são criadas tabelas com uma dimensão fixa ou variável para os diversos sistemas, consoante o ambiente computacional usado. 67
CAPÍTULO 5. AVALIAÇÃO As tabelas são preenchidas de acordo com os parâmetros de configuração desse mesmo ambiente. Dado que nos testes remotos a dimensão das tabelas é superior em relação ao que se verifica nos testes locais, é imprescindível evitar a criação de uma tabela completa ao longo da execução para refletir uma certa quantidade de transformações. Para solucionar este problema eficientemente, são calculadas as diferenças das quantidades das modificações especificadas. Desta forma, assim que a tabela seja instanciada, não existe a repetição da execução de transformações. Aquando da finalização deste processo, dá-se início à leitura propriamente dita, iterando-se no fim o conteúdo da tabela usada. Olhando em primeiro lugar para a Figura 33, onde são exibidos os resultados da execução local da operação de leitura, é possível desde logo destacar a superioridade, em termos de desempenho, dos sistemas Delta Lake e MongoDB . A leitura de uma tabela referente ao sistema MongoDB resume-se exclusivamente à leitura de documentos que exprimem o conteúdo atual da última. Tal como foi dito no início do Capítulo 5, cada linha pertencente à tabela dá origem a um documento distinto, pela que a sua leitura é efetuada rapidamente. Quanto ao sistema Delta Lake , este apresenta um comportamento bastante satisfatório na execução da operação de leitura, sobretudo à medida que o número de transformações aumenta. Ainda que esta solução utilize ficheiros para guardar o estado transacional de uma tabela, esta consegue obter um melhor desempenho do que o do sistema HyLake . O comportamento do último deve-se ao facto da operação de leitura se traduzir na leitura quer dos ficheiros Parquet guardados numa tabela delta quer dos documentos armazenados em MongoDB . Consequentemente, observa-se um pobre desempenho. 1510 25 50 0 0.1 0.2 0.3 0.4 0.5 0.6 0.7 0.8 Número de transformações Tempo de execução (em segundos) Delta Lake HyLake - Alternativa 1 HyLake - Alternativa 3 MongoDB Figura 33: Teste de desempenho relativo à execução local da operação de leitura 68
5.3. DELTA LAKE VS HYLAKE VS MONGODB Passando para a análise da Figura 34, onde são evidenciados os resultados da execução da operação de leitura num ambiente remoto, pode-se observar um comportamento totalmente distinto para o sistema MongoDB . Este contraste em termos de desempenho acontece pelo aumento significativo do número de linhas da tabela considerada. Dado que no sistema MongoDB cada linha presente na tabela dá origem a um documento distinto, produzem-se muitos documentos na respetiva base de dados. No caso deste teste chega-se a atingir a geração de um milhão de documentos. Por conseguinte, a leitura de uma tabela pertencente ao sistema MongoDB torna-se cada vez mais lenta. Tal como se pode ver na Figura 34 e regressando ao exemplo anterior, a leitura de uma tabela com um milhão de linhas demora quase 18 segundos. Em relação ao sistema Delta Lake , este exprime o melhor comportamento na execução da operação de leitura, demonstrando o seu potencial em ambientes analíticos. Ainda assim, para uma tabela com um milhão de linhas, esta solução é ultrapassada, em termos de desempenho, pelo sistema HyLake . Tratando-se de um sistema que adota um comportamento híbrido, este consegue manter um desempenho positivo, independentemente da dimensão da tabela considerada. Como tal, o sistema HyLake apresenta um bom comportamento em ambientes transacionais que envolvem cargas de trabalho mistas. 1000 10000 100000 1000000 0 2 4 6 8 10 12 14 16 18 Número de linhas da tabela Tempo de execução (em segundos) Delta Lake HyLake - Alternativa 1 HyLake - Alternativa 3 MongoDB Figura 34: Teste de desempenho relativo à execução da operação de leitura num ambiente remoto Assim, para cargas de trabalho onde ocorrem leituras de tabelas com uma dimensão reduzida, o sistema Delta Lake é o mais apropriado e, para um número significativo de linhas, deverá optar-se pela utilização do sistema HyLake . 69
BIBLIOGRAFIA [13] IBM Docs . Inglês. 9 de set. de 2021. url: ?iiTb,ffrrrXB#KX+QKf/Q+bf2Mf/#kfRRX8 (ver pp. 3,8). [14] Mashamsft. SQL Server technical documentation - SQL Server . Inglês. 9 de set. de 2021. url: ?iiTb,ff/Q+bXKB+`QbQ7iX+QKf2M@mbfb[Hfb[H@b2`p2`f\pB2r4b[H@b2`p2`@ p2`R8 (ver pp. 3,8). [15] minrk. findspark . Inglês. 9 de set. de 2021. url: ?iiTb,ff;Bi?m#X+QKfKBM`Ff7BM/bT`F (ver p. 36). [16] MySQL :: MySQL 8.0 Reference Manual . Inglês. 9 de set. de 2021. url: ?iiTb,ff/2pXKvb[HX +QKf/Q+f`27KMf3Xyf2M (ver pp. 3,8). [17] Oracle Database 21c - Get Started . Inglês. 8 de set. de 2021. url: ?iiTb,ff/Q+bXQ`+H2X +QKf2Mf/i#b2fQ`+H2fQ`+H2@/i#b2fkRfBM/2tX?iKH (ver pp. 3,8). [18] PostgreSQL 13.4 Documentation . Inglês. 12 de ago. de 2021. url: ?iiTb,ffrrrXTQbi;`2b[HX Q`;f/Q+bf+m``2Mi (ver pp. 3,8). [19] Press Release - Domo Releases Eighth Annual “Data Never Sleeps” Infographic | Domo . Inglês. 9 de set. de 2021. url: ?iiTb,ffrrrX/QKQX+QKfM2rbfT`2bbf/QKQ@`2H2b2b@2B;?i?@ MMmH@/i@M2p2`@bH22Tb@BM7Q;`T?B+ (ver p. 7). [20] Project Jupyter . 5 de ago. de 2021. url: ?iiTb,ffDmTvi2`XQ`; (ver p. 35). [21] PyMongo 3.12.0 Documentation — PyMongo 3.12.0 documentation . Inglês. 15 de jul. de 2021. url: ?iiTb,ffTvKQM;QX`2/i?2/Q+bXBQf2Mfbi#H2 (ver p. 36). [22] PySpark Documentation — PySpark 3.1.2 documentation . 27 de mai. de 2021. url: ?iiTb,ff bT`FXT+?2XQ`;f/Q+bfHi2bifTBfTvi?QMfBM/2tX?iKH (ver p. 36). [23] ScalaTest . 9 de set. de 2021. url: ?iiTb,ffrrrXb+Hi2biXQ`; (ver p. 36). [24] Serviços de Computação na Cloud | Microsoft Azure . Português. 9 de set. de 2021. url: ?iiTb, ffxm`2XKB+`QbQ7iX+QKfTi@Ti (ver p. 2). [25] Serviços de nuvem – Amazon Web Services (AWS) . Português. 7 de set. de 2021. url: ?iiTb, ffrbXKxQMX+QKfTi (ver p. 2). [26] Social Media Statistics: Top Social Networks by Popularity . Inglês. 28 de mai. de 2021. url: ?iiTb, ff/mbiBMbiQmiX+QKfbQ+BH@K2/B@biiBbiB+b (ver p. 1). [27] Spark Connector Scala Guide — MongoDB Spark Connector . Inglês. 14 de jun. de 2021. url: ?iiTb, ff/Q+bXKQM;Q/#X+QKfbT`F@+QMM2+iQ`f+m``2Mifb+H@TB (ver pp. 36,39,46, 47,60). [28] The most popular database for modern apps . Inglês. 9 de set. de 2021. url: ?iiTb,ffrrrX KQM;Q/#X+QK (ver pp. 4,35,45,60,72). 76
BIBLIOGRAFIA [29] The Scala Programming Language . 9 de set. de 2021. url: ?iiTb,ffrrrXb+H@HM;XQ`; (ver pp. 19,35,53,61). [30] Top 20 Big Data Statistics for 2020 | Sigma Computing . Inglês. 6 de jul. de 2021. url: ?iiTb, ffrrrXbB;K+QKTmiBM;X+QKf#HQ;fiQT@ky@#B;@/i@biiBbiB+b@7Q`@kyky (ver p. 6). [31] Total data volume worldwide 2010-2025 | Statista . Inglês. 9 de set. de 2021. url: ?iiTb,ffrrrX biiBbiX+QKfbiiBbiB+bf3dR8RjfrQ`H/rB/2@/i@+`2i2/ (ver p. 6). [32] TPC-H Homepage . Inglês. 9 de set. de 2021. url: ?iiT,ffrrrXiT+XQ`;fiT+? (ver pp. 35,61, 73,78). [33] Welcome to Apache Avro! 17 de mar. de 2021. url: ?iiTb,ffp`QXT+?2XQ`; (ver p. 14). [34] Welcome to Python.org . Inglês. 9 de set. de 2021. url: ?iiTb,ffrrrXTvi?QMXQ`; (ver pp. 19, 35,57). [35] M. Zaharia et al. “Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing”. Em: 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12) . San Jose, CA: USENIX Association, abr. de 2012, pp. 15–28. url: ?iiTb,ffrrrX mb2MBtXQ`;f+QM72`2M+2fMb/BRkfi2+?MB+H@b2bbBQMbfT`2b2MiiBQMfx?`B (ver pp. 3,19,61). 1bi2/Q+mK2MiQ7QB;2`/QmiBHBxM/QQT`Q+2bb/Q`UT/7fs2fGmVG h 1 s-+QK#b2MQi2KTHi2LPoi?2bBb-/2b2MpQHpB/QMQ.2TXAM7Q`KiB+/6*hĜLPoTQ`CQ½QJXGQm`2MÏQX(R) (R) CXJXGQm`2MÏQXh?2LPoi?2bBbG h 1sh2KTHi2lb2`ǶbJMmHXLPolMBp2`bBivGBb#QMXkykRXl_G,?iiTb,ff;Bi?m#X+QKfDQQKHQm`2M+QfMQpi?2bBbf`rfKbi2`fi2KTHi2XT/7Up2`TXddVX 77
Apêndice A Apêndice A Figura 37 evidencia o esquema do conjunto de dados da ferramenta TPC-H [32]. Nesta ilustração são exibidas as chaves primárias e estrangeiras de cada tabela, sendo que sob cada uma das mesmas é apresentada a respetiva cardinalidade. Esta última informação pode ser consultada na Tabela 8. Figura 37: Esquema do conjunto de dados da ferramenta TPC-H 78