This one IP is worth 200 points. Parts 1 and 2 together are worth 100 points. Part 3, which is “summative” is worth 100 points.
Distributed mutual exclusion
Example scenario:
An order from a customer is sent to a group of distributed nodes for processing.
One node must grab the order and commit to handling it.
This means that each node that’s interested in this order must gain exclusive access to it to do the following:
find out if it has been taken on by another node
decide whether to take it on or not
mark it as taken, if necessary.
We shall say that there are locks:
All nodes in the group must agree on which one node – if any – holds a particular lock at each point in time
In the above scenario, the lock could be for a certain customer
12/2/2019
1
Mutual-exclusion algorithms
Centralized mutual exclusion:
A semaphore server handles the requests for the lock
Problem: Single point of failure
Contention-based mutual exclusion:
All nodes compete equally for the lock.
A queue of pending requests for the lock is kept at each node
Request-resolution criteria may be based on the time of each request
We shall look at a history of such algorithms, improving on each other
Controlled mutual exclusion:
A token representing the lock visits all nodes in a regulated fashion
10/26/2020
2
Contention based: Timestamp-based schemes
The nodes use messages to request and release the lock
The requests are totally ordered based on Lamport's clock
(The basic implementation of the clock)
Lamport’s distributed mutual exclusion algorithm (1978):
A node requests the lock:
by multicasting a time-stamped request to each node in the group
Each node maintains a time-ordered queue of pending requests
After receiving a request, a node sends a reply to the requester
10/26/2020
3
3
Getting and releasing the lock
A node multicasts its lock request
It holds the lock when and only when:
It has gathered all the reply messages and
Its own request is at the top of the node’s queue
To release the lock, the node multicasts a release message
On receiving a release message, each receiving node:
Removes the completed request from its queue.
If the node's own request is now at the head of the queue
and it has received all reply messages
it holds the lock
6/10/2020
4
Lamport mutual exclusion: Clock diagram S3 has the lock; S1 wants (and gets) it
12/2/2019
5
5
Ricart and Agrawala’s algorithm (1981)
In Lamport’s algorithm, with N nodes:
Number of messages per use of the lock: 3 (N – 1)
We can optimize this by letting reply messages also signal release:
Glenn Ricart and Ashok Agrawala’s algorithm:
If a node receives a request
while holding the lock, or
if its own request is timestamped before the new request
then the node delays its reply until it releases the lock
Release messages also serve as replies; no separate release msgs
otherwise, it replies immediately
It discards the request once it has replied to it
A node receiving a reply from a node X removes X’s request from the queue if it’s there
When a node has received all the reply messages, it holds the lock
Number of messages per use with this approach: 2 (N – 1)
10/26/2020
6
Ricart-Agrawala clock diagram S3 holds the lock; S1 wants it
9/21/2019
7
R & A protocol-state machine
It describes the local protocol state at any one node in the group
Multiple nodes can be in states Released and Wanted at any given point in time.
At most one node can be in state Held (regarding this lock).
(Better state names may be “Not interested,” “Wanting,” and “Holding.”)
6/10/2020
8
Mamoru Maekawa’s (1985) quorum-based algorithm
With the timestamp-based algorithms, if one node is unavailable, the lock becomes unavailable.
Actually, a node only needs the votes from a subset of all nodes as long as the subsets of any two nodes overlap.
Each node pi has a voting set Vi, of nodes (a.k.a., quorum set)
The voting sets are ideally of equal size, K
Vi includes pi itself.
Any two voting sets have at least one member in common.
Each node only needs the votes of all the nodes in its voting set
Including its own vote
Each node votes in M voting sets
Optimally, K N (N = the number of nodes) and M = K
10/26/2020
9
Voting-set diagram with 7 nodes (Dr. Brandon “Oz” Osborn, DCS, summer 2015)
6/10/2019
10
Explanation of diagram with 7 nodes
Each node points at the two other nodes in its voting set
Color coded
It needs votes from those two and from itself: K = 3
Node 1 needs votes from nodes 1, 2 and 3
Each node is pointed to by two other nodes
It votes in those voting sets and in its own: M = 3
Node 1 is in 4’s and 6’s voting sets
Node 2 is in 1’s and 5’s voting sets
Node 3 is in 1’s and 7’s voting sets
Say that Node 1 wants the lock.
It gets its own vote if it hasn’t voted for itself (1) nor for 4 or 6
It gets 2’s vote if 2 hasn’t voted for itself (2) or for 5
It gets 3’s vote if 3 hasn’t voted for itself (3) of for 7
10/26/2020
11
Maekawa’s algorithm as seen by each node (here called “this” node, i)
Initially, State = Released, Voted = False
This node (i) wants the lock. // The calling thread blocks until State=Released.
State = Wanted
Multicast a request to all nodes in voting set Vi (including self)
Wait for K replies // Including from itself
State = Held // after K replies have been received
This node receives a request: // Possibly from itself
If Voted then
Queue the request without replying
else
reply
Voted = True
This node is done with the lock:
State = Released
Multicast release to all nodes in voting set Vi (including self)
This node receives a release message // Possibly from itself
Voted = False
If the request queue is not empty then
remove the head of the queue, say pk // Possibly itself
reply to pk
Voted = True
6/11/2020
12
Maekawa scenario from the point of view of P1
3/11/2019
13
Controlled approaches: Token based
Alternative to contention-based mutual exclusion with less message overhead.
Different topologies:
Ring
Tree (not covered here)
Broadcast structures (not covered here)
The token (lock) is like the talking stick used by tribes to show who has the floor
10/26/2020
14
Ring structure
The token circulates in a logical ring of nodes.
The node that has the token holds the lock.
When a node is done (or does not need the lock at this time), it passes the on token to the next node.
The wait for the token can be long even if only one node needs the lock.
The performance is better if the request load is high.
The token can be used to carry state information.
Token Bus (IEEE 802.4) and Token Ring (IEEE 802.5) standard protocols for LANs use hardware structures of processors.
6/10/2020
15
Unit 5 DB 1
A lot of literature on group communication is available.
Find an article about any kind of group communication that seems interesting
Post a little summary of it.
6/10/2020
16
Final IP, 200 points
Please refer to:
“CS844 DB and IP Instructions”
6/10/2020
17
Part 1: Style of the answer expected (Slide 11 “2-phase total-order multicast scenario” from group communication chat)
18
Comm. handler buffer at node G1:
Msg Suggested time Commit time
m0 2 Delivered
m1 7 15
m2 9 13
m3 18 Pending
10/26/2020
You don’t have to follow the color-coding scheme. Make all arrows black if you like.
Part 2: Style of the answer expected: (Slide 9 from this presentation)
State machine for Maekawa’s algorithm like the one we had for Ricart and Agrawala:
Note that Maekawa’s algorithm has other states, as for example these:
Released and not voted
Released and voted
Wanted and not voted
Wanted and voted
Held and voted (for self)
Refer to slide “Maekawa scenario” (slide 14)
Warning: There is an incorrect solution for this problem on the Web. Don’t use it!
10/26/2020
19
Part 3: Essay
Summarize the course in your own words
4-7 or so pages.
Make it read nicely from top to bottom
Style according to APA
Reuse your own DB posts, IPs, etc., if you like
Please correct known errors.
Include (at least):
How threads and safe objects interact
Entity-life modeling
Deadlock prevention
Logical clocks
6/10/2019
20
Ricart & Agrawala: Protocol monitor
package Ricart (The application thread sees a semaphore safe object)
protected Lock_PO State-machine PO for distributed mutex
entry Acquire when state = Released Set state to Wanted
entry Confirmed when True Requests have been multicast. Save request time;
then requeue the calling thread on Replies_Collected
private entry Replies_Collected Block until replies are collected. Then set state to Held
procedure Reply_Received
procedure Release; Reply to each queued request. Set State to Released
procedure Request_Received (… Reply) Queue request or let caller reply.
end Lock_PO;
procedure Acquire_Lock Application thread at this node wants lock
Lock_PO.Acquire; State is Wanted upon return from Acquire
Total_Order_Sender.Send ( ) Multicast request
Lock_PO.Confirmed Block on Replies_Collected until lock acquired
procedure Release_Lock Patch through to Lock_PO.Release. (This node releases lock.)
procedure Reply_Received Patch through to Lock_PO. Reply_Received
(Called by some listener thread.)
procedure Request_Received Called by some listener thread when another node wants lock
Lock_PO.Request_Received (… Reply : out Boolean)
Here: If parameter Reply is true, send reply to node
end Ricart
10/26/2020
21
R&A protocol-monitor stack
3/16/2016
22
Maekawa protocol monitor
package Maekawa
protected Maekawa_PO State-machine PO for distributed mutex
entry Acquire when state = Released Set State to Wanted
procedure Reply_Received Increment reply count
entry Replies_Collected when …. Block until replies collected, then set state to Held
procedure Release Set State to Released
procedure Request_Received If Voted, queue request, else set Voted true and
let calling procedure know to reply
procedure Release_Received If request queue not empty, return head, etc.
procedure Acquire_Lock Application thread (this node) wants lock
Maekawa_PO.Acquire
Multicast request to all nodes in Vi (including self)
Maekawa.Replies_Collected
procedure Release_Lock Application thread (this node) releases lock
Maekawa_PO.Release
Multicast release to all nodes in Vi (including self)
procedure Reply_Received Patch through to Maekawa_PO
procedure Request_Received
Maekawa_PO.Request_Received Find out if reply to be sent now
Send reply if necessary
procedure Release_Received
Maekawa_PO.Release_Received Get first pending request (if any)
Send reply if necessary
end Maekawa;
6/10/2020
23
Conclusion
We looked at 3 algorithms for distributed mutual exclusion
They all build on multicast
Discussed in the previous chat
There is a kind of historical progression:
Contention-based approaches:
Lamport 1978
Optimized by Ricart and Agrawala 1981
Quorum-based
Maekawa 1985
10/26/2020
24
References
Coulouris, G., Dollimore, J., Kindberg, T., Blair, G. (2012). Distributed systems: Concepts and Design, 5th Ed., Boston, MA: Addison-Wesley. Section 15.2
Lamport, B. (1978). Time, clocks and the ordering of events in a distributed system, CACM, vol. 21, no. 7, pp 558-565.
Lamport, B. (2015). The computer science of concurrency: The early years, CACM, vol. 58, no. 62, pp 71-76.
Maekawa, Mamoru. (1985). A √N algorithm for mutual exclusion in decentralized systems, ACM Transactions on Computer Systems, vol.3, no. 2, pp 145-159.
Ricart, G. and Agrawala, A. K. (1981). An optimal algorithm for mutual exclusion in computer networks, CACM vol. 24, no. 1, pp. 9-17.
Thiare, O. & Fall, P. A. (2012). Using Maekawa’s algorithm to perform distributed mutual exclusion in quorums. Advances in Computing, 2(4), 54 – 59.
Wu, W., Zhang, J., Luo, A., & Cao, J. (2015). Distributed Mutual Exclusion Algorithms for Intersection Traffic Control. IEEE Trans. Parallel Distrib. Syst, 26(1), 65-74. doi:10.1109/tpds.2013.2297097
Yadav, N., Yadav, S., & Mandiratta, S. (2015). A review of various mutual exclusion algorithms in distributed systems. International Journal of Computer Applications, 129(14), 11 – 16.
3/11/2019
25
Web resources
https://www.cs.columbia.edu/~du/ds/assets/lectures/lecture8.pdf
http://cse.csusb.edu/tongyu/courses/cs660/notes/dmex.php
https://www2.cs.siu.edu/~rahimi/cs420/slides/cs420-part4.pdf
12/17/2018
26
Solution for last week: Group communication: Exercise 1
Messages delivered in this sequence:
Msg J delayed
Msg G delivered: P1’s vector: 0,1,0,0
Msg H delayed
Msg F delivered: P1’s vector: 0,1,1,0
Msg J now delivered: P1’s vector: 0,1,1,1
Msg H now delivered: P1’s vector: 0,1,2,1
9/6/2016
27
S1S2S3S4requestreplyapplication releases lock
S1S2S3S4requestreplyapplication releases lock
S1 S2 S3 S4 request reply application releases lock
Releasedrequest / replyHeldrequest / queueAcquire lock /multicast requestReturn from ReleaseAll replies receivedWantedrequest [earlier than own] / replyrequest [later than own] / queueRelease lock / reply to all queued request
Released request / reply
Held request / queue
Wanted request [earlier than own] / reply request [later than own] / queue
Acquire lock / multicast request
Return from Release
All replies received
Release lock / reply to all queued request
5432716
5 4 3 2 7 1 6
State of Node P1 Event Action New state Queue
Released, not voted Request from P4 Reply to P4 Released, voted P4 gets lock
Relased, voted Release from P4 Released, not voted P4 is done
Released, not voted Request from P2 Reply to P2 Released, voted P2 gets lock
Released, voted Request from P3 Released, voted P3 P3's request is queued
Released, voted Release from P2 Reply to P3 Released, voted P2 done; P3 gets lock
Released, voted Release from P3 Released, not voted P3 done
Released, not voted P1 calls Acquire Multicast request Wanted, not voted Local call to Acquire
Wanted, not voted Request from P1 Reply to P1 Wanted, voted P1 receives and replies to own request
Wanted, voted Replies from P1, P2, P3, P4 Held, voted P1 gets lock
Held, voted P1 calls Release Multicast release Released, voted Local call to Release
Released, voted Release from P1 Released, not voted P1 receives own release message
Released, not voted Request from P2 Reply to P2 Released, voted P2 gets lock
Released, voted Request from P4 Released, voted P4 P4's request is queued
Released, voted P1 calls Acquire Multicast request Wanted, voted P4 Local call to Acquire
Wanted, voted Request from P1 Wanted, voted P4 ,P1 P1 queues own request
Wanted, voted Release from P2 Reply to P4 Wanted, voted P1 P2 done. P4 gets lock
Wanted, voted Release from P4 Reply to P1 Wanted, voted P4 done. P1 replies to own request
Wanted, voted Replies from P1, P2, P3, P4 Held, voted P1 gets lock
Held, voted P1 calls Release Multicast release Released, voted Local call to Release
Released, voted Release from P1 Released, not voted P1 receives own release message
State of Node P1Event ActionNew stateQueue
Released, not votedRequest from P4Reply to P4Released, votedP4 gets lock
Relased, votedRelease from P4 Released, not votedP4 is done
Released, not votedRequest from P2Reply to P2Released, votedP2 gets lock
Released, votedRequest from P3 Released, votedP3P3's request is queued
Released, votedRelease from P2Reply to P3Released, votedP2 done; P3 gets lock
Released, votedRelease from P3 Released, not votedP3 done
Released, not votedP1 calls AcquireMulticast requestWanted, not votedLocal call to Acquire
Wanted, not votedRequest from P1Reply to P1Wanted, votedP1 receives and replies to own request
Wanted, votedReplies from P1, P2, P3, P4 Held, voted P1 gets lock
Held, votedP1 calls ReleaseMulticast releaseReleased, votedLocal call to Release
Released, votedRelease from P1 Released, not votedP1 receives own release message
Released, not votedRequest from P2Reply to P2Released, votedP2 gets lock
Released, votedRequest from P4 Released, votedP4P4's request is queued
Released, votedP1 calls AcquireMulticast requestWanted, votedP4Local call to Acquire
Wanted, votedRequest from P1 Wanted, votedP4 ,P1P1 queues own request
Wanted, votedRelease from P2Reply to P4Wanted, votedP1P2 done. P4 gets lock
Wanted, votedRelease from P4Reply to P1Wanted, votedP4 done. P1 replies to own request
Wanted, votedReplies from P1, P2, P3, P4 Held, voted P1 gets lock
Held, votedP1 calls ReleaseMulticast releaseReleased, votedLocal call to Release
Released, votedRelease from P1 Released, not votedP1 receives own release message
Sheet1
| State of Node P1 | Event | Action | New state | Queue | |
| Released, not voted | Request from P4 | Reply to P4 | Released, voted | P4 gets lock | |
| Relased, voted | Release from P4 | Released, not voted | P4 is done | ||
| Released, not voted | Request from P2 | Reply to P2 | Released, voted | P2 gets lock | |
| Released, voted | Request from P3 | Released, voted | P3 | P3's request is queued | |
| Released, voted | Release from P2 | Reply to P3 | Released, voted | P2 done; P3 gets lock | |
| Released, voted | Release from P3 | Released, not voted | P3 done | ||
| Released, not voted | P1 calls Acquire | Multicast request | Wanted, not voted | Local call to Acquire | |
| Wanted, not voted | Request from P1 | Reply to P1 | Wanted, voted | P1 receives and replies to own request | |
| Wanted, voted | Replies from P1, P2, P3, P4 | Held, voted | P1 gets lock | ||
| Held, voted | P1 calls Release | Multicast release | Released, voted | Local call to Release | |
| Released, voted | Release from P1 | Released, not voted | P1 receives own release message | ||
| Released, not voted | Request from P2 | Reply to P2 | Released, voted | P2 gets lock | |
| Released, voted | Request from P4 | Released, voted | P4 | P4's request is queued | |
| Released, voted | P1 calls Acquire | Multicast request | Wanted, voted | P4 | Local call to Acquire |
| Wanted, voted | Request from P1 | Wanted, voted | P4 ,P1 | P1 queues own request | |
| Wanted, voted | Release from P2 | Reply to P4 | Wanted, voted | P1 | P2 done. P4 gets lock |
| Wanted, voted | Release from P4 | Reply to P1 | Wanted, voted | P4 done. P1 replies to own request | |
| Wanted, voted | Replies from P1, P2, P3, P4 | Held, voted | P1 gets lock | ||
| Held, voted | P1 calls Release | Multicast release | Released, voted | Local call to Release | |
| Released, voted | Release from P1 | Released, not voted | P1 receives own release message |
S1G1G2S22175 67 8989 1011 12 15 16 1713 1410 17 18 1915 16m1 (7)m1(15)m2m2(9)(13)171819(13) 20(15)
Released request / reply
Held request / queue
Wanted request [earlier than own] / reply request [later than own] / queue
Acquire lock / multicast request
Return from Release
All replies received
Release lock / reply to all queued request
Ricart
protocol monitor
Application
protocol monitor
Total-order sender
protocol monitor
Sending
PO
Lock PO
Ricart.api
Acquire_LockRelease_Lock
Send
Ack_Received
Reply_ReceivedReq_Received
Ricart.api
Acquire_Lock
Release_Lock
Send
Ack_Received
Reply_Received
Application protocol monitor
Ricart protocol monitor
Total-order sender protocol monitor
Lock PO
Sending PO
Req_Received