antonio leandro

enciclopédia

sistemas distribuídos

coordenação, consenso e escala: os papers do google e da amazon que viraram a infraestrutura de todo mundo.

33 verbetes, 17 no caminho mínimo, 22 lidos na íntegra · do que mais pesa para o que menos

muda como você pensa

  1. Time, Clocks, and the Ordering of Events in a Distributed Systemnum sistema distribuído não existe "a" ordem dos eventos: existe a ordem parcial do que pôde causar o quê, e um contador por processo mais um timestamp na mensagem bastam para estendê-la a uma ordem total qualquer · núcleo
  2. Paxos Made Simpleconsenso cabe em duas rodadas de mensagem: maiorias sempre se cruzam, então basta cada proposta nova prometer carregar o valor da proposta mais alta já aceita, e o resto é eleger um líder para não travar · núcleo
  3. The Google File Systemfalha de hardware é o regime normal, não o acidente: se o sistema assume isso e a aplicação aceita um append "pelo menos uma vez" no lugar de consistência forte, um master único e discos baratos aguentam centenas de terabytes. · núcleo
  4. MapReduce: Simplified Data Processing on Large Clustersrestringir o programador a duas funções, map e reduce, é o que paga a conta: se a computação é determinística sobre a entrada, re-executar a tarefa vira toda a tolerância a falha que o cluster precisa. · núcleo
  5. Bigtable: A Distributed Storage System for Structured Dataum mapa ordenado, esparso e versionado, com atomicidade só dentro da linha, sustenta sessenta produtos: o que o cliente ganha não é schema relacional, é controle sobre o layout físico dos dados · núcleo
  6. The Chubby lock service for loosely-coupled distributed systemsum serviço central de locks grosseiros resolve eleição de líder melhor que uma biblioteca de paxos — não por ser mais correto, mas porque cabe em duas linhas de código num sistema que já existe · núcleo
  7. Spanner: Google's Globally-Distributed Databaseexpor a incerteza do relógio em vez de escondê-la: se o banco sabe quanto pode estar errado sobre a hora, ele espera essa janela passar e entrega transação linearizável em escala planetária · núcleo
  8. Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computingdá para ter memória distribuída tolerante a falha sem replicar nada: se a escrita só acontece em bloco, o sistema pode guardar a receita da computação e recomputar só a partição que caiu · núcleo
  9. The Tail at Scalea cauda da latência é problema de confiabilidade, não de otimização: quando uma requisição consulta milhares de servidores, o lento eventual vira o comportamento normal — e dá para tolerá-lo como se tolera falha · núcleo · pelo resumo
  10. In Search of an Understandable Consensus Algorithmcompreensibilidade dá para tratar como requisito de projeto: com líder forte e menos estado não-determinístico, raft entrega o mesmo que o multi-paxos usando só quatro tipos de mensagem · núcleo
  11. Large-scale cluster management at Google with Borgutilização alta num datacenter não vem do escalonador esperto: vem de admitir menos, empacotar melhor, prometer mais recurso do que existe e botar serviço e batch na mesma máquina · núcleo · pelo resumo
  12. Zanzibar: Google's Consistent, Global Authorization Systempermissão vira uma tupla só — objeto#relação@usuário, em que o usuário pode ser outra tupla — e todo o resto é regra de reescrita; o que segura a coisa é o timestamp que o cliente guarda junto com o conteúdo · núcleo
  13. Millions of Tiny Databasesconsistência forte e disponibilidade sob partição só brigam se você exigir disponibilidade para todos os clientes; cada chave precisa viver em três pontos da rede, então não construa um banco — construa milhões · núcleo
  14. Firecracker: Lightweight Virtualization for Serverless Applicationsdá para ter isolamento de vm com overhead de container: trocar o qemu por 50 mil linhas de rust rende cerca de 3mb de overhead por microvm e boot na casa de 125ms — é isso que roda o lambda · núcleo

vale o tempo

  1. Brewer's Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Servicesquando a rede pode perder mensagem, nenhum algoritmo entrega consistência atômica e disponibilidade ao mesmo tempo — e no modelo assíncrono ele falha até nas execuções em que nada se perde · núcleo
  2. Web Search for a Planet: The Google Cluster Architectureconfiabilidade é um problema de software: com réplica e balanceamento, 15.000 pcs de prateleira entregam mais busca por dólar do que um punhado de servidores caros com fonte redundante e raid
  3. Dapper, a Large-Scale Distributed Systems Tracing Infrastructurerastrear um sistema distribuído inteiro sai barato se você desistir de rastrear tudo: uma amostra de 1 em 1.024 requisições, colhida dentro das bibliotecas de RPC e threading, basta — e o autor da aplicação nem fica sabendo · núcleo
  4. Pregel: A System for Large-Scale Graph Processingalgoritmo de grafo é iterativo, e o mapreduce cobra o grafo inteiro em disco a cada iteração — a resposta do google não foi otimizar o framework genérico, foi construir um sistema só para grafo · pelo resumo
  5. Kafka: a Distributed Messaging System for Log Processinglog é a abstração certa para dado de evento: com o broker sem estado de entrega e o offset guardado no consumidor, a vazão para de depender de quanto dado já está acumulado no disco · núcleo
  6. Megastore: Providing Scalable, Highly Available Storage for Interactive Servicesdá para ter acid de verdade com replicação síncrona entre datacenters, desde que você aceite quebrar o banco em milhões de bancos minúsculos e só prometer transação dentro de cada pedaço
  7. Challenges to Adopting Stronger Consistency at Scaleconsistência forte não trava por teoria: trava porque o mesmo dado vive em centenas de serviços com sharding diferente, e aí um único nó lento vira latência de todo mundo
  8. The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processingparar de tratar fluxo infinito como lote que um dia fica completo: o sistema deve assumir que nunca saberá se viu todo o dado, e entregar ao engenheiro a escolha explícita entre correção, latência e custo · pelo resumo
  9. Real-time Data Infrastructure at Uberadotar kafka, flink e pinot é a parte barata: o que sustenta o tempo real na uber é a camada de indireção construída em cima deles — federação, proxy, sql — porque é ela que absorve dados, casos de uso e usuários

para aprofundar

  1. Epidemic Algorithms for Replicated Database Maintenanceréplica não precisa de entrega garantida nem de estrutura de controle: se cada site conversa com um parceiro aleatório de vez em quando, a atualização contagia a rede inteira em tempo logarítmico
  2. Large-scale Incremental Processing Using Distributed Transactions and Notificationso índice da web não precisa ser reconstruído em lote: se você tiver transações distribuídas e notificações sobre um repositório de petabytes, dá para atualizar documento por documento e cortar a idade do índice pela metade · pelo resumo
  3. Tenzing: A SQL Implementation On The MapReduce Frameworksql quase completo em cima do mapreduce não é gambiarra de analista: em 2011 isso já servia 10.000 consultas ad hoc por dia sobre 1,5 petabyte comprimido, sem trocar o motor de execução do cluster · pelo resumo
  4. F1: A Distributed SQL Database That Scalesa escolha entre escalar como um nosql e manter sql com consistência forte era falsa: dá para ter os dois desde que você aceite pagar a conta em latência de commit e desenhar o schema para reduzi-la · pelo resumo
  5. Web-scale Job Schedulingdatacenter da web não é supercomputador grande: em 2013 ele já podia ser maior que qualquer máquina de hpc e rodar carga mais bagunçada, e mesmo assim quase ninguém tinha escrito sobre escalonar aquilo · pelo resumo
  6. Optimizing Google's Warehouse Scale Computers: The NUMA Experienceem máquina de datacenter a memória não é uniforme: onde a página mora em relação ao core decide o desempenho, e num warehouse-scale computer essa decisão é do scheduler, não de quem escreveu o programa · pelo resumo
  7. Mesa: Geo-Replicated, Near Real-Time, Scalable Data Warehousingwarehouse não precisa ser a foto de ontem: dá para ingerir milhões de atualizações de linha por segundo, servir bilhões de consultas por dia e ainda garantir a mesma resposta quando um datacenter inteiro cai · pelo resumo
  8. Magnet: Push-based Shuffle Service for Large-scale Data Processingo mapper empurra seus blocos para um shuffle service remoto que os mescla por partição: leitura aleatória de dezenas de kb vira leitura sequencial de mb, e o push nem precisa terminar para o job ganhar
  9. Kora: A Cloud-Native Event Streaming Platform for Kafkapara rodar kafka como serviço, o protocolo não precisou mudar — precisou mudar tudo em volta: storage em dois níveis, células por tenant, cota coordenada e uma métrica de carga tirada de teoria de filas

de nicho

  1. Yedalog: Exploring Knowledge at Scaleo gargalo da análise em escala deixou de ser a máquina e virou o programador — e parte da culpa é a costura entre a linguagem do pipeline e a linguagem em que se escreve a lógica · pelo resumo