2. What is a Distributed System?
Patricio Galeas edited this page on August 4, 2017.
I) Introduction
A distributed system is one in which its components are located in
Networked computers communicate and coordinate their actions through messages.
These systems are characterized by:
Concurrent Operation: That is, the execution of programs
each component of the system can be simultaneous, for this reason the
coordination of tasks between components that share or have
shared resources is a recurring theme.
There is no global definition of time: When the programs
they need to cooperate, they coordinate their actions through messages. The
specific coordination often depends on a shared idea
about the time in which the program tasks are executed. Without
however, this coordinated notion of time is not always present
as part of a distributed system.
Independent faults: All computational systems can
to fail, and in the case of distributed systems, there are new ways to
failures that add to these individual failures. For example, the failures in the
network can cause the isolation of the computers from the system, without
embargo this does not mean that each of the computers
involved stopped working and sometimes it's difficult for these programs
to determine whether the network has failed or is just functioning
slower than normal. On the other hand, the failure of a component is not
immediately known by the other components of the system that
they cooperate with him.
The main motivation of a distributed system is to share resources, which
which includes hardware resources such as printers, units of
storage, etc. and software, such as files, databases, and objects.
Lecturas recomendadas:
A brief introduction to distributed systems
II) Examples of Distributed Systems
As can be seen, distributed systems are part of the most
recent and significant technical developments of the last few years and are
associated with a wide range of applications in our daily routine, starting from
from more localized systems like the ones we find in our
cars or airplanes to global scale systems that involve thousands
of nodes, from data-centric services to tasks that require high
computation, from small systems based on simple sensors to systems
that allow advanced data processing, from embedded systems
to systems that provide sophisticated interfaces for the user.
II.1) Internet Search Engines
Search engines for the web is the industry that has grown the most in
last decade. Some recent statistics indicate that the number of
Global searches have increased to more than 10 billion per month.
the main task of a search engine is to index the content of the web, the
which contains a wide range of resources of different types and styles, which
translates to an approximate amount of about 63 billion
pages and a trillion web addresses (URLs). Given this scenario and
considering that most search engines try to analyze the
total content of the web, it becomes necessary to have mechanisms for
sophisticated processing that can operate on this huge database,
which translates into a great design challenge for a distributed system.
Google, the leader in web search technologies, has invested heavily
resources for the design of a distributed infrastructure that supports its
search applications. This is currently one of the largest and
complex distributed systems in the history of computing. The main ones
components of your infrastructure include:
Its hardware platform consists of a large number of
networked computers located in different parts around the
world.
A distributed file system designed to support files of
large size, which is also specially optimized for
operate with the different search applications.
A related storage system that provides fast access to
large datasets.
A locking system that offers distributed locking functions and
negotiation.
A programming model that supports task management
mass computing for a parallel and distributed architecture across
on its hardware platform.
Recommended readings:
Designing Distributed Systems: Google Case Study
II.2) Online Games
Massively Multiplayer Online Games (MMOG)
they offer an immersive experience where a large number of users
connected through the Internet interact in a virtual world. Some
emblematic examples of these games are Sony's EverQuest and EVE Online
the Finnish company CCP Games. These virtual worlds have increased
considerably in sophistication and currently offer high scenarios
complexity. The number of followers of these games has also grown, and the
systems that support this type of games must handle loads of up to
50,000 simultaneous online players. The engineering of the MOOGs represents
also a great challenge for distributed systems technologies,
especially due to the need for quick response times for
guarantee a good user experience. Another important challenge has to
see with the propagation of real-time events to multiple users, in order to
maintain a consistent image of shared virtual worlds. Like
It can be appreciated that this type of system is a good example of the challenges.
what a designer of modern distributed systems faces.
Recommended readings:
Massive Multiplayer Online Game Architectures
II.3) Financial Systems
The finance industry has been a market characterized by its use of
the most modern distributed systems, due to their requirements of
real-time operation and with different sources of information. In addition
this industry includes the use of automated applications for commerce and
monitoring. This type of applications focuses on communication and
processing of items of interest, known as ‘events’ in the
terminology of distributed systems, which must be transmitted to
time and reliably to a large number of interested users in
this information. Some examples of these events are the price drop of
the actions, the release of the latest unemployment data, etc. This type of
requirements need a different architecture than the mentioned ones
previously. Financial systems commonly employ what is known
as distributed event-based systems. Figure 2 shows a typical
financial trading system, where a series of events enter a
financial institution. These events share the following characteristics: (1)
Information sources come in different formats and are generated from
of different technologies, which illustrates the problem of heterogeneity
present in most distributed systems. Figure 2 also
show the use of adapters that translate from formats
heterogeneous towards a common internal format. (2) The trade system must
handle a wide variety of event flows, all entering the system to
a great speed, which often requires real-time processing
real, in order to detect patterns of potential business opportunities. Without
embargo, these processes are usually carried out manually, but the large
market competition will lead to the development of automated services
known as Complex Event Processing or Complex Event
Processing (CEP), which will offer a way to package events.
concurrent in logical, temporal, or spatial patterns. This technology is
mainly used to program commercial strategies in the form of
algorithms for buying and selling stocks, particularly monitoring the
patterns that indicate business opportunities and responding appropriately
automatically executing orders to make transactions.
Recommended readings:
Evolution and Practice: Low-latency Distributed Applications in Finance
Algorithmic Trading System Architecture
Figure 2: A typical financial trading system.
II.4) Space Observatories (ALMA Common
Software
complete
III) Trends of the Systems
Distributed
Current distributed systems are going through a phase of changes.
significant associated mainly with the following trends:
The emergence of pervasive networking technologies.
The emergence of ubiquitous computing.
The growing demand for multimedia services
The vision of distributed systems as a commodity
III.1) Pervasive Networks: The Modern Internet
The modern Internet is a massive collection of different networks of
interconnected computers, whose heterogeneity is increasing
including a wide range of technologies such as WiFi, WiMAX,
Bluetooth and mobile phone networks. The result is a network that has
transformed into a pervasive resource where the devices of
communication can be permanently connected regardless
of its location. Figure 3 illustrates a typical network configuration that
is part of the Internet. The programs that run on this network
interact through the sending of messages using a common medium of
communication. The design and construction of communication mechanisms of
Internet (Internet protocols) is a great technical challenge as they allow
that applications from a program running anywhere on the network
can send messages to another program located in a different place and
regardless of the thousands of technologies involved in the
communication. The Internet is certainly also a huge distributed system,
which allows its users to use the web, email,
file transfer (ftp), regardless of where they are located.
And these services can be extended by adding, for example, new
servers and/or new types of services. The figure shows a set of
intranets operated by companies and other organizations which are
normally protected through Firewalls, which prevents unauthorized messages
authorized to enter or leave the network. This type of intranets are usually
connected through Backbones, which is simply a connector of
networks that have a high transmission capacity that are normally used
satellite connections, fiber optic cables or other circuits that have a
large bandwidth.
Figure 3: A typical network that is part of the Internet.
III.2) Mobile and Ubiquitous Computing
The recent technological advances in the miniaturization of devices and networks.
Wireless technologies have led to the integration of a wide range of devices.
small laptops towards a distributed system. These devices include:
laptops, mobile phones, smartphones, watches
smart, GPS devices, etc. Similar devices also exist.
incorporated into appliances such as washing machines, music equipment,
automobiles and refrigerators. The portability of these devices, along with their
the ability to connect to networks in different locations makes it possible to
Mobile computing. Mobile computing is defined as the execution of
computing tasks, while the user is on the move, or in
different places from their usual work environment. This mobility introduces
a series of challenges for distributed systems, including the task of
handle different types of connectivity and potential disconnections
while the service is being used. Ubiquitous computing is the
utilization of the various computational devices available in
the user's environment, including their home, the office, and even environments
natural. The term 'ubiquitous' suggests that devices will be omnipresent.
in daily life, and they will not even be noticed by the user at the moment
which interact with them.
Figure 4: Ubiquitous computing scheme.
III.3) Distributed Multimedia Systems
Another need that has gained great relevance lately is the
the need to store and transmit audio and video to large amounts of
users through the network. The most important thing in this type of services is
preserve a continuous flow of data that ensures a good experience for
the user. The benefits of distributed multimedia systems can be
appreciate in applications that offer live television broadcasts or
pre-recorded, access to libraries of movies and music, video conferencing,
webcasting, VoIP, etc.
III.4) Distributed Systems as a
Commodity
Commodities are generic goods, meaning they are not specified.
differentiation between themselves. Normally, when talking about commodities, it is discussed
from raw materials or primary goods, such as wheat, which is planted in
anywhere in the world and that will have the same price and the same quality.
Nowadays there are companies that view distributed resources as a
commodity, where services are not purchased but leased by the
clients to specialized companies, which considers both physical resources
(storage capacity, processing nodes, etc.) as logical
(email, online calendars, data processing software, etc). Big
companies like Google, Amazon, Yahoo, Microsoft, etc. already offer this type
of services, under the concept of Cloud Computing.
The term Cloud Computing also promotes the vision of presenting a
offer everything under the label of a service, generating a pricing model
based on usage and not on the definitive purchase, which facilitates scalability of the
service and reduces the fixed costs for users. The Cloud is usually
implemented with computer clusters that ensure scalability of
service, which is presented to the user as a single integrated high service
performance.
IV) Challenges of Systems
Distributed
Here are some of the most important challenges related to the
distributed systems.
IV.1) Heterogeneity
Although the Internet is made up of different types of networks, its
differences are masked under a common protocol that is used by everyone
the computers connected to communicate. Likewise, the systems
operating systems of all computers connected to the Internet must implement
Internet protocols, however, not all have the same interface for
these protocols. For example, the message exchange used by the
UNIX machines are different from those used by Windows machines. On the other hand,
the different programming languages use different representations
for characters and data structures such as arrays and records. These
differences must be handled if programs written in
different languages communicate. In addition, the programs written by the
different developers cannot establish communication unless
they use a standard, for example, for communication over the network,
representation of basic data and data structures in messages.
IV.1.1) Middleware
The term Middleware refers to a layer of software that provides a
programming abstraction that masks the heterogeneity associated with the
networks, hardware, operating systems, and programming languages. CORBA
(Common Object Request Broker) is a classic example of middleware.
Other middleware such as RMI (Java Remote Method Invocation) only supports one
programming language. Most middleware is implemented
about Internet protocols, which also mask the differences
mentioned earlier. In addition to solving the problems of
heterogeneity, a middleware provides programmers a model
computational for the creation of distributed servers and applications. These
models can include remote invocation of objects, remote notification
of events, remote access to SQL and distributed processing of
transactions. For example, CORBA provides remote invocation of objects, the
which allows an object of a program running on a computer
invoke a method of another object in a program that runs in another
computer. Its implementation hides the fact that messages of
The request and the response are passed through a network.
IV.1.2) Heterogeneity and Mobile Code
The term mobile code is used to refer to the code that can be
transferred from one computer to another and is executed on the destination computer.
Java applets with a good example. This is not trivial, as it usually
The code that runs on one machine cannot always be executed on another
computer, due to the differences between the instruction sets of each
processor and the possible differences of the operating systems. The
Virtual Machine technologies offer a mechanism for
that the same code can be executed on different computers: the
compiler of a particular language generates code for a virtual machine
instead of doing it for a specific hardware. For example, the compiler of
Java generates code for the Java virtual machine, only requiring that
the virtual machine is implemented only once in the different systems
operational so that the programmed code can run on all of them. Today in
day, the most widely used mobile code is programs written in Javascript that
they are executed on web pages and loaded in different browsers
platforms.
IV.2) Opening
The openness of a computational system is related to its capacity.
to be extended and re-implemented. In the case of the systems
distributed, the opening is mainly associated with the degree to which a new
service can be added and made available to the different programs
clients. This opening cannot be achieved unless the specifications and
the documentation of the most relevant software interfaces of the system is
available for software developers. In other words, it corresponds
to the publication of the most relevant interfaces. However, the publication of
interfaces is just the starting point in the incorporation and extension of the
services in a distributed system. The challenge for designers is the management
of the complexity of the system associated with many developed components
by different people. To face this difficulty, the designers of the
Internet protocols introduced a series of documents called
Request for Comments or RFCs, each of them is identified by a
number. The specifications of Internet communication protocols
were published in these series in the early 80s, followed by the
specifications for the applications that ran on it, such as the
file transfer, email, and telnet in the mid-80s. This practice has
continued and forms the basis of the technical documentation of the Internet. This
series of documents also includes discussions in addition to the specification
of protocols([Link]). Distributed systems designed for
sharing resources are referred to as open distributed systems
distributed systems), to highlight the fact that they are extensible. These
can be extended at the hardware level, through the aggregation of new
computers to the network, and at the software level, incorporating new services or
re-implementing old services, thus enabling resource exchange
among different applications. Another benefit commonly mentioned of the
open systems have their independence from particular vendors.
IV.3) Security
Many of the information resources that are managed in a system
distributed has a high intrinsic value for its users. For this reason its
security is of vital importance. Information security consists of three
components: confidentiality (protection against disclosure to unauthorized persons
authorized), integrity (protection against alteration or corruption), and
availability (protection against interference to access resources).
En un sistema distribuido, los clientes envían solicitudes para acceder a los
data managed by servers, which translates into messages sent to
través de la red. Por ejemplo:
when a doctor requests access to a patient's data in a
hospital, or when you want to add information to the file of the
patient.
e-commerce and online banking services, customers send
your credit card number via the internet.
In both examples, the challenge translates into sending sensitive information to
through the network in a secure manner. However, security is not just about
not only involves hiding the content of the messages sent, but also involves
to know with certainty the identity of both the user sending and the one
receive the message. Both challenges can be faced using techniques
data decryption developed especially for these purposes,
which are widely used on the Internet today. Despite the great development
of security techniques, there are still some problems that have not been able
to be solved efficiently:
Massive attacks on services, consisting of the bombardment of services
with a large number of requests, it is also commonly referred to
as Denial of Service Attack.
Mobile code security. This type of code must be handled with
Be careful, as many times the origin of the executable code
it may not be reliable, possibly being the effects of the execution of the
unpredictable programs. It is for these that normally the execution of
Java Applets are blocked by default in most of the
internet browsers.
IV.4) Scalability
Distributed systems must operate effectively and efficiently
regardless of the scale of the system. Starting from small intranets
up to the very Internet. A system is called scalable when its efficiency
is not affected by a significant increase in its resources or number of
users. A clear example of this is the Internet.
IV.5) Error Handling
Computer systems sometimes fail both in hardware and in
software. When these failures occur, programs may produce
incorrect results or they may stop before finishing their task.
Failures in distributed systems are typically partial, meaning some
components may fail while others continue to function, making the
management of particularly complex errors. Some of the techniques for the
error handling are:
Fault Detection: Some faults can be detected, for example
Checksums are typically used to detect corrupted data.
in a message or in a file. However, it is difficult or sometimes
impossible to detect some failures, such as detecting the fall of a
server somewhere on the Internet.
Failure Masking: Some of the failures that have been
detected can be concealed or transformed into less severe.
Examples of failure concealment: (1) The messages can be
transmitted when their sending fails. (2) The data in files can be
duplicates on several disks, so that if one disk fails, the others
they still contain the correct information. Just the fact of discarding a
corrupt message and resend it, it is a simple mechanism to do
a less severe fault.
Fault Tolerance: Most services offered on the Internet
they present failures at some point. For this reason, the client programs of
these services can be designed to tolerate these failures. By
for example, when the browser cannot connect to a server, it
it does not last forever trying to establish communication,
but informs the user about the problem.
Failure Recovery: Recovery in this case relates to
the recovery of information after the system fails, given
that normally the calculations executed by some programs remain
incomplete, causing corrupt or incomplete data in the system.
Redundancy: Fault-tolerant systems can also be designed.
through the use of redundant components. Some examples of
these are: (1) There should be at least two distinct paths between two
routers on the Internet. (2) In the domain name systems
(DNS), each name table is replicated on at least two servers.
A database should be replicated on different servers to
ensure that the information is accessible after a failure in some
of these servers. These servers can also be designed to
detect failures in some of their peers; when a failure is detected in
a server, clients are automatically redirected to the
remaining servers.
Distributed Systems provide a high degree of availability in the case of
hardware failures. System availability is measured in terms of the
proportion of the time that is available for use. If a component of
distributed system failure only affects the task that was using it
component. A user in this case could switch to another computer, if
the computer that is used fails, or likewise, the processes of a server
they could be executed on another if the first one fails.
IV.6) Concurrence
Both services and applications offer resources that can be shared.
for clients in a distributed system. Therefore, there is the possibility that
several users attempt to access a resource simultaneously
specific shared. For example, a data structure that stores the
offers in an auction, is accessed very frequently when
about the term of the auction. The process that manages this resource
shared could only attend to one query at a time, but this modality
this would greatly limit the functioning of the system. For this reason, the services and
applications generally accept multiple queries from users
they are processed simultaneously. Specifically, suppose that each resource
is encapsulated as an object and each request is executed in threads
(concurrent) threads. In this case, it is possible that several threads may be
executed concurrently on an object, which could generate conflicts
and inconsistent results. For example, if two concurrent offers in the
The bids are 'Juan:$122' and 'Pedro:$111', and the corresponding operations are
interleaved without any kind of control, these should be stored as
Pedro: $111 and Juan: $122. The moral of this example is that each object that
represents a shared resource in a distributed system must ensure that
It can operate correctly in a concurrency context. This does not apply.
Not only to servers, but also to objects existing in applications. For
this, any programmer who wants to integrate code from an object that was not
designed for use in a distributed environment, it must make the changes
necessary to ensure safe operation in situations of
concurrency. In general, to make an object safe in an environment
concurrent ones, their operations must be synchronized in such a way that
stay consistent, which can be achieved with some techniques
standard like traffic lights, which are used in most systems
operational.
IV.7) Transparency
Transparency consists of hiding, both from the user and the programmer,
applications, the separation existing between the system components
distributed, so that the system is perceived as a "whole" and not
as independent components. The ANSA reference manual and the
International Organization for Standardization’s Reference Model for Open
Distributed Processing (RM-ODP) identifies eight forms of transparency:
1. Access transparency: the operations to access resources in
local or remote forms must be identical.
2. Location transparency: allows access to resources without
need to know their physical position or within the network.
3. Transparency of concurrency: allows various processes to operate
with shared resources in a concurrent manner without interfering.
[Link] transparency: allows multiple instances of a
resource can be used to improve availability and performance without
that the user or the programmer notices it.
5. Failure transparency: enables the concealment of failures, allowing for
users and applications completing their tasks regardless of
hardware or software component failures.
6. Mobility transparency: allows the movement of resources and
clients within the system, without affecting the operation of users or
programs.
7. Performance transparency: allows reconfiguring the system to
improve their performance as the load varies.
8. Scalability Transparency: allows for the system's expansion and its
applications without changing their structure or their algorithms.
The two most important forms of transparency are transparency in the
access and localization, given that their presence or absence affects
considerably to the utilization of distributed resources. Both are
usually referred to as network transparency. An example of
transparency in access is when a graphical interface of a folder
file system, in which it does not matter if the folder is located in the
local computer or on a remote computer. The same example in a system
without access transparency, it would be the same file system that does not allow
access to the remote folder unless you use some FTP application or
similar. An example of network transparency is email addresses
electronic. The email address:[Link]@[Link] finished it in
username and domain name. Sending an email to the user
does not imply knowing its physical location or the network. Nor the procedure to
Sending the email message will depend on the recipient's location.
Therefore, email within the Internet offers transparency
both in location and access, that is, network transparency. In the case of the
transparency of failures, this can also be illustrated in the context of
email, which is delivered even when the servers or links are
communication fails. The failures are masked when trying to retransmit the
messages until they are delivered correctly, which can even take
several days. In general, middleware converts network and process failures
in exceptions at the programming level. To illustrate transparency
mobility, let’s consider the case of mobile phones. Let’s assume that the
the caller and the person receiving the call are traveling by train in different
places of a country, moving from one cellular environment to another. We consider that
the caller's phone number as the client and the recipient's phone number
from the call as a resource. The two phone users who make the
calls are not aware of the mobility of phones between the different
coverage areas.
IV.8) Service Quality
Once users obtain the required functionality from a service,
for example, the file service of a distributed system, we can
continue to advance and ask ourselves about the quality of such service. The
main non-functional properties of systems that affect quality
of the service experienced by clients and users
son: reliability, safety, and performance. The ability to adapt to the
changes in system configuration and resource availability, also
has been recognized as an important aspect of service quality. The
Reliability and safety are critical points in the design of most
computer systems. Performance has been initially defined as
response speed terms, but lately it has been redefined in
terms of ensuring the timely execution of tasks. Some
applications, including multimedia, work with data streams that
they need to be processed or transferred from one process to another within a
limited time. The success of this task will depend on the availability of the
communication networks and the processing capacity of the system. The
networks today offer these levels of performance, but when they are
overload, the performance of these can be affected generating a
deterioration in the quality of the services that use it. For this reason, each resource
critical should be reserved for applications that need to guarantee a
certain quality of service functions properly.