Objetivos dos Sistemas Distribuídos
Objetivos dos Sistemas Distribuídos
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Sistema Distribuído: Definição e Características
❑ Características fundamentais:
– Vários computadores autónomos
– Sistema único e coerente
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Middleware
❑ Um dos problemas fundamentais em sistemas distribuídos é o
tratamento da heterogeneidade (SOs, linguagens de programação,
hardware e bibliotecas diferentes)
❑ Solução mais comum: Middleware
– Esconde as diferenças entre diferentes plataformas
– Introduz transparência, tentando esconder a distribuição
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Objectivos na Construção de Sistemas Distribuídos
❑ Objectivos principais:
1. Acessibilidade de Recursos
2. Transparência
3. Abertura
4. Segurança de funcionamento (Dependability)
5. Segurança
6. Escalabilidade
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
1. Acessibilidade de Recursos
❑ O principal objectivo de um sistema distribuído é dar acesso e
partilhar recursos remotos de uma forma controlada e eficiente
– Exemplos de recursos: impressoras, computadores,
armazenamento, ficheiros (páginas web), dados (ex., emails, base
de dados), etc.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
2. Transparência
❑ Objectivo: simular o comportamento de um sistema não-distribuído, i.e.,
esconder distribuição
❑ Existem vários tipos de transparência
Tipo Esconder…
Acesso a forma de acesso e a representação dos dados
Localização a localização dos recursos
Migra\ção a mudança de localização dos recursos
Re-localização a mudança de localização dos recursos, enquanto são usados
Replicação a replicação dos recursos
Concorrência o acesso simultâneo por outros clientes aos recursos
Falhas as falhas de um recurso bem como a recuperação do mesmo
❑ Exemplos :
– o nome do recurso (e.g., impressora, ou URL) não inclui a sua localização
– replicar os ficheiros mais acedidos em vários servidores de ficheiros
– acessos concorrentes a um ficheiro são automaticamente protegidos por
locks
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
2. Transparência
❑ Grau de transparência, i.e., De quanta transparência preciso?
❑ A transparência é desejável numa aplicação distribuída, no entanto, por
vezes não faz sentido esconder a distribuição, pois:
– Não é possível:
» Exemplo: esconder latência de comunicação na Internet
– Não é boa ideia:
» Exemplo: imprimir um documento em alguma impressora escolhida
pelo sistema, quando seria melhor o utilizador escolher a
impressora mais próxima
– Tem impacto no desempenho:
» Transparência exige autonomia que requer mais código e algoritmos
mais complexos. Tipicamente temos:
Transparência
Desempenho
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
3. Abertura
❑ Um sistema é aberto se existem regras e interfaces bem
definidas que especificam a sintaxe e a semântica dos
serviços fornecidos
– Possibilita a inclusão de componentes de diferentes
fabricantes
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
3. Abertura
❑ Conceitos relacionados:
– Interoperabilidade – funcionamento em conjunto de
componentes de fabricantes diferentes
– Portabilidade – funcionamento de componente sem
modificação noutro sistema
– Extensibilidade – funcionalidades de componentes podem ser
trocadas ou melhoradas sem modificar o sistema como um todo
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
4. Segurança de funcionamento
❑ Até que ponto podemos confiar que um sistema funcione como é
esperado.
❑ Num sistema distribuído é mais difícil que num sistema não
distribuído devido a falhas parciais de componentes.
❑ Requisitos:
– Disponibilidade (availability) – probabilidade de o sistema
funcionar corretamente em qualquer momento (prontidão para
fornecer o serviço)
– Confiabilidade (reliability) – intervalo de tempo em que o
sistema pode funcionar continuamente sem falhas (continuidade
do serviço)
– Segurança (safety) – quando ocorrem falhas o sistema não causa
catástrofes (com o sentido de proteção, não de cibersegurança)
– Capacidade de manutenção (maintainability) – capacidade de
um sistema em falha ser reparado e continuar em
funcionamento
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
5. Segurança
❑ Se o sistema não é seguro (cibersegurança), então não oferece
segurança de funcionamento (não cumpre os requisitos desta)
❑ Requisitos de segurança:
– Confidencialidade – informação e recursos do sistema só são
disponibilizados a componentes autorizados
– Integridade – informação e recursos do sistema só podem ser
alterados de forma autorizada
❑ Mecanismos:
– Autenticação – verificação de uma identidade apresentada ao
sistema
– Autorização – verificação de permissões de acesso a recursos
por entidades verificadas
❑ Criptografia é a técnica principal:
– Cifrar e decifrar a comunicação e a informação armazenada
– Assinaturas digitais
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
6. Escalabilidade
❑ Um sistema distribuído que funciona bem com 100 máquinas vai
funcionar bem com 1000?
❑ A escalabilidade tem pelo menos três dimensões:
1- tamanho do sistema (número de máquinas e de utilizadores)
Quanto maior, mais carga deve suportar… temos de evitar os
“gargalos”
2- distribuição geográfica
Há que lidar com elementos ou utilizadores longínquos, que
podem gerar latências grandes (e imprevisíveis)
3- em termos de administração
Administradores diferentes tem ideias diferentes sobre o que
deve ou não ser usado na sua rede
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
4. Escalabilidade
❑ Deve-se evitar todo o tipo de “centralização”
Aspeto Exemplo
Componentes Um servidor único (e.g. de mail ou web) para todos os
utilizadores
Dados Uma tabela única com os nomes e telefones dos utilizadores
Algoritmos Algoritmos de encaminhamento que necessitam de
informação sobre o estado completo do sistema
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
4. Escalabilidade
Centralizado
Descentralizado
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Técnicas para Escalabilidade: Esconder Latência
❑ Comunicação Assíncrona
– Evitar bloquear a execução do programa enquanto espera por
respostas de servidores
– Enviar requisição e fazer outras coisas até a resposta chegar
– Exemplo: RPC (Remote Procedure Call)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Técnicas para Escalabilidade: Esconder Latência
❑ E quando não há “outras tarefas” a fazer? Evita-se a comunicação o
máximo possível!
– Tentar resolver as coisas na máquina local
– Dois benefícios adicionais em termos de desempenho:
1. O envio de uma mensagem grande é muito mais rápido que várias
pequeninas
2. Traz o trabalho para o lado do cliente e alivia a carga no servidor
Exemplo:
Validação da entrada
de dados em
formulários na Web
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Técnicas para Escalabilidade: Particionar e Distribuir
❑ Dividir um componente em diversos componentes menores e espalhá-los
pelo sistema
❑ Cada máquina tenta resolver as suas requisições nos componentes mais
próximos de si
Observação: A isto, juntamente com a ideia de fazer o máximo de trabalho
localmente, chamamos explorar a localidade
5 - E assim continuaria…
4
Exemplo:
Resolução 3
de nomes
no DNS 2
1
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Técnicas para Escalabilidade: Replicar
❑ Replicar para quê?
– Balancear carga
» Distribuir o trabalho de forma inteligente por várias máquinas
– Aumentar a disponibilidade
» Ter réplicas dos dados próximo dos sítios onde eles serão usados
❑ Problemas com a replicação:
– (in)consistência, i.e., versões diferentes dos dados em diferentes
réplicas…
– Os protocolos para replicação com consistência forte (i.e., sistema com
transparência de replicação) são, em geral, lentos, e não escaláveis!
❑ Exemplo: Cache
– Dados são copiados para sítios mais próximo de onde serão usados
– Os problemas de consistência aparecem na mesma
– Obs: Outra razão para a escalabilidade do DNS é o uso de cache: um
endereço só é resolvido num nível superior se ele não estiver na cache
do nível atual
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Ideias erradas sobre Sistemas Distribuídos
❑ Existe um conjunto de ideias erradas (ou simplistas) que os
engenheiros de sistemas distribuídos “caloiros” geralmente têm.
Segundo Peter Deutsch, as mais comuns são:
1. A rede é fiável
2. A rede é segura
3. A rede é homogénea
4. A topologia da rede não muda
5. Latência é próxima a zero
6. Largura de banda é infinita
7. O custo do transporte é próximo a zero
8. Existe apenas um administrador
❑ Estas ideias são geralmente válidas em redes locais, mas não o são
na Internet
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Tipos de Sistemas Distribuídos
(Capítulo 1, páginas 32-52)
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Tipos de Sistemas Distribuídos
❑ Existem três tipos de sistemas distribuídos:
– Sistemas de Computação Distribuídos
Computação de alto desempenho com vários computadores
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Tipos de Sistemas Distribuídos
❑ Existem três tipos de sistemas distribuídos:
– Sistemas de Computação Distribuídos
Computação de alto desempenho com vários computadores
+ ou - estáticos
– Sistemas de Informação Distribuídos
Integração de diversos sistemas de informação
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Computação Distribuída
Clusters
Computação em Grid
Computação na nuvem
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Clusters
❑ Definição: conjunto de várias máquinas “comuns” ligadas por uma rede de alto
desempenho de tal forma a trabalharem como um super-computador
❑ Usado para processamento paralelo (i.e., criam-se processos em diferentes nós)
❑ Melhor relação custo benefício que um super-computador:
– Super-computador com 500 processadores é muito mais caro que 500
computadores “normais” do tipo PC
❑ Exemplo: Cluster Beowulf (baseado em Linux)
– Middleware: (1.) gestão e distribuição de tarefas
(2.) comunicação eficiente entre processos
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Clusters
❑ Definição: conjunto de várias máquinas “comuns” ligadas por uma rede de alto
desempenho de tal forma a trabalharem como um super-computador
❑ Usado para processamento paralelo (i.e., criam-se processos em diferentes nós)
❑ Melhor relação custo benefício que um super-computador:
– Super-computador com 500 processadores é muito mais caro que 500
computadores “normais” do tipo PC
❑ Exemplo: Cluster Beowulf (baseado em Linux)
– Middleware: (1.) gestão e distribuição de tarefas
(2.) comunicação eficiente entre processos
Distribui
tarefas e
coordena
computação
distribuída.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Clusters
❑ Definição: conjunto de várias máquinas “comuns” ligadas por uma rede de alto
desempenho de tal forma a trabalharem como um super-computador
❑ Usado para processamento paralelo (i.e., criam-se processos em diferentes nós)
❑ Melhor relação custo benefício que um super-computador:
– Super-computador com 500 processadores é muito mais caro que 500
computadores “normais” do tipo PC
❑ Exemplo: Cluster Beowulf (baseado em Linux)
– Middleware: (1.) gestão e distribuição de tarefas
(2.) comunicação eficiente entre processos
Executa
Distribui tarefas
tarefas e definidas
coordena pelo master.
computação
distribuída.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Clusters
❑ Definição: conjunto de várias máquinas “comuns” ligadas por uma rede de alto
desempenho de tal forma a trabalharem como um super-computador
❑ Usado para processamento paralelo (i.e., criam-se processos em diferentes nós)
❑ Melhor relação custo benefício que um super-computador:
– Super-computador com 500 processadores é muito mais caro que 500
computadores “normais” do tipo PC
❑ Exemplo: Cluster Beowulf (baseado em Linux)
– Middleware: (1.) gestão e distribuição de tarefas
(2.) comunicação eficiente entre processos
Executa
Distribui tarefas
tarefas e definidas
coordena pelo master.
computação
distribuída.
Usada para
comunicação
inter-nós.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Em 11/2024
Lista de super-
computadores
mais rápidos:
[Link]
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação em Grid
❑ Uma das principais facilidades em se lidar com clusters é o ambiente
homogéneo, i.e., todas as máquinas são iguais e correm o mesmo software
❑ Grids são sistemas distribuídos usados para computação em larga escala, mas
com um alto grau de heterogeneidade:
– Cada nó pode ter hardware, sistema operativo, rede, administração, e
políticas de segurança diferentes
– Esta heterogeneidade tem de ser tratada pelo middleware
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação em Grid
Retirado de [Link]
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação em Grid
❑ O middleware dá acesso aos recursos a utilizadores da sua organização, ainda
que os recursos estejam em domínios administrativos diferentes
❑ Arquitectura em 4 níveis para o middleware de grids:
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação em Grid
❑ O middleware dá acesso aos recursos a utilizadores da sua organização, ainda
que os recursos estejam em domínios administrativos diferentes
❑ Arquitectura em 4 níveis para o middleware de grids:
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação na Nuvem
❑ A computação na cloud (ou nuvem) concretiza o conceito de utility computing, onde
um utilizador envia sua computação para um serviço e é cobrado pelo uso de
recursos.
❑ Do ponto de vista de concretização, a cloud se caracteriza por oferecer acesso a um
conjunto de recursos virtualizados, que podem ser disponibilizados de várias formas.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Computação na Nuvem
❑ A computação na cloud (ou nuvem) concretiza o conceito de utility computing, onde
um utilizador envia sua computação para um serviço e é cobrado pelo uso de
recursos.
❑ Do ponto de vista de concretização, a cloud se caracteriza por oferecer acesso a um
conjunto de recursos virtualizados, que podem ser disponibilizados de várias formas.
Aplicações na
nuvem
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Informação Distribuídos
Sistemas de Processamento de Transações
Integração de Aplicações Corporativas
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Transações são usualmente empregues em operações com bases de dados
mas não só… no futuro, talvez o sejam em todo lado (memória transaccional)
❑ O objectivo fundamental das transações é
– proteger um recurso (e.g., base de dados) contra acessos simultâneos por
parte de diversos processos
– permitir que um processo faça diversas operações sobre um ou mais
recursos, como se estas fossem uma única operação atómica (indivisível)
❑ Modo de funcionamento
1. indicar início de transação
2. aceder uma ou mais vezes ao recurso (e.g., leitura e/ou escrita)
3. indicar que as alterações deverão ser permanentes (commit) ou que se
pretende abortar (i.e., esquecer) as alterações (abort)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações: Interfaces
Primitive Description
BEGIN_TRANSACTION Make the start of a transaction
END_TRANSACTION Terminate the transaction and try to commit
ABORT_TRANSACTION Kill the transaction and restore the old values
READ Read data from a file, a table, or otherwise
WRITE Write data to a file, a table, or otherwise
T1: T2:
BEGIN_TRANSACTION BEGIN_TRANSACTION
reservar Lisboa → Londres reservar Paris → NY
reservar Londres → NY reservar NY → SF
reservar NY → SF END_TRANSACTION
END_TRANSACTION
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações: Interfaces
Primitive Description
BEGIN_TRANSACTION Make the start of a transaction
END_TRANSACTION Terminate the transaction and try to commit
ABORT_TRANSACTION Kill the transaction and restore the old values
READ Read data from a file, a table, or otherwise
WRITE Write data to a file, a table, or otherwise
T1: T2:
BEGIN_TRANSACTION BEGIN_TRANSACTION
reservar Lisboa → Londres reservar Paris → NY
reservar Londres → NY reservar NY → SF
reservar NY → SF END_TRANSACTION
END_TRANSACTION
As duas não podem ser confirmadas se há
apenas um lugar vago no voo NY → SF
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
Propriedades ACID das Transações
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Transações Aninhadas:
– Transações que afetam vários recursos são divididas em subtransações
– O uso de transações aninhadas pode levar a alguns problemas:
» Suponha uma transação t com duas subtransações t1 e t2, executadas em paralelo
» Se t1 se confirmar e t2 abortar, o que devemos fazer?
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Transações Aninhadas:
– Transações que afetam vários recursos são divididas em subtransações
– O uso de transações aninhadas pode levar a alguns problemas:
» Suponha uma transação t com duas subtransações t1 e t2, executadas em paralelo
» Se t1 se confirmar e t2 abortar, o que devemos fazer?
➢ Opção 1: os resultados de t1 valem (respeita-se a persistência de t1)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Transações Aninhadas:
– Transações que afetam vários recursos são divididas em subtransações
– O uso de transações aninhadas pode levar a alguns problemas:
» Suponha uma transação t com duas subtransações t1 e t2, executadas em paralelo
» Se t1 se confirmar e t2 abortar, o que devemos fazer?
➢ Opção 1: os resultados de t1 valem (respeita-se a persistência de t1)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Transações Aninhadas:
– Transações que afetam vários recursos são divididas em subtransações
– O uso de transações aninhadas pode levar a alguns problemas:
» Suponha uma transação t com duas subtransações t1 e t2, executadas em paralelo
» Se t1 se confirmar e t2 abortar, o que devemos fazer?
➢ Opção 1: os resultados de t1 valem (respeita-se a persistência de t1)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas de Processamento de Transações
❑ Monitor de Processamento de Transações:
– Middleware usado na integração de vários sistemas de informação através
do uso de transações
– Oferece uma interface de transações para sistemas distribuídos
» Exemplo de uso: concretizar transferência de dinheiro entre bases de
dados de bancos diferentes
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Integração de Aplicações Corporativas
❑ Nem todas as aplicações são bases de dados, e em alguns casos é preciso integrar
aplicações sem aceder às bases de dados que estão por trás
❑ As aplicações devem ser capazes de comunicar umas com as outras, e não permitir que
os “clientes” acedam às suas bases de dados (como nos monitores de transacções)
❑ Existem vários modelos de comunicação, mas em geral, a integração de aplicações é
feita através de middleware orientado a mensagens
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas “Pervasivos” Distribuídos
Sistemas Ubíquos
Sistemas Móveis
Redes de Sensores
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas “Pervasivos” Distribuídos
❑ Caracterizados pelo facto dos dispositivos/sistemas fazerem parte do
ambiente”
– A separação entre utilizadores, sistema e ambiente não é muito clara
❑ Fazem uso de todo tipo de dispositivos
❑ Hoje em dia estão muito ligados a ideia de Internet das Coisas
❑ Um exemplo:
Sistema
de Saúde
Eletrónico
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas Ubíquos
❑ Do dicionário Priberam ([Link]
❑ Requisitos fundamentais:
– Distribuição: dispositivos distribuídos e conectados de maneira transparente
– Interação: a interação entre os utilizadores e o sistema deve ser “natural”
– “Consciência” do contexto: o sistema tem em conta o contexto do utilizador
(e.g., localização, conectividade, bateria) para otimizar interações
– Autonomia: dispositivos operam de maneira autônoma, com o mínimo de
configuração e esforço de operação
– Inteligência: sistema como um todo pode se adaptar a uma variedade de
situações
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo de Sistemas Ubíquos: Home Systems
❑ Construídos sobre redes domésticas, que usualmente integram:
– Um ou mais PCs
– Uma série de aparelhos electrónicos como TVs, equipamento de áudio,
consolas, telemóveis, TV box, etc.
– Dispositivos como câmaras de vigilância, controlos de luz, etc.
❑ A ideia é que todo esse equipamento seja interligado para que se possa:
– Controlar o volume da TV e do equipamento de áudio pelo telemóvel
– Ver a saída da câmara de vigilância no PC
– Exportar fotos no PC para as visualizar na TV
❑ Requisito número 1: o sistema deve ser auto-configurável e auto-gerido
– Não vamos esperar que nossa avó consiga ligar todas essas coisas na sua
casa nova ☺
– É preciso uma interface plug and play (e.g., UPnP)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sistemas Móveis
❑ Satisfazem apenas alguns dos requisitos dos sistemas ubíquos
(quais?)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Redes de Sensores
❑ Redes de sensores são fundamentais para muitas aplicações “pervasivas”
❑ Uma rede de sensores consiste em muitos pequenos sensores (temperatura,
movimento, radiação, etc.) com algum poder de processamento e comunicação
– Os elementos da rede são em geral muito limitados em termos de
processamento, comunicação e energia
– Mas custam pouco: de cêntimos a centenas de euros
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Redes de Sensores
❑ As redes de sensores normalmente operam como bases de dados
❑ Exemplo: monitorização de estradas
– Considere uma rede de sensores nas estradas de Portugal
– Lançam-se consultas do tipo: “Como está o tráfego na A2?”
– Para responder a esta consulta, os sensores espalhados pela A2 devem
colaborar com os seus dados. Há duas soluções possíveis:
(a.) Cada sensor envia os seus dados para uma base de dados centralizada
e uma “aplicação comum” consolida o resultado
(b.) Encaminha-se a consulta aos sensores relevantes que calculam a
resposta, sendo a aplicação/operador responsável por consolidar um
resultado final
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Redes de Sensores
Solução (a.)
Problema: envia muita informação –
desperdício de rede e energia.
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Redes de Sensores
Solução (a.)
Problema: envia muita informação –
desperdício de rede e energia.
Solução (b.)
Problema: desperdiça a capacidade de
agregação da rede (podia ser retornada
muito menos informação)
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Redes de Sensores
Solução (a.)
Problema: envia muita informação –
desperdício de rede e energia.
Solução (b.)
Problema: desperdiça a capacidade de
agregação da rede (podia ser retornada
muito menos informação)
❑ Uma terceira solução melhor ainda, seria fazer o processamento dos dados na rede:
– Define-se uma arvore de distribuição na rede de sensores
– Difunde-se a consulta nessa arvore
– Os resultados são enviados de volta (devidamente agregados) até a raiz da arvore
© DI-FCUL Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Comunicação
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Revisão de Redes de Computadores
❑ Protocolos em Camadas
– Aplica o estilo arquitetural das
camadas na concretização de pilhas
de protocolos
– Boa forma de lidar com a
complexidade inerente ao software
necessário para comunicação
❑ Protocolos de Baixo Nível
– Físico
– Link
– Rede (IP)
❑ Protocolos de Transporte
– Transporte (TCP, UDP)
❑ Protocolos de Alto Nível
– Sessão
– Apresentação
– Aplicação
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Protocolos de Middleware
❑ Nem sempre as aplicações “falam” através dos protocolos de transporte.
– Exemplo: Web Services utilizam SOAP, que é encapsulado em HTTP (em
geral) que por sua vez é encapsulado em TCP
❑ Os protocolos de alto nível (não específicos de uma aplicação), são chamados
protocolos de middleware:
– Podem suportar serviços como transacções, autenticação, autorização e
sincronização.
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Protocolos de Middleware: tipos de Comunicação
Middleware
❑ Persistência:
– Transitória: a mensagem é armazenada no sistema de comunicação
apenas enquanto o emissor e o recetor estão ativos
– Persistente: a mensagem é armazenada no sistema de comunicação até
ser entregue ao recetor
❑ Sincronização na comunicação:
– Síncrona: cliente bloqueia à espera da resposta do servidor
– Assíncrona: cliente não bloqueia à espera do servidor, e recebe uma
notificação quando a resposta está disponível
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Chamadas a
Procedimentos Remotos
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Comunicação por mensagens
❑ A comunicação por mensagens (UDP, TCP,…) tem duas limitações:
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Chamada de Procedimentos Remotos
❑ Nas linguagens procedimentais existe
PROCESSO 1 PROCESSO 2
um fluxo de actividade que chama
sincronamente os procedimentos do
programa
Invocação do
Bloqueia-se Procedimento
Remoto
❑ A chamada a procedimentos permite a
transferência de controlo e dados Execução do
dentro do programa Procedimento
Retoma a
❑ A Chamada de Procedimentos execução Devolução dos
Parâmetros de
Remotos (Remote Procedure Call - Resposta
RPC) pode ser vista como uma
extensão deste modelo para o caso
dos sistemas distribuídos
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Chamada a Procedimentos entre Dois Processos
❑ Resumo dos passos
– chamada à rotina de adaptação do cliente (stub)
– a rotina de adaptação constrói a mensagem e passa-a ao SO
– o SO cliente transmite a mensagem pela rede
– o SO remoto passa a mensagem à rotina de adaptação do servidor (skeleton)
– a rotina desempacota os parâmetros e chama o procedimento local
– o servidor executa o pedido e devolve o resultado pela rotina de adaptação
– a rotina empacota o resultado numa mensagem e passa-a ao SO
– o SO remoto transmite a mensagem pela rede
– o SO cliente passa mensagem à rotina de adaptação do cliente
– a rotina retorna o resultado
.. def max():
.
max()
.. ..
. .
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Rotinas de adaptação (stubs)
❑ Rotina de adaptação do cliente
[Link](P)
R = [Link](….) //bloquear à espera da resposta
return r
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Rotinas de adaptação (stubs)
❑ Rotina de adaptação do servidor
if procedimento == ‘max’:
r = max(a, b) # executar o procedimento local
# serializar r na mensagem R
[Link](R)
elif … # código semelhante para tratar outros procedimentos servidos
…
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Empacotamento dos Parâmetros
❑ Para que o RPC funcione é fundamental que o cliente e o servidor concordem
sobre o formato das mensagens de requisição e resposta. Estas são
transmitidas byte a byte, em série.
❑ Devem concordar sobre tudo: estrutura da mensagem, formato dos dados
(boolean, float, int, short, string, etc)
❑ É possível gerar estas rotinas de
empacotamento/desempacotamento
4 bytes {
automaticamente através de certas ferramentas.
(e.g., rpcgen do Sun RPC, CORBA idl).
procedimento mensagem
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Empacotamento dos Parâmetros–formato
❑ Exemplos de formatos dos dados:
– char é ASCII (Linguagem C) ou unicode (linguagem Java)?
– float segue o padrão IEEE 754?
❑ Big Endian (bytes da direita para a esquerda) vs Little Endian (ao contrário)
– Exemplo: enviar mensagem <5,”Jill”>
– Mensagens são enviadas bit a bit (no nível físico)
– Se não invertermos há problema (fig. b) e inverter tudo também não funciona (fig. c)
» A solução é tratar apenas os tipos numéricos
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Geração das Rotinas de Adaptação (stubs)
Ficheiro de
IDL - Interface Description Language cabeçalhos
Empacota
app.h
pedidos e
Definição desempacota
da Interface Ficheiro com respostas.
Compilador
em IDL Stub Cliente
IDL
app_client.c
[Link]
Ficheiro com
Stub Servidor
app_server.c
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Localização do Servidor
❑ Codificar diretamente no cliente o endereço do servidor
– Problema: pouco flexível, falta de transparência
– Solução: servidor de nomes (SN)
❑ Associação dinâmica ao servidor
– Operações a suportar pelo Servidor de Nomes
Operação Argumentos Resultado
register() nome, versão, handle, identificador Servidores
deregister() nome, versão, identificador (procedimentos remotos)
lookup() nome, versão handle, identificador Clientes
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
RPC Assíncrono
❑ Os RPCs assíncronos aumentam o desempenho dos RPCs tradicionais (síncronos) pois
evitam o bloqueio do cliente enquanto o pedido se encontra a ser processado
O servidor remoto
não precisa de devolver
qualquer resultado
RPC tradicional
RPC assíncrono
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
RPC Assíncrono
❑ Os RPCs assíncronos aumentam o desempenho dos RPCs tradicionais (síncronos) pois
evitam o bloqueio do cliente enquanto o pedido se encontra a ser processado
O servidor remoto
Existem variantes assíncronas não precisa de devolver
em que o cliente não espera qualquer resultado
pela confirmação da aceitação
da chamada: one-way RPCs
RPC tradicional
RPC assíncrono
© DI-FCUL – Pedro Ferreira, Vinicius Cogo. Reprodução proibida sem autorização prévia.
Sincronização e relógios
(Páginas 249-264)
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Sincronização em Sistemas Distribuídos
❑ O que são problemas de sincronização/coordenação?
São problemas em que os processos de um sistema distribuído precisam cooperar
uns com os outros
❑ Exemplos de problemas deste tipo (que estudamos na cadeira):
1. Sincronização de Relógios
2. Relógios Lógicos
3. Exclusão Mútua
4. Eleição de Líder
5. Acordo e Ordenação de Mensagens
❑ Como resolver problemas de sincronização?
– Através de algoritmos de sincronização, que funcionam sobre pressupostos do
ambiente distribuído
– Exemplos de pressupostos:
» Não há falhas de máquinas
» Os canais de comunicação são fiáveis
» A latência de comunicação é no máximo T
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Sincronização de Relógios
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Motivação
❑ Exemplo : compilador e editor num processador
Tempo baseado
1234 1256 1278 1290 no relógio local
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Funcionamento de um Relógio Físico
Contador é
decrementado contador valor inicial
em cada oscilação
do cristal
Quando o contador chega a 0,
é gerada uma interrupção e o valor
inicial é colocado no contador
(clock tick)
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Relógios Físicos
❑ Pretendemos que os processos cheguem a um acordo em relação ao tempo, e
eventualmente também que esse tempo seja próximo do tempo real
❑ Maior problema:
– os relógios físicos têm taxas de deriva (drift rates) diferentes
Tempo do dC/dt > 1 dC/dt = 1
Relógio, C dC/dt < 1
especificação:
rápido: C(t) > t
dC
1− 1+
dt lento: C(t) < t
drift máx.
ideal: C(t) = t
Tempo Real, t
– O erro máximo entre dois relógios após t é t × 2
– Para garantir um erro entre relógios (clock skew) menor que , é preciso
resincronizar os relógios periodicamente com período t < / 2
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo de Cristian (Centralizado)
OBJECTIVO: sincronizar os relógios das máquinas
clientes pelo relógio da máquina servidora
Qual é o valor
T0 do teu relógio? Servidor de Tempo; pode ter acesso a
relógio de melhor qualidade (GPS)
T1 Processamento
C
local
C_local = C?
Máquinas
Clientes
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo de Cristian (Centralizado)
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
NTP: Network Time Protocol
❑ Algoritmo usado para sincronização de relógios na Internet
– Consegue uma precisão de aproximadamente 1-50 milissegundos
❑ Porque funciona?
– Se fizermos muitas medições e escolhermos a de menor latência, usamos o
valor de θ sujeito a menor erro
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
NTP: Network Time Protocol
❑ Problema: Se aplicarmos a metodologia “indiscriminadamente” na rede, pode
ocorrer que uma máquina com um relógio mais preciso se acerte pelo relógio
menos preciso de outra máquina (com a qual tenha latência de comunicação
baixa)!
❑ Para evitar este problema, as máquinas são dividas em estratos:
– Menor estrato → Melhor relógio (i.e., menor drift)
» E.g., estrato 1 - máquina com relógio atómico ( = 0)
– Uma máquina A (com estrato nA) só sincroniza com B (estrato nB) se nA > nB
❑ Após a sincronização, nA = nB + 1 (A passa ao estrato acima de B)
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo de Berkeley (Descentralizado)
Qual o valor do
teu relógio? Daemon
OBJECTIVO : sincronizar os relógios de Tempo
OBJECTIVO
entre : sincronizar
as várias máquinasosderelógios
forma
entre as várias
distribuída (nestemáquinas de temos
caso, não forma
distribuída
uma máquina (neste
pelacaso,
qual não temos uma
as outras
máquina pela
máquinas qual as outras
sincronizam máquinas
os seus Qual o valor
sincronizam
relógios). os seus érelógios)
A solução acertar para o do teu relógio?
valor médio dos relógios.
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo de Berkeley (Descentralizado)
Qual o valor do
OBJECTIVO teu relógio? Daemon
OBJECTIVO :: sincronizar
sincronizar os
os relógios
relógios de Tempo
entre
entre as
as várias
várias máquinas
máquinas dede forma
forma
distribuída
distribuída (neste
(neste caso,
caso, não
não temos
temos
uma
uma máquina pela qual
máquina pela qual as
as outras
outras
máquinas
máquinas sincronizam
sincronizam os
os seus
seus Qual o valor
relógios)
relógios). A solução é acertar para o do teu relógio?
valor médio dos relógios.
No meu No meu
são 3:25 são 2:50
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo de Berkeley (Descentralizado)
Qual o valor do
teu relógio? Daemon
OBJECTIVO : sincronizar os relógios de Tempo
entre as várias máquinas de forma
distribuída (neste caso, não temos
uma máquina pela qual as outras
máquinas sincronizam os seus Qual o valor
relógios). A solução é acertar para o do teu relógio?
valor médio dos relógios.
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Relógios Lógicos
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Relógios Lógicos
❑ Para muitas aplicações é suficiente que os processos cheguem a um acordo em
relação à ordem com que ocorrem certos eventos, não sendo necessário que os
processos cheguem a um acordo sobre o tempo real
transitividade : se a → b e b → c então a → c
eventos concorrentes : nem x → y, nem y → x
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Relógios Lógicos
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Concretização do Relógio Lógico de Lamport
❑ Os relógios lógicos são usualmente concretizados na camada de middleware
– Os eventos, e.g., mensagens enviadas, incrementam o contador e anexam o
valor do mesmo aos eventos
– Os instantes das mensagens recebidas são usados para ajustar o contador
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Forma de Construir o Relógio (Lamport)
1. Manter um contador C(e) em cada processo, que é inicializado a 0
2. Sempre que ocorre um evento interno (i.e., no mesmo processo), incrementa-
se C(e) = C(último_evento_local) + 1, e atribui-se o seu valor ao evento
– O envio de uma mensagem também conta como evento interno
3. Quando se transmite uma mensagem, anexa-se aos dados o valor atual do
contador C(e_mesg)
4. Sempre que se recebe uma mensagem, atualiza-se se o contador:
© DI-FCUL Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira. Reprodução proibida sem autorização prévia
Uso de Relógios Lógicos: Difusão com Ordem Total
Réplica 1 Réplica 2
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Motivação
❑ Já conhecemos o problema da exclusão mútua de sistemas operativos
– O objectivo é garantir que dois ou mais processos não acedem ao mesmo
recurso simultaneamente, o que poderia causar inconsistências
(e.g., um ficheiro, uma estrutura de dados, um dispositivo de entrada de dados, etc…)
❑ Em sistemas distribuídos o problema é o mesmo…
– Que aplicação de uma grid vai usar a máquina M entre os dias X e Y?
» Alocação de recursos
– Num cluster, que thread (computador) de uma aplicação vai executar uma
determinada tarefa que está preparada para execução?
» Produtor/consumidor
– Como garantir que vários processos em máquinas diferentes acedem a um
recurso partilhado que não suporta concorrência?
(e.g., ficheiro, impressora)
» Generalização: exclusão mútua
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua em Sistemas Distribuídos
Definição do Problema
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua em Sistemas Distribuídos
Propriedades a serem satisfeitas:
Safety e Liveness
são conceitos fundamentais
em qualquer algoritmo
distribuído
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Como resolver o problema?
❑ Algumas observações sobre as propriedades de vivacidade:
– Processo que tenta aceder sem concorrência a um recurso é bem sucedido com V1
– V3 é mais forte que V2 (que é mais forte que V1): V2 não garante que todos os
processos que tentam, conseguem aceder ao recurso (dando origem a algoritmos
não equitativos).
– Quando V1 ou V2 são suficientes?
❑ Hipóteses:
– Os processos não falham e as mensagens não são perdidas
– Um dos processos é selecionado como coordenador (usando, p. ex., um
dos algoritmos de eleição de líder que veremos na próxima aula)
❑ Algoritmo:
1. Sempre que um processo quer entrar numa região crítica, envia uma
mensagem ao coordenador com o identificador do recurso
2. O coordenador devolve uma mensagem indicando que o processo pode
continuar (OK), se nenhum outro processo estiver nesse momento na
região crítica
3. Caso contrário, o coordenador não devolve qualquer mensagem (ou
retorna uma mensagem NOK indicando que o processo não tem
permissão), e coloca o processo numa fila de espera
4. Quando um processo deixa a região crítica, envia uma mensagem ao
coordenador libertando o recurso
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Centralizado
❑ Ilustração:
OK 2
2 C 2 Posso entrar ? 2 OK C
❑ Propriedades satisfeitas:
– S1 e V3
❑ Discussão:
– Vantagens: o algoritmo é equitativo (satisfaz V3); fácil de realizar; apenas
três mensagens para aceder ao recurso
– Problemas: o coordenador é um ponto único de falha; o coordenador pode
trazer problemas de desempenho (bootleneck do sistema)
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Descentralizado
Algoritmo Descentralizado
❑ Hipóteses:
– Os processos não falham e as mensagens não são perdidas
– Um conjunto de n processos são coordenadores
❑ Algoritmo:
1. Sempre que um processo quer entrar numa região crítica, envia uma
mensagem aos coordenadores com o identificador do recurso a ser acedido
2. Cada coordenador devolve uma mensagem OK indicando que o processo pode
continuar, ou uma mensagem NOK indicando que o acesso ao recurso está
cedido a outro processo
3. Um processo apenas entra na sua secção critica se recebe pelo menos m > n/2
mensagens OK de diferentes coordenadores
4. Se a maioria das mensagens for NOK, o processo avisa os coordenadores que
votaram OK que ele não tem acesso e espera uma quantidade de tempo
aleatória para voltar a tentar executar o algoritmo
(Similar ao backoff usado para transmissão de dados nas redes Ethernet)
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Descentralizado
❑ Ilustração:
1 2
C1 C2 C3 C4 C5
Envia OK ao 1 Envia OK ao 2
❑ Propriedades satisfeitas:
– S1 e V1
❑ Discussão:
– Vantagens: é fácil de realizar; funciona em sistemas de larga escala como
por exemplo P2P (os coordenadores podem ser dispostos numa DHT); pode
tolerar faltas se aumentarmos o numero de votos OK esperados
– Problemas: Satisfaz apenas V1
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Descentralizado II
❑ Como satisfazer vivacidade? (i.e., como fazer o algoritmo centralizado
descentralizado?)
❑ Se usarmos difusão com ordem total para enviar os pedidos de acesso e
libertação de recursos dos processos aos coordenadores…
(Difusão com ordem total: todos recebem a mesma sequência de mensagens)
❑ No resto, o algoritmo é igual… e satisfaz V3
1, 2
C1
1
enter(r) 1, 2
C2
2
Difusão com Ordem 1, 2
enter(r) C3
Total (fila virtual)
1, 2
C4
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Distribuído
Algoritmo Distribuído
❑ Hipóteses:
– Assume que é possível ordenar todos os eventos no sistema (com relógios lógicos)
– Os processos não falham e as mensagens não são perdidas
❑ Algoritmo:
– Quando quer entrar numa região crítica o processo envia uma mensagem a todos
os outros processos, com a seguinte informação
< identificador do recurso; identificador do processo; instante temporal >
– Quando um processo recebe a mensagem faz o seguinte
1. devolve OK se não está na região crítica e não pretende entrar
2. guarda a mensagem numa fila e não responde se estiver na região crítica
3. se quer entrar na região crítica
➢ compara os instantes temporais do seu pedido e da mensagem
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo Distribuído
❑ Ilustração: A quer
entrar A vai
2 entrar
A C quer A A
2 2 entrar OK
OK OK
6 C vai
B C 6 B C B C entrar
6 OK
❑ Propriedades:
– S1 e V3
❑ Discussão:
– Vantagens: o algoritmo é equitativo (satisfaz V3);
– Problemas: num sistema com n processos precisamos de 2(n-1)
mensagens para entrar na região crítica; temos n pontos de falha; cada
processo precisa de manter uma lista dos processos que potencialmente
podem entrar na região crítica; n pontos em que podem haver problemas
de desempenho.
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo em Anel
Algoritmo com um Anel
❑ Hipóteses:
– Os processos não falham e as mensagens não são perdidas
– Utiliza-se os identificadores dos processos para os organizar num anel lógico
❑ Algoritmo:
– Quando o anel é iniciado, associa-se a um processo o testemunho (token)
– O testemunho é trocado entre os processos, sendo passado pela ordem
lógica que define o anel
– Um processo apenas pode entrar numa região crítica quando detém o
testemunho
– Um processo deve passar o testemunho quando sai da região crítica
– O testemunho continua a ser trocado entre os processos ainda que nenhum
deles queira entrar numa região critica
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Exclusão Mútua: Algoritmo em Anel
❑ Ilustração 0
7 1
❑ Propriedades:
5 3
– S1 e V3
4
❑ Discussão:
– Vantagens: o algoritmo é equitativo (satisfaz V3); entre 1 a n mensagens
para se entrar na região crítica
– Problemas: perda do testemunho (token); falha de um processo; requer
envio de mensagens ainda que nenhum processo queira entrar numa
região crítica
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Comparação entre os Algoritmos para Exclusão Mútua
❑ Considerações:
– Mensagens enviadas sequencialmente
– Multicast/Broadcast não é executado
– O algoritmo descentralizado não usa difusão com ordem total
– Decentralized usa m coordenadores e distributed/token considera n processos
© DI-FCUL - Alysson Neves Bessani, Nuno Ferreira Neves, Pedro Ferreira, Reprodução proibida sem autorização prévia
Eleição de Líder
(Capítulo 6, páginas 283-287; 294-297)
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Algoritmos de Eleição
❑ Muitos algoritmos distribuídos necessitam de seleccionar um processo entre
um conjunto de processos iguais: coordenador, líder, primário…
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmos de Eleição
❑ Problema da Eleição de Líder
– Definição:
» Existe um sistema com n processos;
» Cada processo tem um identificador único conhecido por todos;
» Para eleger um líder, cada processo invoca a função lider(), que retorna
o identificador do processo eleito como líder.
– Propriedades:
» Segurança (Safety):
S: Após a execução do algoritmo, existe apenas um líder que é
conhecido por todos.
» Vivacidade (Liveness):
V: A execução do algoritmo termina.
– Solução:
» Por exemplo, localizar o processo correto com o identificador mais alto.
Existem vários algoritmos para fazer essa localização.
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully
Hipóteses:
Existem tempos máximos conhecidos para a comunicação
Os canais são fiáveis
Algoritmo:
Um processo P inicia o algoritmo quando detecta que o líder não responde
(através de um timeout, por exemplo)
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully
Algoritmo (continuação):
3. Quando um processo percebe que vai ser o próximo líder (i.e., não recebe OK de
nenhum processo com identificador superior a ele), envia uma mensagem
COORDENADOR a todos os processos
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully - exemplo
Grupo com 8 processos (identificadores de 0 a 7), onde o processo 4 deteta que o processo 7 falhou.
0 0 0
7 1 a) Envia ELEIÇÃO 7 1 7 1 c) Enviam ELEIÇÃO
6 2 6 2 6 2
5 3 5 3 b) Enviam OK 5 3
4 4 4
0 0
7 1 7 1
6 2 6 2
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully
❑ É importante perceber porque os algoritmos distribuídos funcionam
❑ Implica conhecer, ou desenvolver, os argumentos por trás das provas de correcção
❑ Algoritmos distribuídos são em geral complexos (o nosso cérebro não é muito bom
a lidar com várias “threads”), e apenas através de provas rigorosas podemos ter
certeza que os algoritmos estão correctos
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully
❑ Prova (informal) de correção do Algoritmo Bully (continuação):
Safety (segurança): Após a execução do algoritmo, existe apenas um líder que é
conhecido por todos.
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo Bully
❑ Prova de correção do Algoritmo Bully (continuação):
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo em Anel
Hipótese:
Existem tempos máximos conhecidos para a comunicação
Os canais são fiáveis
Algoritmo:
os processos encontram-se organizados num circulo lógico, o anel
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo em Anel
Algoritmo (continuação):
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo com um Anel
Exemplo: coordenador (5) falhou.
A falha é detectada pelos processos 4 e 1
0
7 1
6 2
5 3
4
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo com um Anel
7 Exemplo: coordenador (5) falhou.
6 A falha é detectada pelos processos 4 e 1
4 0
7 1
6
4 1
6 2
4 2
1
5 3
4 4
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Algoritmo com um Anel
7 Exemplo: coordenador (5) falhou.
6 A falha é detectada pelos processos 4 e 1
4 0
7 0
1 7
6 0 6
4 1 4
7 1
3 3
6 2 2 2
1 1
4 2 0
1 6 7 2
5 3 6
4 4 4
5 3
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Os algoritmos anteriores baseiam-se em hipóteses que não são válidas em redes ad hoc
sem fio (entre outros ambientes)
– Canais fiáveis; topologia fixa;
❑ Vamos ver um algoritmo que não requer essas hipóteses, e portanto pode ser utilizado
em redes ad hoc sem fio (sistemas pervasivos)
❑ O algoritmo consiste em organizar os processos em árvore, sendo o processo que inicia
a eleição a raiz da árvore
❑ Uma importante propriedade deste algoritmo é que o melhor nó (de acordo com
alguma métrica pré-definida) é escolhido como líder
– Exemplo: o nó da rede de sensores que tem maior nível de energia
Nós ao
alcance de e.
Nós ao
alcance de a.
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
(a) O nó a deteta que o antigo líder falhou ou saiu da rede (e.g., moveu-se e não está
mais acessível), e resolve iniciar o algoritmo de eleição de líder
- O critério de escolha do novo líder é a “capacidade” do nó
(b) O nó a envia uma mensagem ELEIÇÃO aos seus vizinhos (b e j) e fica a espera de
confirmações; os nós que recebem essa mensagem definem a como seu pai
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
(c) Os nós b e j enviam mensagens ELEIÇÃO a seus vizinhos (menos ao seu pai - a) e ficam a
espera de confirmações; os nós que recebem a mensagem ELEIÇÃO definem seu
emissor como pai (note que g recebe primeiro de b e portanto ignora a mensagem de j)
(d) Os nós c e g seguem o algoritmo da mesma forma; note que e define g como pai porque
recebe a mensagem ELEIÇÃO dele antes da de c
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Como lidar com eleições iniciadas por processos diferentes?
– Cada mensagem ELEIÇÃO vem com o identificador da eleição (e.g., o
identificador do processo que a iniciou)
– Cada nó executa apenas a eleição com menor identificador, ignorando (ou
parando) a execução das demais
❑ Na prática o que acontece é que a árvore da eleição de menor identificador
sobrepõe-se às de maior identificador
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Como lidar com eleições iniciadas por processos diferentes?
– Cada mensagem ELEIÇÃO vem com o identificador da eleição (e.g., o
identificador do processo que a iniciou)
– Cada nó executa apenas a eleição com menor identificador, ignorando (ou
parando) a execução das demais
❑ Na prática o que acontece é que a árvore da eleição de menor identificador
sobrepõe-se às de maior identificador
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Como lidar com eleições iniciadas por processos diferentes?
– Cada mensagem ELEIÇÃO vem com o identificador da eleição (e.g., o
identificador do processo que a iniciou)
– Cada nó executa apenas a eleição com menor identificador, ignorando (ou
parando) a execução das demais
❑ Na prática o que acontece é que a árvore da eleição de menor identificador
sobrepõe-se às de maior identificador
Ignora eleição de d e
segue a executar a
de a (a < d)
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Como lidar com eleições iniciadas por processos diferentes?
– Cada mensagem ELEIÇÃO vem com o identificador da eleição (e.g., o
identificador do processo que a iniciou)
– Cada nó executa apenas a eleição com menor identificador, ignorando (ou
parando) a execução das demais
❑ Na prática o que acontece é que a árvore da eleição de menor identificador
sobrepõe-se às de maior identificador
Ignora eleição de d e
segue a executar a
de a (a < d)
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Eleição de Líder em Redes Ad Hoc sem Fios
❑ Como lidar com eleições iniciadas por processos diferentes?
– Cada mensagem ELEIÇÃO vem com o identificador da eleição (e.g., o
identificador do processo que a iniciou)
– Cada nó executa apenas a eleição com menor identificador, ignorando (ou
parando) a execução das demais
❑ Na prática o que acontece é que a árvore da eleição de menor identificador
sobrepõe-se às de maior identificador
Ignora eleição de d e
segue a executar a
de a (a < d)
© DI-FCUL - Nuno Ferreira Neves, Miguel Correia, Alysson Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Serviços de Coordenação
(material não está totalmente no livro – ver referências)
Pedro Ferreira
Serviços de Coordenação
❑ Comunicação de grupo é uma forma de concretizar coordenação entre
processos num sistema distribuído
❑ Mas existe outra: serviços de coordenação
Biblioteca de Serviço de
Coordenação Coordenação
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Serviços de Coordenação
❑ Algoritmos distribuídos
– Difíceis de concretizar e não tão simples de usar
– Tipicamente pouco escaláveis
– Modelo de programação pouco intuitivo
❑ Serviços de coordenação
– Construídos com segurança de funcionamento uma
única vez, para depois serem usados por múltiplas
aplicações
» Os engenheiros do serviço tem de ser peritos em
tolerância a faltas e replicação, mas seus
utilizadores (os engenheiros de aplicação) não
– Um único serviço pode ser partilhado por muitas
aplicações
» E.g.: 90k+ clientes comunicavam com uma única
réplica mestre do serviço Chubby, no Google
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Serviços de Coordenação
❑ Os serviços de coordenação permitem
– Programação num modelo parecido com “memória partilhada” (i.e., os
processos comunicam através de uma zona de memória partilhada)
“memória partilhada”
» Muito mais simples de usar que passagem de mensagens
– Concretização de algoritmos de coordenação
» Eleição de líder, consenso, exclusão mútua, etc.
– Armazenar pequenas quantidades de informação de configuração processos a
coordenar
» Entradas relativamente pequenas (e.g., max. 1MB no Zookeeper)
» Espera-se que os dados caibam na memória RAM do servidor
– Permite a deteção de falhas nos clientes do serviço
❑ No entanto, os serviços devem ser concretizados com algum cuidado
– Os algoritmos distribuídos que vimos são usados na replicação do serviço
– Tolerância a faltas e segurança são requisitos fundamentais
– Baixa latência e alta capacidade de processamento também são requisitos
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Serviços de Coordenação
❑ Existem muitos serviços de coordenação
[Link]
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Apache ZooKeeper
❑ Serviço para coordenação de Sistemas Distribuídos
– Fiável, escalável e com alta disponibilidade
– Não serve para guardar dados, mas sim configurações e meta-dados
» Ex: endereços, timestamps e versões
– Usado por outros serviços, como:
» Base dados replicadas, filas de mensagens, blockchains, etc.
❑ Pode ser usado por múltiplos clientes simultaneamente
❑ Está concretizado em Java (código disponível)
❑ Bibliotecas e APIs em múltiplas linguagens
– C, Java, Perl, Python
– Command Line Interface (CLI) também disponível
❑ Concretiza tolerância a faltas através da replicação passiva com primário fixo
– Usando um protocolo de replicação similar ao Paxos chamado Zab
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Modelo de Dados do ZooKeeper
❑ Árvore de Znodes
– Organização hierárquica, i.e., um nó pode ter vários “filhos”
– Mantém os dados em memória, mas também persiste as operações e
nós em armazenamento estável, para recuperar falhas
– Nomes usam notação estilo sistema de ficheiros UNIX
» Ex: /app1/p_1/...
❑ Tipos de Znodes:
– Normal: armazena dados
– Efémero: desaparece se seu criador desconectar
– Sequencial: está associado a um número de
sequência único (e.g., : /app1/p_0000000001)
❑ Znodes suportam watches, que permitem que um
cliente seja notificado da próxima modificação no nó
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
ZooKeeper - Exemplos de Utilização
❑ Gestão de configuração distribuída
– Um cliente cria um Znode c para guardar configuração de um SD
– Outros clientes fazem watch de c
– Quando há uma atualização de c, os clientes são notificados e vão ao
sistema verificar o que mudou
Configuration data
ZooKeeper
Service
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
ZooKeeper - Exemplos de Utilização
❑ Gestão de membros de grupo
– Cliente que faz a gestão do grupo cria um Znode G
– Cada um dos restantes clientes criam um Znode efémero filho de G
– Para receber atualizações sobre o grupo, clientes fazem watch a G
/nodes
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
ZooKeeper - Exemplos de Utilização
❑ Consenso
– Previamente criamos um Znode consenso_id
– Algoritmo em dois passos:
» Um cliente cria um nó sequencial filho de consenso_id com sua
proposta para esta execução do consenso
» Depois lista os Znodes filhos de consenso_id e escolhe o valor
associado ao Znode com menor número de sequência
/consenso_id
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Referências
❑ M. Burrows. The Chubby lock service for loosely-coupled distributed systems.
In Proceedings of the 7th Symposium on Operating Systems Design and
Implementation (OSDI ’06), pages 335–350, 2006.
❑ P. Hunt, M. Konar, F. P. Junqueira, and B. Reed. ZooKeeper: Wait-free
coordination for Internet-scale systems. In Proceedings of the 2010 USENIX
Annual Technical Conference (ATC ’10), pages 145–158, 2010.
❑ T. Distler, C. Bahn, A. Bessani, F. Fischer, and F. Junqueira. Extensible
Distributed Coordination. In Proceedings of the 10th ACM SIGOPS/EuroSys
European Systems Conference (EuroSys’15). Bordeux, France. April 2015.
© DI-FCUL - Alysson Bessani, Bernardo Ferreira. Reprodução proibida sem autorização prévia
Consistência de Dados
(Capítulo 7, páginas 392-405, 406-407, 415-423)
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Porquê Replicar?
❑ Dados são normalmente replicados para melhorar a fiabilidade ou o desempenho
– fiabilidade : como os dados se encontram em várias máquinas, ainda que
uma delas falhe, é possível ir buscar uma cópia a outra
➢ Replicação tolerante a faltas (próximas aulas)
– desempenho :
» se o número de utilizadores crescer muito, a existência de várias réplicas
permite que mais pedidos possam ser tratados por unidade de tempo
➢ Cluster para balanceamento de carga num datacenter
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Replicação e Consistência
❑ A replicação traz o problema da consistência dos dados porque é preciso garantir
que quando uma das réplicas é alterada, as outras também são atualizadas
– Fazer as atualizações em todas as réplicas na mesma ordem pode ser custoso
– Uma solução mais leve será ordenar apenas operações em conflito
» Conflito Escrita-leitura: operações de leitura e escrita concorrentes
» Conflito de escritas: duas escritas concorrentes
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelos de Consistência de Dados
❑ Definição: Modelo de consistência
É um contrato entre os processos/clientes e o sistema que armazena os dados, que
obriga a que os processos obedeçam a um conjunto de regras e que o sistema de
armazenamento se comporte corretamente
❑ Ideias principais
– contrato entre processos e o sistema de armazenamento
– os processos ficam obrigados a fazer os acessos de acordo com certas regras
» Exemplo: “enviar as escritas a todas as réplicas”
– os processos esperam observar um sistema de armazenamento consistente
» Exemplo: “quando se efetua uma leitura obtém-se o valor da última escrita”
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência Estrita, ou Consistência “Perfeita”
❑ Definição
Uma leitura no objeto x devolve o valor da escrita mais recente nesse objeto
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplos de Consistência Estrita
❑ Notação para diagramas de execução
– Considere um sistema com dois processos P1 e P2 que efectuam operações de
leitura e escrita em objetos num sistema de armazenamento replicado
– Cada processo tem acesso a uma réplica local
– Considere um objeto x com o valor inicial NIL
– Um processo ao fazer uma escrita no objeto x com o valor a, W(x)a, começa por
escrever localmente e depois propaga as alteração para as outras réplicas
– Um processo ao fazer uma leitura num objeto x obtém o valor a, R(x)a, faz a leitura
localmente
❑ Exemplo
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência Sequencial
❑ Definição :
O resultado de qualquer execução é equivalente ao que se observaria se as
operações (de leitura e escrita) dos diversos processos fossem executados numa
ordem sequencial e as operações de cada processo individual aparecessem nessa
sequência na ordem especificada pelo seu programa
❑ Ideias principais
– não existe qualquer referência ao tempo
– quando os processos EXECUTAM em paralelo em máquinas diferentes,
qualquer sequencia pela qual são efectuadas as operações é válida, desde
que todos os processos “vejam” a mesma sequência
– essa sequência deve respeitar a ordem local dos eventos nos processos
– Um processo só “vê” as suas leituras
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo 1: Consistência Sequencial
NOTA: Fazer x=1 é
Existem 2 processos P1 P2 equivalente a executar
que se vão executar x = 1; print(x); um W(x)1, e fazer print(x)
em paralelo print(x); é equivalente a R(x)a
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo 2: Consistência Sequencial
Exemplo
Existem 4 processos P1 P2 P3 P4
que se vão executar x = 1; x = 2; print(x); print(x);
em paralelo print(x); print(x);
P1 W(x)1
Sistema de armazenamento que verifica a
P2 W(x)2
consistência sequencial, embora não verifique a
P3 R(x)2 R(x)1 consistência estrita
P4 R(x)2 R(x)1
Exemplo de sequência: W(x)2, R(x)2, R(x)2, W(x)1, R(x)1, R(x)1
P1 W(x)1
W(x)2 Sistema de armazenamento que não verifica a
P2
consistência sequencial e não verifica a
P3 R(x)2 R(x)1
consistência estrita
P4 R(x)1 R(x)2
Exemplo de sequência: ???
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência Causal
❑ Definição :
Escritas potencialmente relacionadas em termos de causa/efeito devem ser vistas
por todos os processos do sistema na mesma ordem. Escritas concorrentes
podem ser vistas em ordens diferentes por processos diferentes.
❑ Ideias principais
– Para suportar estas relações de causa/efeito temos de construir um grafo de
dependência entre as diversas escritas no sistema
» Isto pode ser feito através de relógios lógicos, mais precisamente, de
vetores de relógios (assunto do mestrado)
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo: Consistência Causal
concorrentes
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo: Consistência Causal
concorrentes
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Exemplo: Consistência Causal
concorrentes
concorrentes
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Agrupamento de Operações
❑ Os modelos já vistos são definidos em termos de operações de leitura e escrita em itens de
dados, que são operações de baixo nível
– Definidos originalmente para sistemas multiprocessador com memória partilhada
❑ O nível de granularidade destas operações muitas vezes não corresponde ao nível usado
nas aplicações distribuídas
– Muitas vezes a operação a ser executada num servidor replicado consiste numa série
de leituras e escritas atómicas (ex. append numa fila, incremento de um contador)
– Como fazemos para fazer com que uma sequência de operações seja executada de
uma forma atómica?
» Transacções (sequência de operações entre begin e commit)
» Exclusão mútua (sequência de operações entre enter() e exit())
– Em geral, o que se faz é associar locks a cada uma das variáveis: operações de leitura e
escrita só são executadas na variável enquanto o processo tem o lock
– A partir daí, as sequências de operações satisfazem algum modelo de consistência.
Não adquiriu
lock Ly.
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Consistência para utilizadores móveis
Exemplo:
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelos de Consistência Centrados no Cliente
❑ Modelos de consistência centrados no cliente fornecem garantias de
consistência a cada utilizador em relação aos acessos efetuados por ele
❑ Além de lidar bem com a mobilidade, este tipo de modelo de consistência é útil
para sistemas em que
– a conectividade da rede é pouco fiável (algumas réplicas podem ficar
temporariamente desligadas das outras)
– potencialmente podem surgir problemas de desempenho (a ligação entre
duas réplicas pode ficar muito lenta)
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência Eventual
❑ Modelos centrados nos clientes e sem conflitos de escrita usualmente satisfazem algum
modelo de consistência eventual, cujas principais características são:
– as escritas são (tipicamente) feitas apenas por uma entidade, o dono dos dados
– as réplicas dos dados não precisam de ser atualizadas imediatamente após uma
alteração, mas apenas gradualmente ao longo do tempo
❑ Mesmo com este tipo de consistência, podemos ainda dar algumas garantias fracas:
– Leituras Uniformes
– Escritas Uniformes
– Ler suas Escritas
– Escritas seguem Leituras
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência das Leituras Uniformes (Monotonic-Reads)
❑ Definição :
Se um processo lê um valor de um objeto x, então qualquer posterior operação
de leitura nesse mesmo objeto por esse processo deverá sempre devolver o
mesmo valor ou um valor mais recente
❑ Ideias principais
– se um processo vê um valor de x no instante t, então ele nunca verá um
valor mais antigo num instante posterior
– é útil por exemplo num sistema de armazenamento de emails em que a
mailbox dos utilizadores se encontra replicada em várias máquinas
» os novos emails podem ser adicionados em qualquer uma das réplicas;
» serão propagados para as restantes réplicas ao longo do tempo;
» ainda que o utilizador mude de localização, ele observa sempre uma
versão igual ou mais recente da mailbox
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência das Escritas Uniformes (Monotonic-Writes)
❑ Definição :
Uma operação de escrita de um processo num objeto x é terminada antes de
qualquer outra operação de escrita posterior executada pelo processo em x
❑ Ideias principais
– As escritas executadas por um processo num objeto x são executadas na
mesma ordem FIFO nas diversas cópias de x
– é útil por exemplo na aplicação de patches numa biblioteca do servidor
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência “Ler suas Escritas” (Read Your Writes)
❑ Definição :
O efeito de uma operação de escrita de um processo num objeto x, será sempre
visto pelas operações de leitura posteriores executadas por este processo em x
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Modelo de Consistência “Escritas seguem Leituras”
(Writes Follow Reads)
❑ Definição :
Uma operação de escrita de um processo num objeto x que sucede uma leitura
executada por este processo em x vai afetar sempre o mesmo valor lido pelo
processo ou um valor mais recente
❑ Exemplo: para a resposta a posts num forum, é preciso que toda gente primeiro
receba o post original para depois receber sua resposta
© DI-FCUL, Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira. Reprodução proibida sem autorização prévia
Gestão de Replicação
(Capítulo 7, páginas 423-434, 437-443)
Aplicações Distribuídas
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Onde Colocar as Réplicas e os Conteúdos
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Propagação das Actualizações
❑ O que propagar?
1. notificação da modificação do dado
2. transferir uma cópia do dado actualizado
3. propagar a operação que fez a actualização
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Propagação das Actualizações (cont.)
❑ Quem inicia as propagações?
– servidor
– cliente
❑ Protocolos push: o servidor inicia a propagação das atualizações (ou a réplica
mais atualizada contacta as outras réplicas); usado quando:
– é necessário um elevado nível de consistência
– número de leituras é muito superior às escritas rácio leituras/escritas alto
– as réplicas são usadas por muitos clientes
❑ Protocolos pull: o cliente inicia (ou as réplicas desatualizadas contactam as mais
atualizadas); usado quando:
– relação do número de leituras para as escritas é reduzida
– cada réplica é usada por um único cliente (ex., cache)
❑ Protocolos push provisório (lease): o servidor só envia as actualizações durante
um período; passado esse intervalo, o cliente volta a pedir mais tempo ou passa
a usar um método de pull
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Realização da Replicação com os
Modelos de Consistência
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Protocolos de Consistência
❑ O protocolo de consistência é responsável por manter as réplicas atualizadas de
acordo com um dado modelo de consistência
NOTA: como em ambos os casos apenas existe uma cópia a ser actualizada,
é fácil garantir a consistência sequencial
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação Passiva com Primário Fixo – Escrita Remota
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação Passiva com Primário “Móvel” – Escrita Local
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação Passiva
❑ O que acontece se o primário falha? (por paragem)
❑ A falha tem de ser detectada:
– por ex., cada T unidades de tempo envia uma mensagem “estou vivo” aos
secundários;
– estes, quando não a recebem percebem que o primário falhou
❑ Tem de se eleger um novo primário
– Usando um protocolo de eleição (Bully, Anel, etc…)
❑ Garantir a coerência do estado dos servidores não é simples
– Como saber se uma actualização enviada pelo servidor foi executada por
todas as réplicas?
– Como garantir que um pedido de um cliente não se perde?
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Protocolos com Escritas Replicadas
❑ As actualizações podem ser efectuadas num conjunto de réplicas
– Replicação activa
» cada actualização (escrita) tem de ser enviada para todas as réplicas
– Replicação com um Quórum (quórum = subconjunto das réplicas)
» as actualizações têm de ser feitas apenas em quóruns, não em todas as
réplicas
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação Activa
❑ Solução genérica para serviços tolerantes a faltas
❑ Propriedades básicas:
– Todas as réplicas começam no mesmo estado
– Todas as réplicas têm de executar os mesmos
comandos pela mesma ordem
– As réplicas têm de ser deterministas: o
mesmo comando executado no mesmo
estado em réplicas diferentes têm de
modificar o estado da mesma forma
❑ Se uma réplica falha não é preciso fazer nada
❑ Se uma falha e reinicia é preciso fazer uma
transferência de estado
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação com Quóruns
❑ Quórum = conjunto de réplicas
❑ Sistemas de Quóruns são sistemas replicados onde as operações de leitura e escrita
são executadas em subconjuntos (quóruns) de réplicas
– Funciona devido às intersecções não vazias entre os quóruns usados nas
operações de leitura e escrita
❑ Nos sistemas de quóruns, cada objeto x tem um número de versão associado
❑ Antes de se poder efectuar uma escrita ou leitura é preciso obter um conjunto de
votos das réplicas; por exemplo (quóruns maioritários):
– escrita: contactam-se e atualizam-se metade mais uma das réplicas; essas réplicas
passam a ter um número de versão maior que o anterior
– leitura: contactam-se metade mais uma das réplicas e obtém-se os seus números de
versão; faz-se a leitura da réplica com o maior número de versão
– (os quóruns podem ter diferentes construções para além de maiorias)
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Replicação com Quóruns (cont.)
Em geral tem-se
N = número total de réplicas
NR = número de elementos do quórum de leitura
NW = número de elementos do quórum de escrita
Precisamos de garantir que:
1) NR + NW > N (existe sempre intersecção entre quóruns de leitura e de escrita)
2) NW > N/2 (sempre há intersecção entre quóruns de escrita)
© Nuno Ferreira Neves, Miguel Correia, Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Reprodução proibida sem autorização prévia
Tolerância a Faltas
(Capítulo 8, páginas 461-471, 474-477 e 491-495)
Aplicações Distribuídas
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Motivação
❑ Muitos sistemas informáticos são críticos, não podem falhar:
– aviões, centrais nucleares, navios, comboios, foguetões, etc.
❑ Outros, não sendo tão críticos têm custos de falha elevados:
– servidores de web, bases de dados empresariais, etc.
❑ Como fazer com que os sistemas detetem, tolerem ou recuperem destas falhas
parciais, de tal forma que o sistema como um todo não seja comprometido?
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Noções Fundamentais
❑ Em geral, o que se quer num sistema critico é que ele nos forneça as
propriedades de segurança de funcionamento
– Intuição: o sistema deve cumprir o que promete, mesmo em caso de falhas
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Faltas e Suas Manifestações
❑ Falta (fault): causa do estado incorrecto do sistema (software ou hardware):
avaria de um dos componentes, interferência do ambiente, desenho incorrecto
(uma tempestade interfere com a transmissão de dados numa rede)
❑ Erro: manifestação da falta no programa ou na estrutura de dados
(uma mensagem fica com alguns bits alterados)
❑ Falha (failure): ocorre quando o serviço prestado pelo sistema difere da
especificação
(o programa lê a mensagem e toma decisões incorretas)
Tipos de Faltas
❑ Permanentes: falta que se mantém ao longo do tempo
❑ Intermitentes: falta que ocorre ocasionalmente devido a instabilidades no
hardware ou software
❑ Transitórias: falta que resulta de condições temporárias, e que em princípio
não se vai repetir outra vez
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Tipos de Falhas
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Redundância
❑ Redundância é o meio pelo qual podemos tolerar faltas no sistema
❑ Existem vários tipos de redundância:
– Redundância de informação
» Incluem-se bits extra nos blocos de dados de tal forma que, mesmo que alguns
bits do bloco sejam corrompidos, ainda é possível ler o bloco correctamente
» Exemplo: códigos de apagamento (armazenamento de dados), códigos de
hamming (mensagens)
– Redundância temporal
» Repetir uma acção várias vezes para ter certeza que ela tem efeito
» Exemplo: numa rede com comunicação não fiável, enviar uma mensagem
várias vezes dá maior probabilidade que ela chegue ao destino
– Redundância física
» Usar várias réplicas de um mesmo componente físico
» Exemplo: utilizar vários servidores para disponibilizar um serviço critico no
sistema distribuído
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Redundância: Exemplo de redundância física
❑ TMR: Triple Modular Redundancy
– Muito usado no projeto de hardware fiável
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Sistemas Replicados Tolerantes a Faltas
❑ Um sistema replicado com n replicas é dito tolerante a k faltas se funciona bem
mesmo que k das n réplicas falhem
❑ Qual o n para que um sistema seja tolerante a k faltas?
– Falhas por crash
» Basta que uma réplica envie respostas corretas aos clientes
» Devemos ter pelo menos n ≥ k + 1
– Falhas Arbitrárias (Bizantinas)
» A maioria das réplicas devem ser corretas para evitar que as que falham
enviem o mesmo valor incorreto de uma resposta
» Devemos ter pelo menos n ≥ 2k + 1
❑ Este tipo de solução tipicamente requer uma primitiva de difusão com ordem
total para disseminar as mensagens às réplicas
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Acordo em Sistemas Sujeitos a Faltas
❑ Problemas de acordo:
– Ordenação de mensagens
– Sincronização (exclusão mútua, sincronização de relógio, eleição de líder)
– Confirmação de transacções distribuídas
❑ Considerações
– Em sistemas livres de falhas (como os considerados nas aulas sobre sincronização)
existem algoritmos relativamente simples e eficientes para fazer com que os
processos estabeleçam acordo
– Num ambiente com falhas por crash, mas onde os canais de comunicação são fiáveis
e apresentam latência de transmissão limitada e conhecida, o problema tem
solução razoavelmente simples
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Consenso baseado em Flooding
❑ Algoritmo funciona em rondas. Em cada ronda, cada processo:
– Envia a lista propostas que conhece (inicialmente só a sua) a todos
– Espera propostas dos outros e as junta a sua lista:
» Se recebeu de todos:
➢ Decide deterministicamente por uma delas
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Consenso baseado em Flooding (cont.)
❑ Note que...
– processos que recebem a mesma lista, decidem o mesmo valor
– todos decidem numa ronda sem falha
– se alguém decide e não falha, todos vão receber sua lista na ronda seguinte
❑ Só funciona porque conseguimos detetar falhas perfeitamente (impossível na
internet). Existem algoritmos que eliminam esse requisito, como o Paxos.
Funcionamento do Paxos…
cada processo pode correr
três algoritmos (papéis).
(ver detalhes no livro)
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Problema dos Generais Bizantinos
❑ Problema fundamental de acordo em sistemas distribuídos
– A formulação do problema foi motivada pelas dificuldades em se construir um
algoritmo que consolidasse sinais de sensores de altitude replicados num avião
❑ Descrição do problema:
– Um grupo de n generais bizantinos estão acampados ao redor da cidade inimiga;
– Os generais leais devem decidir um plano de ação comum;
– Alguns dos generais podem ser traidores (t);
– Todos os generais podem comunicar com todos os outros;
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2k+1 processos
❑ Formalização (um general envia a sua ordem aos outros):
– Considere n generais, sendo um deles o comandante. O comandante
precisa enviar uma ordem para os n-1 subordinados de tal forma
que:
» IC1: Todos os subordinados leais obedecem a mesma ordem;
» IC2: Se o comandante é leal, então todos os generais subordinados leais
obedecem a ordem enviada por ele.
– (O problema dos generais consiste em n execuções em paralelo
desse algoritmo: cada general dissemina sua opinião de forma
confiável e, decidem a ação por voto)
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2k+1 processos
❑ Com mensagens orais nenhuma solução é possível a não ser que
mais de dois terços dos generais sejam leais, i.e., com três generais
(um deles traidor) nenhuma solução é possível (como veremos nos
próximos slides)
– Isto contradiz nossa idéia de que n ≥ 2k +1 seriam suficientes
– Simplificações:
» Consideraremos apenas dois planos de ação: ‘atacar’ ou
‘retirar’;
» Se o general não enviar nada em tempo util assume-se que
ele vota por uma ‘retirada’ (talvez ele até já se tenha
retirado…).
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
atacar atacar
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
atacar atacar
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
atacar atacar
traidor
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
atacar atacar
IC1: Todos os subordinados leais obedecem a mesma ordem;
IC2: Se o comandante é leal, então todos os generais subordinados leais
traidor
obedecem a orem enviada por ele
ele disse ‘retirada’
subordinado 1 subordinado 2
Atacar ou retirada?
Como satisfazer IC2 Caso 1: subordinado Traidor
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
comandante
atacar
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
traidor
comandante
atacar retirada
subordinado 1 subordinado 2
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
traidor
comandante
atacar retirada
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
traidor
comandante
atacar retirada
IC1: Todos os subordinados leais obedecem a mesma ordem;
IC2: Se o comandante é leal, então todos os generais subordinados leais
obedecem a orem enviada por ele
ele disse ‘retirada’
subordinado subordinado
1 2
ele disse ‘atacar’
Atacar ou retirada?
Como satisfazer IC1? Caso 2: comandante traidor
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Generais Bizantinos: Impossibilidade com 2t+1 processos
traidor
comandante
atacar retirada
Nos dois casos o
subordinado 1 não
sabe quem é o
traidor e nem qual
das ordens seguir.
ele disse ‘retirada’
subordinado 1 subordinado 2
Atacar ou retirada?
Como satisfazer IC1? Caso 2: comandante traidor
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
❑ Com n ≥ 2k+1 não é possível resolver… agora vamos usar n ≥ 3k+1
❑ Assume-se que:
Limite comum em
– A1: Toda a mensagem enviada é entregue corretamente; tolerância a faltas
Bizantinas
(canal fiável)
– A2: O receptor da mensagem sabe quem a enviou;
(id do emissor na mensagem)
– A3: A perda de uma mensagem pode ser detetada (‘retirada‘ é a ordem em
caso de omissão).
(sistema síncrono)
❑ Observações:
– A1 e A2 previnem que um traidor interfira na comunicação entre dois
generais leais;
– A3 previne que um traidor atrapalhe o acordo por omissão.
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
maioria(v1 , v2, ..., vn-1)
o valor vi que é maioria (se existir) ou ‘retirada’.
Algoritmo OM(0):
1. O comandante envia a sua ordem a todos os subordinados.
2. Cada subordinado usa o valor recebido pelo comandante ou ‘retirada’ caso não
receba nenhum valor.
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais e k = 1
comandante
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais e k = 1
Passo 1:
OM(1)
comandante
v v
v
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais e k = 1
Passo 1:
OM(1)
comandante
v v
v
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais e k = 1
Passo 1:
OM(1)
comandante
v v
v
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
Passo 1:
OM(1)
comandante
v v
v
Passo 1:
subordinado 1: (v,v,x) = v OM(1)
subordinado 2: (v,v,x) = v comandante
subordinado 3: traidor
v v
v
comandante
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
Passo 1:
OM(1)
comandante
x y
y
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
Passo 1:
OM(1)
comandante
x y
y
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
Passo 1:
OM(1)
comandante
x y
y
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Algoritmo com Mensagens Orais
Passo 1:
subordinado 1: (x,y,y) = y OM(1)
subordinado 2: (x,y,y) = y comandante
subordinado 3: (x,y,y) = y
x y
y
Passo 1:
subordinado 1: (x,y,z) = ? OM(1)
subordinado 2: (y,x,z) = ? comandante
subordinado 3: (z,x,y) = ?
x z
y
Passo 1:
subordinado 1: (x,y,z) = retirar OM(1)
subordinado 2: (y,x,z) = retirar comandante
subordinado 3: (z,x,y) = retirar
x z
y
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Condições para Acordo em Sistemas Distribuídos
❑ Deteção de falhas perfeita (i.e., sem suspeitas erradas) é impossível na prática!
– Usamos timeouts para detetar falhas, mas isso permite que suspeitemos de
processos que ainda estão vivos em caso de problemas na rede.
❑ Isso faz com que o problema de consenso tolerante a faltas (mesmo de crash) não
admita solução em todo tipo de ambiente. As hipóteses consideradas são:
– Tempos de processamento e comunicação são limitados ou não?
– As mensagens são ordenadas pela rede?
– A comunicação fiável é feita por unicast ou multicast?
Message ordering
Unordered Ordered
Communication delay
Process behavior
X X X X Bounded
Synchronous
X X Unbounded
X Bounded
Asynchronous
X Unbounded
Unicast Multicast Unicast Multicast
Message transmission
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Condições para Acordo em Sistemas Distribuídos
❑ Deteção de falhas perfeita (i.e., sem suspeitas erradas) é impossível na prática!
– Usamos timeouts para detetar falhas, mas isso permite que suspeitemos de
processos que ainda estão vivos em caso de problemas na rede.
❑ Isso faz com que o problema de consenso tolerante a faltas (mesmo de crash) não
admita solução em todo tipo de ambiente. As hipóteses consideradas são:
– Tempos de processamento e comunicação são limitados ou não?
– As mensagens são ordenadas pela rede?
– A comunicação fiável é feita por unicast ou multicast?
Message ordering
Unordered Ordered
Communication delay
Process behavior
X X X X Bounded
Synchronous
X X Unbounded
Mundo X Bounded
Asynchronous
Real X Unbounded
Unicast Multicast Unicast Multicast
Message transmission
© Alysson Neves Bessani, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia.
Condições para Acordo em Sistemas Distribuídos
❑ Deteção de falhas perfeita (i.e., sem suspeitas erradas) é impossível na prática!
– Usamos timeouts para detetar falhas, mas isso permite que suspeitemos de
processos que ainda estão vivos em caso de problemas na rede.
❑ Isso faz com que o problema de consenso tolerante a faltas (mesmo de crash) não
admita solução em todo tipo de ambiente. As hipóteses consideradas são:
– Tempos de processamento e comunicação são limitados ou não?
– As mensagens são ordenadas pela rede?
– A comunicação fiável é feita por unicast ou multicast?
Message ordering
Unordered Ordered
Communication delay
Process behavior
X X X X Bounded
Synchronous
X X Unbounded
Mundo X Bounded
Asynchronous
Real X Unbounded
Aplicações Distribuídas
Departamento de Informática
Faculdade de Ciências da Universidade de Lisboa
Tópicos
❑ Semântica do RPC na Presença de Falhas
– Falhas no Servidor
– Falhas no Cliente
❑ Confirmação Atómica
– Confirmação em Duas Fases (2PC)
– Confirmação em Três Fases (3PC)
❑ Recuperação de Falhas
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Semântica do RPC na Presença de Falhas
❑ Na presença de falhas, é difícil fazer com que as RPCs tenham um
comportamento idêntico às chamadas locais a procedimentos
CLIENTE SERVIDOR
❑ Tipos de falhas a considerar
1
1. cliente não consegue localizar o servidor
2. perda da mensagem/pedido do cliente para Invocação do
Procedimento
o servidor Remoto
3. o servidor tem uma falha após ter recebido 2
a mensagem com o pedido Execução do
3 Procedimento
4. perda da resposta do servidor Remoto
5. o cliente tem uma falha antes de receber a resposta 5 4
Devolução dos
Parâmetros de
Resposta
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Erro na Localização do Servidor (1)
❑ Causas do erro
– servidor não está em execução
– o cliente tenta utilizar uma versão desatualizada da rotina de adaptação
– servidor de nomes tem algum problema
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Perda do Pedido/Resposta (2, 4)
❑ Perda do pedido do cliente
– utilizar os mecanismos típicos para recuperação da perda de mensagens
» temporizadores e retransmissão + ACKs (como no TCP)
– existe o problema de se considerar erroneamente que o servidor está em baixo
por se perderem muitas mensagens
Problema: Será que o pedido pode ser executado mais do que uma vez ?
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Perda do Pedido/Resposta (2, 4)
❑ Perda do pedido do cliente
– utilizar os mecanismos típicos para recuperação da perda de mensagens
» temporizadores e retransmissão + ACKs (como no TCP)
– existe o problema de se considerar erroneamente que o servidor está em baixo
por se perderem muitas mensagens
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Perda do Pedido/Resposta (2, 4)
❑ Perda do pedido do cliente
– utilizar os mecanismos típicos para recuperação da perda de mensagens
» temporizadores e retransmissão + ACKs (como no TCP)
– existe o problema de se considerar erroneamente que o servidor está em baixo
por se perderem muitas mensagens
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
O Servidor tem uma Falha (3)
❑ O maior problema consiste na determinação do momento em que a falha ocorre
recebe_pedido() recebe_pedido()
Caso 1 crash Caso 2 executa_pedido()
crash
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
O Cliente tem uma Falha (5)
❑ A falha de um cliente após ter enviado um pedido ao servidor pode criar
computações órfãs, que podem resultar:
– no desperdício de ciclos do processador
– na inutilização temporária de certos recursos (e.g., lock de um ficheiro)
❑ Soluções
1. exterminação: antes do cliente enviar o pedido, anota-o num log em disco;
quando o cliente é reiniciado os órfãos são terminados
» escrever em disco por cada RPC pode ser caro; partição da rede pode
impedir a terminação dos órfãos
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
O Cliente tem uma Falha (5) (cont.)
❑ Soluções (cont.)
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação Atómica: Motivação
❑ Transacções Distribuídas
– Execução de uma sequência de operações em recursos distribuídos
– Com atomicidade, i.e., ou se executa tudo, ou não se executa nada
Exemplo
begin_transaction
reservar Lisboa → Londres /* base de dados da companhia A */
reservar Londres → NY /* base de dados da companhia B */ A B
end_transaction
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
O Problema da Confirmação Atómica
❑ Este problema pode ser generalizado como um problema de acordo em que
todas as máquinas confirmam ou não a execução da ação
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
O Problema da Confirmação Atómica
Participante
Coordenador
Cliente (A)
(coordena os participantes
(Executa transação T)
na transação)
Participante
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Uma Fase (One-Phase Commit)
❑ Num sistema sem falhas basta que uma máquina coordenadora envie uma
mensagem definindo se haverá commit ou abort da transação. Essa máquina
repetirá a mensagem até que todos os participantes confirmem a operação.
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases (Two-Phase Commit)
❑ Quando a possibilidade de falhas existe, o protocolo de confirmação em duas-
fases (2PC - Two-Phase Commit) resolve este problema, sendo uma das
máquinas a coordenadora e as restantes participantes
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases (Two-Phase Commit)
❑ Na primeira fase o coordenador pede aos participantes um voto por commit ou
abort.
– Depois de votar por commit, um participante não pode decidir abort.
– Antes de votar commit o participante assegura que pode executar a sua parte
da transação, mesmo em caso de falha.
– Um participante diz-se no estado preparado se eventualmente puder fazer
commit.
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases (cont.)
• pelo menos 1 abort => abort
coordenador participante • todos commit => commit
coordenador
vote-request
Fase 1
Inicial
Commit / Vote-request
vote-commit Vote-abort /
vote-abort Escrever Esperar Vote-commit /
Escrever no Log Abortar Global
Confirmar Global
no Log
Abortar Confirmar
confirmar ou
abortar global
participante
Inicial
ack Vote-request /
Abortar Vote-request / Vote-commit
Fase 2
Preparado
Confirmar Global / ACK
Abortar Global
Abortar / ACK Confirmar
Problema: quando há falhas podem bloquear!
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases: Acções do Coordenador
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases: Acções dos Participantes
write INIT to local log;
wait for VOTE_REQUEST from coordinator;
Espero pelo
if timeout {
write VOTE_ABORT to local log; Pedido de Voto
exit;
}
if participant votes COMMIT {
write VOTE_COMMIT to local log; Envio Voto e
send VOTE_COMMIT to coordinator; Espero Decisão
wait for DECISION from coordinator;
if timeout { Commit
multicast DECISION_REQUEST to other participants;
wait until DECISION is received; /* remains blocked */
write DECISION to local log; PODE BLOQUEAR!
} Se o coordenador tiver uma
if DECISION == GLOBAL_COMMIT falha antes de enviar a decisão
write GLOBAL_COMMIT to local log; aos participantes, estes podem
else if DECISION == GLOBAL_ABORT ficar bloqueados à espera da
write GLOBAL_ABORT to local log;
recuperação do coordenador.
} else {
write VOTE_ABORT to local log;
send VOTE ABORT to coordinator; Abort
}
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Duas Fases: Acções dos Participantes
Tratamento de Pedidos de Decisão
Inicial
Vote-request /
Abortar Vote-request / Vote-commit
Preparado
Confirmar Global / ACK
Abortar Global
Abortar / ACK Confirmar
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Confirmação em Três Fases (Three-Phase Commit)
❑ Confirmação em duas fases pode não chegar a qualquer decisão se o
coordenador falhar em certas fases da execução do protocolo, o que resulta no
bloqueio dos participantes até que o coordenador recupere
– isto ocorre quando os participantes enviaram o seu voto mas ainda não
obtiveram o resultado da votação do coordenador (todos em READY)
❑ O protocolo de confirmação em três fases evita esses bloqueios
Se der timeout neste estado,
Coordenador Participante pode abortar pois ninguém
ainda fez commit
Se der timeout pode fazer
commit pois o coordenador
já decidiu commit.
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia
Recuperação de Falhas
❑ Armazenamento estável
– Se queremos recuperar falhas, temos de armazenar (parte do) estado do
processo em memória estável (disco)
– Desta forma, após uma falha, as operações interrompidas podem ser
retomadas
❑ Duas técnicas complementares:
– Checkpointing: armazenamento periódico do estado do sistema
– Logging de mensagens: armazenamento de mensagens entre checkpoints
© Alysson Neves Bessani, Nuno Ferreira Neves, Miguel Correia, Pedro Ferreira – DI/FCUL. Publicação, reprodução, ou difusão, proibidas sem autorização prévia