HomeNotesDevelopment Blog

Logical Clock in a Distributed System

Overview

What is a distributed system?

A key challenge of a distributed system is the delay between messages. If message delay is not negligible compare to the time between events in a single thread, the system should be considered distributed.

In this sense, the common example of two threads performing increment on a common memory location should be considered distributed. Because the event of increment happens so quickly in comparison to message passing between threads.

Main Goal

Give a algorithmic way to discover total ordering of events in a distributed system that is consistent with the user's perception. Hopefully with an acceptable bound.

The Partial Ordering of Events

We have the following definition:

A system consist of a collection of processes
\(\{ p_1, p_2, \dots , p_n\}\)
A process \(p_i\) consist of a sequence of events
\(\langle e_{i1}, e_{i2} , \dots , e_{im_i}\rangle\)
A class of distinguished events for sending and receiving a message
\( s, r \)

"Happened before"

We define the relation "happened before", denoted "→" as follow:

  1. Within a process, events that comes before also happened before \[ j < k \Rightarrow e_{ij} \to e_{ik} \]
  2. For a fixed message, if \(s\) sends the message and \(r\) receives the message, then \( s \to t \).
  3. "→" is transitive
  4. "→" is irreflexive

By this definition, we define \(e_1, e_2\) as concurrent if \((\neg a \to b) \land (\neg b \to a)\).

Logical Clock

We have the following definition:

  • Each process gets a clock, which is a mapping from the events of process to a number: \[ C_i(e) \]
  • There exists a system clock, which for each events in the system, agree with the clock of the process the events belong in: \[ C(e_{ij}) = C_i(e_{ij}) \]

The paper only assigns number to the collection of sequences of events that actually happened.

Clock Condition

There are many assignments of events to numbers, which one is correct?

Author gives the Clock Condition:

For all events \(e, e'\), \( e \to e' \Rightarrow C(e) < C(e') \)

Above is the endgoal, we can obtain the endgoal by satisfying two subconditions:

Clock Condition (C1)

Within a process, events that come before should be assigned a smaller number: \[ j < k \Rightarrow C_i(e_{ij}) < C_i(e_{ik}) \]

Clock Condition (C2)

For a fixed message, if \( s \) sends the message from process \(i\) to process \(j\) and \( r \) receives the same message, then \( s \) is assigned a smaller number: \[ C_i(s) < C_j(r) \]

A term:

On page 3 (or 560 of the journal), left column, the final paragraph, the author suppose we mark each timeline with natural numbers, with the number getting bigger as it goes up. Then the term "like-numbered" means drawing a tick line between the marked positions with the same natural number.

A way to satisfy the Clock Condition

We wish to satisfy C1, C2. More importantly, this needs to be done distributively, because there is no master clock we can refer to. So in practice, each process keeps their own clock, and we use an algorithm to sync the clock up.

The algorithm consist of two parts:

Algorithm: IR1

Each process \(p_i\) increment \(C_i\) between any two successive events. In practice, can simply increment after an event.

Algorithm: IR2

For a fixed message \(m\), send message event \(s\) done by process \(p_i\) contains a time stamp of the of the time of the send event, \(T_s = C_i(s) \). On receiving the message, the receipient process \( p_j \) increment its clock \( C_j \) to a number greater than or equal to its current value or \( T_s \).

WARNING: There is a question to be asked here, should there still be an increment after choosing the max of current value and \( T_s \)? I think you should still increment as instructed by IR1, because receiving constitute an event.