Hierarchical Probabilistic Multicast

Patrick Th. Eugster, Rachid Guerraoui · Infoscience (Ecole Polytechnique Fédérale de Lausanne) · 2001

In our DACE project [6], diverging requirements expressed through QoS are mainly explored by a variety of different delivery semantics implemented through different dissemination algorithms ranging from "classic" Reliable Broadcast [18], to new and original algorithms, like the broadcast algorithm we introduce in [7], and which ensures reliable delivery of events despite network failures.While striving for strong scalability, we have invested considerable effort in exploring probabilistic (gossip-based) algorithms.These appear to be more adequate in the field of large scale event dissemination than traditional strongly reliable approaches like [18].Basically, probabilistic algorithms trade the strong reliability guarantees against very good scalability properties, yet still achieve a "pretty good degree of reliability" [12].Until now, most work on gossip-based algorithms considers broadcasting information to all participants in a system, paying little or no attention to individual and dynamic requirements, as typically encountered in contentbased dissemination.We present here Hierarchical Probabilistic Multicast (hpmcast [9]), a novel gossip-based algorithm which deals with the more complex case of multicasting an event to a subset of the system only.Requirements, such as limiting the consumption of local memory resources by view and message buffering, as well as exploiting locality (the proximity of participants) and redundancy (commonalities in interests of these participants), are all addressed.Though hpmcast has been motivated by our specific context of TPS, it is general enough to be applied to any context in which a strongly scalable primitive for event, message, or information dissemination is required.* This work is partially supported by Agilent Laboratories and Lombard Odier & Co.these acknowledgements converge.1 Moreover, such protocols hide any form of membership [2, 24], making them difficultly exploitable with more dynamic dissemination (filtering). Probabilistic AlgorithmsGossip, or rumor mongering algorithms [5], are a class of epidemiologic algorithms, which have been introduced as an alternative to such reliable network-level broadcast protocols.They have first been developed for replicated database consistency management, and have been mainly motivated by the desire of trading the strong reliability guarantees offered by costly deterministic algorithms against weaker reliability guarantees, but in return obtaining very good scalability properties.The analysis of such algorithms is usually based on stochastics similar to the theory of epidemics [3], where the execution is broken down into steps.Probabilities are associated to these steps, and such algorithms are therefore sometimes also referred to as probabilistic algorithms. Reliability DegreeThe "degree of reliability" is typically expressed by a probability; like the probability 1-α of reaching all participants in the system for any given message, or by a probability 1-β of reaching any given participant with any given message.Ideally, α and β are precisely quantifiable.A more precise measure, called ∆-Reliability, based on the distribution of the probability of reaching a fraction of participants, is given in [12]. Basic ConceptsDecentralization is the key concept underlying the scalability properties of gossip-based broadcast algorithms, i.e., the overall load of (re)transmissions is reduced by decentralizing the effort.Participants are viewed as peers, symmetric in role, which are all equally eligible to forward information.2 This makes of gossip-based algorithms ideal candidates for systems with an underlying peer-to-peer model. ParametersMore precisely, retransmissions are initiated in most gossip-based algorithms by having every participant periodically, i.e., every P ms (step interval ), send information to a randomly chosen subset of participants inside the system (gossip subset ).The size F of the subset is usually fixed, and is commonly called fanout.Gossip algorithms differ in the number of times the same information is gossiped.Every participant might gossip the same information the same number of times, meaning that the number of repetitions is fixed.Alternatively, the same information might be forwarded only once by a same participant, and the longest causal chain of message forwards can be limited by fixing the number of hops H (or forwards).Also, the number of rounds T (step intervals) that a message remains in the system can be limited. ApproachesGossiping techniques have been proposed in a broad spectrum of contexts.Consequently, these algorithms vary a greate deal of further characteristics. Messages.Gossip-based algorithms differ in the kind of information that is shipped by gossiped messages (gossips).In early gossip algorithms, gossips reflect the sender's message buffer, including information about missing messages.Gossips have also been used to directly propagate the multicast payload (e.g., [10]), like events in the case of TPS.1 Similarly, the scalability offered by other reliable network-level protocols, like Reliable Multicast Protocol (RMP) [35], Log-Based Receiver-Reliable Multicast (LBRM) [19], or Scalable Reliable Multicast (SRM) [13] is not sufficient for many current application scenarios.2 Note that the SRM protocol also relies on a peer-based approach.A retransmitted message is however rebroadcast to the entire system.

Read the paper · More papers on PaperTik