Show Table of Contents

Queue Replication and High Availability

Chronicle Queue Enterprise has been integrated with Chronicle Services to support state management, messaging between services, and message persistence. In addition, for replication of queues on remote hosts using TCP/IP, Chronicle Services framework uses Chronicle Queue Enterprise. This document first explains why queue replication is used, and describes the key features of queue replication as well as how this applies to Chronicle Services.

Background

A key requirement for many enterprise class software applications is the ability to prevent data loss. The financial sector, as an example, may be highly sensitive to such issues since they may affect the availability of auditing data as well as calculation processes.

Chronicle Services uses Chronicle Queue Enterprise as the means of messaging between application components. This will naturally imply persistent storage of data since Chronicle Queue Enterprise is persistent. However, guarantees that data is successfully persisted introduces delays (since persistent storage access times are significantly greater than memory access times), and these have been deemed unacceptable over many years even for normal operations in operating systems. Unix has implemented its normal I/O in an asynchronous manner for this reason, and database systems utilise sophisticated strategies to guarantee reliable persistence of data. In systems and applications where latencies must be minimised this approach would be completely unacceptable.

Chronicle’s use of memory mapped files to represent queues means that "writes" to queues will be very fast, since these will involve writes to memory. However, it also presents a window of time when the data written to the mapped memory is not synchronised with the mapped file. If the system were to stop before this output had been completed then we would be left with inconsistencies arising from lost messages.

To address these issues, Chronicle Services supports the replication of queues across multiple hosts. Messages written to a queue will be copied to replicas of the queue on different hosts, ensuring that multiple copies of the queue exist. This reduces the likelihood of all queues being subject to the data loss problem described here to almost zero.

It can be shown that the time taken to replicate writes across a TCP network to another queue is substantially less than the time taken to perform a synchronous write operation to a local disk, so replication can be achieved without a significant impact on write performance.

Chronicle Queue Enterprise Replication

Replication within Chronicle Services is based on Chronicle Queue Enterprise Replication, which is designed to ensure that the contents of a source queue are copied to one or more sink queues on other hosts, using TCP/IP network connections.

All write operations are performed on one queue, the source or master. These are copied in real time to the sink queues. In this way the ordering of messages is maintained, and the sink queues are exact replicas of the source queue.

During service startup, the replication component locks the specified source queue and handshakes with the sink queues to ensure the source has the most up-to-date queue in the cluster. If there are messages in a sink which are not in the source queue, it would indicate that the source queue had failed in some way (eg TCP/IP connection or power failure) and a failover to using a sink queue had occurred. When any such messages are replayed to the source queue the queue will be unlocked, and the application may proceed.

Figure 1. Basic Chronicle Queue Replication

A single Chronicle Queue Enterprise can be replicated across a set of N hosts (the replicaSet) according to the following:

  • The replicaSet can contain up to 127 hosts.

  • Within a replicaSet, one host (the master) acts as the source of data and is the only instance which may add data to the queue.

  • The remaining N-1 hosts (the sinks) within the replicaSet receive updates from the source.Updates are applied atomically, and a sink may be used as a live read-only copy of the data (but is not able to write to the queue).

  • Each host has an assigned hostId, and host-host connections are always initiated by the host with higher hostId.

  • Sinks can optionally send acknowledgments back to the source to confirm receipt of a replication update.

  • The source node can be changed dynamically, hence a sink queue is able to take over from a source if a failure occurs on the source.

  • A backfill mechanism ensures a source node has no gaps prior to becoming the master node within a cluster.

  • The ordering of messages is identical across all instances of a queue within a replicaSet.

An example configuration which sets up Queue replication across 3 hosts (host1, host2, host3) forming a replica set group1, with host1 as the master instance, and with acknowledgements enabled is as follows:

!ChronicleQueueReplicationCfg {
  context: {
    baseSourcePath: "replica/source",
    baseSinkPath: "replica/sink",
  },
  hosts: {
    host1: { hostId: 1, connectUri: host.port1 },
    host2: { hostId: 2, connectUri: host.port2 },
    host3: { hostId: 3, connectUri: host.port3 }
  },
  replicaSets: {
    group1: [ host1, host2, host3 ]
  },
  queues: {
    queue1: { path: queue1, replicaSets: [ group1 ], masterId: 1, acknowledge: true }
  }
}

In the above configuration, the queue queue1 is set to be replicated from host1 (as indicated by masterId) to all other hosts which are defined for replicaSet group1. Queues will use storage paths defined by baseSourcePath/baseSinkPath for source and sink, respectively, followed by the value of path variable. For the example above, the source queue will be at replica/source/queue1 while the sink will be written to replica/sink/queue1. acknowledge: true activates acknowledged replication.

Both the input queues and the output queue of a microservice can be replicated.

See runnable examples of queue replication in Chronicle-Services-Demo/Replication.

Acknowledged Replication

Queue Enterprise replication can optionally be configured for each sink instance to send an acknowledgement for each message back to the source, which in turn enables the source to keep track of the highest message which has been replicated to at least one other node.In the event of a network outage there can be no guarantee that any in-flight (unacknowledged) message has been delivered/handled by a sink, so the gap between the last sent and last acknowledged index represents the largest number of messages which could potentially be lost.

User-defined strategies can be used to control this number of "potentially losable" messages, with smaller permitted gaps coming at the cost of higher net latencies as the source message rate will be increasingly sensitive to the network round-trip for acknowledgements.

See a runnable example of processing events after acknowledged replication in Chronicle-Services-Demo/Replication/Example2.

Replication and Replay Strategy

Chronicle Services can ensure that an output message is not written to the output queue until the replication acknowledgment for the corresponding input message has been received. This is important if you are using the output queue to determine replay order (Replay Strategies) and both input and output queues are replicated. If this is the case and an input queue message was not replicated before a failure, but the corresponding output message was, the services framework may not be able to restart the service.

For more information about enabling this feature and other options, see Startup Strategies and Replay Strategies.

Configuration

Services Replication builds directly on Queue Enterprise Replication, wrapping it in the Services framework to provide additional features. Services Replication is multi-way, unidirectional with switchable master: at any given time the master host is defined and data flows from the source queues on the master host to all sinks. The master can be switched on-the-fly in order to failover, changing the direction of replication flows where appropriate.

The configuration of Services Replication closely follows that for Queue Replication, with one or more replicaSets configured across two or more hosts. Any host can be a member of one or more replicaSets. For example the following configuration file, configures replication of queues that Services are built on.

Queue replication configuration

!ChronicleServicesCfg {
  context: {
    baseSourcePath: "replica/source",
    baseSinkPath: "replica/sink$hostId",
  },

  queues: {
    queue-1: { path: queue-1, sourceId: 1 },
    queue-2: { path: queue-2, sourceId: 2, replicaSets: [ prod ], master: hostA, acknowledge: true },
    queue-3: { path: queue-3, sourceId: 3, replicaSets: [ prod ], master: hostA },
    queue-4: { path: queue-4, sourceId: 4, replicaSets: [ dr ], master: hostB }
  },

  # services can be run individually or together in or out of process
  services: {
    # bespoke services
    ...
  },

  # replication config
  hosts: {
    # higher hostId connects to lower
    hostA: { hostId: 1, connectUri: address1:7001 },
    hostB: { hostId: 2, connectUri: address2:7002 },
    hostC: { hostId: 3, connectUri: address3:7003 }
  },
  replicaSets: {
    prod: [ hostA, hostB ],
    dr: [ hostB, hostC ]
  },
}

The replicaSets block defines sets of hosts that are in a replication group. The above configuration shows two separate replicaSets - prod and dr - spanning three hosts, with hostA & hostB in prod, and hostB & hostC in dr.

The hosts block defines hosts for the queues and assigns them hostId and socket number (for hosts' communication over network) using parameters hostId and connectUri respectively.

The queues block configures the queues so that:

  • queue-1 is not replicated.

  • queue-2 and queue-3 are normally writeable on hostA (master: hostA), and any writes are replicated to the other hosts in the prod replicaSet ie hostB. The master can be changed to hostB when needed.

  • queue-2 uses acknowledged replication.

  • queue-4 is normally writeable on hostB, with any writes replicated to hostC (all the other queues in dr).The master can be changed to hostC when needed.

Alternatively, masterId: 1 can be used instead of master: hostA in the above configuration.

The context block holds the context or state about the cluster such storage path for the queues. See more information in Chronicle Queue Configuration Parameters.

Heartbeats

Chronicle Queue Enterprise replication sends heartbeats periodically to detect whether a connection has stalled. The interval of how often these heartbeats are checked, and how often they are sent, can be changed by configuration as follows:

!ChronicleServicesCfg {
  context: !QueueClusterContext {
    ...
    heartbeatIntervalMs: 2000,
    heartbeatTimeoutMs: 3000,
  },
  queues: {
    ...
  },

  services: {
    ...
  },

  hosts: {
    ...
  },
  replicaSets: {
    ...
  },
}

The queue replication heartbeat differs from the service heartbeat, which can be configured in the services block of the configuration file using parameter heartbeatMS. See the Services section in the reference guide.