1 of 69

DISTRIBUTED COMPUTING

Sunita Mahajan, Principal, Institute of Computer Science, MET League of Colleges, Mumbai

Seema Shah, Principal, Vidyalankar Institute of Technology, Mumbai University

© Oxford University Press 2011

2 of 69

Chapter - 6�Distributed System Management

© Oxford University Press 2011

3 of 69

Topics

  • Introduction
  • Resource management
  • Task assignment approach
  • Load balancing approach
  • Load sharing approach
  • Process management in a distributed environment
  • Process migration
  • Threads
  • Fault tolerance

© Oxford University Press 2011

4 of 69

Introduction

© Oxford University Press 2011

5 of 69

Categories of Distributed System management

  • Resource management
  • Process management
  • Fault tolerance

© Oxford University Press 2011

6 of 69

Resource Management

© Oxford University Press 2011

7 of 69

Process scheduling techniques

  • Task assignment approach
  • Load balancing approach
  • Load sharing approach

© Oxford University Press 2011

8 of 69

Example: Google system

  • Load balancing by using least loaded server
  • Proximity routing
  • Fault masking

© Oxford University Press 2011

9 of 69

Desirable features of a good global scheduling algorithm

  • No apriori knowledge about processes to be executed
  • Ability to make dynamic scheduling decisions
  • Flexible
  • Stable
  • Scalable
  • Unaffected by system failures

© Oxford University Press 2011

10 of 69

Task Assignment Approach

© Oxford University Press 2011

11 of 69

Task assignment

  • Minimize IPC costs
  • Less turnaround time for process completion
  • High degree of parallelism
  • Efficient usage of all system resources

© Oxford University Press 2011

12 of 69

Graph theoretic deterministic algorithm

A system with m CPUs and n processes has any of the following three cases:

  • m=n: Each process is allocated to one CPU
  • m<n: Some CPUs may remain idle or work on earlier allocated processes
  • m>n: There is a need to schedule processes on CPUs, and several processes may be assigned to each CPU.

© Oxford University Press 2011

13 of 69

Example of graph theoretic deterministic algorithm-1

  • Weighted graph
    • Each node is a process
    • Each arc is message flowing between two processes

© Oxford University Press 2011

14 of 69

Example of graph theoretic deterministic algorithm-2

© Oxford University Press 2011

15 of 69

Centralized heuristic algorithm

  • Also called Top down algorithm
  • Allocated processing capacity fairly

2

© Oxford University Press 2011

16 of 69

Hierarchical algorithm

  • Works between two levels in a group
  • Top of the tree is truncated into a committee which manages fault tolerance

© Oxford University Press 2011

17 of 69

Load Balancing Approach

© Oxford University Press 2011

18 of 69

Load balancing Taxonomy

  • Improve resource utilization

© Oxford University Press 2011

19 of 69

Issues in designing in load balancing algorithms

  • Deciding policies for:
    • Load estimation
    • Process transfer
    • Static information exchange
    • Location
    • Priority assignment
    • Migration limitation

© Oxford University Press 2011

20 of 69

Policies for Load estimation

  • Parameters:
    • Time dependent
    • Node dependent

© Oxford University Press 2011

21 of 69

Policies for Process transfer

  • Threshold policy
    • Static
    • Dynamic

© Oxford University Press 2011

22 of 69

Location policies

  • Used to select destination node

© Oxford University Press 2011

23 of 69

State information exchange

  • Dynamic policy
  • Decision based on state information

© Oxford University Press 2011

24 of 69

Priority assignment

  • To schedule local and remote processes at a node

© Oxford University Press 2011

25 of 69

Migration limiting policies

  • Uncontrolled policy
  • Controlled policy

© Oxford University Press 2011

26 of 69

�Load Sharing Approach �

© Oxford University Press 2011

27 of 69

Issues in designing load sharing algorithms

  • Load estimation policies
  • Process transfer policies
  • Location policies
  • State information exchange policies

© Oxford University Press 2011

28 of 69

Location policies-1

  • Decides whether sender or receiver node process is to be migrated

© Oxford University Press 2011

29 of 69

Location policies-2

  • Sender initiated algorithms make scheduling decisions at process arrival epoch
  • Receiver initiated algorithms make scheduling decisions at process departure epochs

© Oxford University Press 2011

30 of 69

State information exchange policies

  • Broadcast
  • Poll

© Oxford University Press 2011

31 of 69

Process Management In A Distributed Environment

© Oxford University Press 2011

32 of 69

Functions of distributed process management

  • Process migration
    • change of location and execution of a process from current processor to the destination processor

© Oxford University Press 2011

33 of 69

Desirable features of a good process migration mechanism

  • Transparency
  • Minimal interference
  • Minimal residual dependencies
  • Efficiency
  • Robustness
  • Ability to communicate between co processes of the job

© Oxford University Press 2011

34 of 69

Process Migration

© Oxford University Press 2011

35 of 69

Steps involved in process migration

  • Freezing process on the source node
  • Starting process on the destination node
  • Transporting process address space on destination node
  • Forward the messages addressed to migrated processes

© Oxford University Press 2011

36 of 69

Mechanism

© Oxford University Press 2011

37 of 69

Freezing process on source node

  • Blocking sequence:
    • Blocking the process immediately
    • Wait for I/O operations to complete and then block the process.
  • Track information about open files
  • Create an empty process on the destination node
  • Transfer the migrant process and address space
  • Restart process on destination node

© Oxford University Press 2011

38 of 69

Address space transport mechanisms-1

  • Process address space:
    • Process state: PCB information
    • Process address space: Program code, data and stack

© Oxford University Press 2011

39 of 69

Address space transport mechanisms-2

  • Total freezing:
    • Process execution stopped during address space transfer

© Oxford University Press 2011

40 of 69

Address space transport mechanisms-3

  • Pre transfer:
    • Address space is transferred while process continues to run on source node
    • Highest priority in scheduling

© Oxford University Press 2011

41 of 69

Address space transport mechanisms-4

  • Transfer-on –reference:
    • Process state is transferred while address space is transferred on demand

© Oxford University Press 2011

42 of 69

Message forwarding

  • Track and forward messages which have arrived on source node after process migration

© Oxford University Press 2011

43 of 69

Handle communication between cooperating processes

    • Avoid separation of coprocesses
  • Home node concept
    • Deployed in Sprite system

© Oxford University Press 2011

44 of 69

Process migration in heterogeneous systems

    • Handling floating point numbers
  • Different sized exponents in XDR format
  • Handling overflow and underflow
  • Handling Mantissa
  • Handling signed infinity and zero representations

© Oxford University Press 2011

45 of 69

Advantages of process migration

  • Reduce average response time of heavily loaded nodes
  • Speed up of individual jobs
  • Better utilization of resources
  • Improve reliability of critical processes
  • Improving system security

© Oxford University Press 2011

46 of 69

Threads

© Oxford University Press 2011

47 of 69

Process v/s threads

  • Analogy:
    • Thread is to a process as process is to a machine

© Oxford University Press 2011

48 of 69

Comparison

© Oxford University Press 2011

49 of 69

Thread models

  • Dispatcher worker model
  • Team model
  • Pipeline model

© Oxford University Press 2011

50 of 69

Thread: Dispatcher worker model

© Oxford University Press 2011

51 of 69

Thread: Team model

© Oxford University Press 2011

52 of 69

Thread: Pipeline model

© Oxford University Press 2011

53 of 69

Design issues in threads

  • Thread semantics
    • Thread creation, termination
  • Thread synchronization
  • Thread scheduling

© Oxford University Press 2011

54 of 69

Thread synchronization

  • Execution in Critical region
    • Use binary semaphore

© Oxford University Press 2011

55 of 69

Threads scheduling

  • Priority assignment facility
  • Choice of dynamic variation of quantum size
  • Handoff scheduling scheme
  • Affinity scheduling scheme
  • Signals used for providing interrupts and exceptions

© Oxford University Press 2011

56 of 69

Implementing thread package

  • User level approach

Kernel level approach

© Oxford University Press 2011

57 of 69

Comparison of thread implementation-1

© Oxford University Press 2011

58 of 69

Comparison of thread implementation-2

© Oxford University Press 2011

59 of 69

Threads and Remote execution

  • RPC
  • RMI and Java threads

© Oxford University Press 2011

60 of 69

RPC execution

© Oxford University Press 2011

61 of 69

Threads are created on the fly

© Oxford University Press 2011

62 of 69

�Fault Tolerance �

© Oxford University Press 2011

63 of 69

Component faults

  • Transient faults
  • Intermittent faults
  • Permanent faults

  • Mean time to failure = ∑ kp (1-p) k-1

k=1

  • Mean time to failure = 1/p

© Oxford University Press 2011

64 of 69

System failures

  • Fail silent faults / fail stop faults
  • Byzantine faults

© Oxford University Press 2011

65 of 69

Use of redundancy

  • Information redundancy
  • Time redundancy
  • Physical redundancy
    • Active replication
    • Primary backup methods

© Oxford University Press 2011

66 of 69

Active replication-1

  • State machine approach

(TMR -Triple Modular Redundancy)

© Oxford University Press 2011

67 of 69

Active replication-2

© Oxford University Press 2011

68 of 69

Primary backup

  • Uses two machines :
    • Primary and backup
  • Uses limited number of messages such that these messages go only to the primary server and no ordering is required

© Oxford University Press 2011

69 of 69

Summary

  • Introduction
  • Resource management
  • Task assignment approach
  • Load balancing approach
  • Load sharing approach
  • Process management in a distributed environment
  • Process migration
  • Threads
  • Fault tolerance

© Oxford University Press 2011