Skip to content

Latest commit

 

History

1 Commit

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 

Repository files navigation

Summary

In an on-site demonstration in Atlanta, Ditto continuously synchronized data across more than 100 radios with sub-second latency while using only 35% of the available channel capacity. This represented a more than fourfold increase over the scale supported by previous attempts, and this write-up is a recollection of the engineering effort behind the improvements.

The writing below has the following structure:

  1. describing the peer-to-peer synchronization system at Ditto,
  2. investigating the hardware characteristics of the targeted workload,
  3. arguing the bad fit of pairwise sync,
  4. leveraging a multicast-oriented protocol for group data transport,
  5. designing a multicast-native sync protocol on top of it, and
  6. describing comprehensive testing crucial to the project's success.

Naturally, many details are omitted for the sake of brevity. Among them: the true complexity of Ditto's peer-to-peer sync system (that even allows for dynamic subscriptions); the internal workings of the NORM protocol (including relevant congestion control mechanisms); and implementation details and optimizations of the resulting multicast-native sync protocol.

The scalability problem

The peer-to-peer replication system at Ditto focuses on reliability. Precise state-tracking allows all data to be synchronized, assuming peers remain connected for long enough. This pairwise communication mode, however, comes with disadvantages regarding scalability. A peer can only connect to a limited number of devices, and thus an update might require multiple "hops" to synchronize throughout the entire device mesh. That fact, combined with the frequent small "control packets" exchanged between peers for bookkeeping, means that network usage grows very significantly with increased mesh sizes.

A critical workload required efficient data synchronization for a large fleet of specialized radio devices. In this scenario, tests with the standard peer-to-peer synchronization protocol supported meshes of at most 25 devices. Higher scalability would require major improvements, and "the multicast project" was the name given to the efforts behind these improvements: investigation, design, implementation, testing, and demonstration. But that's a bit of a leap: why "multicast", exactly?

Understanding the previous limit and the driver for multicast requires examining the characteristics of such special radio hardware.

Modern radios use many special technologies, but most relevantly here, these radios specifically use an aspect of "time-division multiple access (TDMA)", where time is divided into slots, and each slot is designated for a single message to be transmitted. It's like radios "take turns" broadcasting their messages, and that's all internally scheduled by the radio firmware. During each time slot, a message is transmitted to all radios in the group at once. The radios are, by nature, "broadcast-y". Additionally, the slots are of a fixed size: each transmission opportunity can carry a message of up to a certain number of bytes, and any unused capacity is simply wasted.

The above clarifies why peer-to-peer synchronization didn't scale well. A peer-to-peer message is directed at a single other peer — containing diff-like conflict-free replicated data types (CRDTs) based on the receiving peer's believed state. Other peers receiving that same message can't process it. In a broadcast-y radio mesh, a peer-to-peer message could reach dozens of devices, only to be discarded by all but one of them. Even worse, while that single message is in-flight (and thus using a time slot), no other message can be exchanged — slowing throughput very significantly. The frequent control packets (for bookkeeping) add another inefficiency to the process, as they consume a whole time slot despite their small payloads.

Happily, the same analysis also point the way forward: leveraging the specifics of radio communication to unlock higher scale.

Leveraging multicast

The selected transport protocol was NORM (NACK-Oriented Reliable Multicast), along with its C++ implementation by the U.S. Naval Research Laboratory. In particular, negative-acknowledgments (NACKs) could prevent ineffective bandwidth usage from frequent control packets. Messages are assumed to be transmitted correctly, and in that case, no extra acknowledgement is needed (no small packet consuming a slot). In less-happy cases, a message is "repaired", either through preemptive "error correction codes", or through explicit repair messages handled by the protocol. NORM is a particularly good match for the modern radios we're talking about, as the radios themselves add yet another reliability layer, through automatic retransmissions. This means that only a small percentage of messages will fail to be delivered, and repairing is rarely necessary.

NORM provides generally reliable group transport communication, but by itself it is not enough to reliably sync data. For example, peers can still miss messages: suppose a peer has been disconnected long enough such that it misses all repair attempts. Or, simpler, suppose a new device joins the mesh. That device would never hear of older messages that have already been exchanged, and some sort of "backfilling" mechanism is needed.

In the critical workload, operators make decisions based on current data, and the group is extremely sensitive to new-data sync latency. Thus, the multicast-native protocol targeted two key objectives:

  1. "new" data must be favored over "historical" data, and synced promptly;
  2. "backfilling" must not significantly slow down the mesh; Where "new data" means updates that have been recently created, and "backfilling" means the process through which peers "catch-up" on older data that hasn't been exchanged in a while.

Implementation-wise, these objectives led to a type of "quality of service" for syncing data, where peers immediately start receiving new updates, while old data re-synchronization is specially triggered on occasion, and batched. Configurable parameters governed many aspects of the system, such as NORM behavior, backfilling circumstances, and miscellaneous experimental optimizations. This flexibility allowed for safely testing and tuning in preparation for the demonstration.

Testing and results

Crucially important, significant effort was devoted to comprehensive testing. We developed a testing harness that could simulate:

  • multiple devices running in parallel;
  • different network topologies;
  • variable network conditions like packet loss, link latency, and bandwidth constraints; and
  • a flexible range of workloads, including extreme scenarios. The system also measured metrics like bandwidth usage, sync rates and latency, and convergence times, as well as retained logs and artifacts for further analysis.

Even without access to the specific radio hardware (let alone in the targeted scale), this extensive simulation testing allowed for very rapid iteration, trade-off exploration, tuning opportunities, and building confidence in the solution throughout development.

As mentioned before, the effort culminated in a successful on-site demonstration with greatly increased performance. Metrics suggest that even higher scale would still be supported, but unfortunately the testing was limited by the number of devices available.

About

A recollection of the engineering effort behind 'the multicast project'

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors