Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 72 additions & 0 deletions docs/_docs/clustering/multi-data-center.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,78 @@ Client nodes also use the data center ID during discovery.
If a client has a data center ID, it tries to connect through a server node from the same data center.
If no server node from the same data center is available, the client connects through any available server node.

== Timeouts

The default timeouts of discovery, communication and partition map exchange are sized for nodes in one data center.
In a cluster that spans data centers, review the timeouts described below.

=== Join Timeout

A joining server node waits until the cluster confirms its join.
The confirmation comes only after the join request has reached the coordinator and the new node has been announced
around the whole discovery ring, so it takes more than one pass of the ring.
The wait is bounded by `TcpDiscoverySpi.networkTimeout`, which is 5000 milliseconds by default.
In a ring that spans data centers, every cross-data-center hop adds its round trip, so the join can take longer than the default.
The node then stops waiting, logs `Node has not been connected to topology and will repeat join process`, and sends
its join request again, while the rest of the cluster may already be waiting for it in partition map exchange.
The node keeps repeating its join this way until it succeeds.

This timeout is a property of `TcpDiscoverySpi`. `IgniteConfiguration.networkTimeout` does not affect discovery.
Set it on the discovery SPI to a value that covers the join with your round trip between data centers, with a margin:

[source,java]
----
TcpDiscoverySpi discoverySpi = new TcpDiscoverySpi();

discoverySpi.setNetworkTimeout(15000);

IgniteConfiguration cfg = new IgniteConfiguration();

cfg.setDiscoverySpi(discoverySpi);
----

Keep the IP finder addresses precise as well: a repeated join tries the configured addresses one by one, so a long list
of addresses where no node runs (for example, wide port ranges) makes each retry slower.

=== Partition Map Exchange After a Split

When the network between data centers fails, the nodes on each side remove the nodes on the other side,
and every remaining part of the cluster runs a partition map exchange.
The exchange waits for the transactions that hold locks on the remaining nodes.
A transaction that a client node on the other side started just before the split may never finish, because the messages
that would complete it are lost.
Such a transaction is finished only when the client node is removed from the cluster, which can take up to
`IgniteConfiguration.clientFailureDetectionTimeout` (30 seconds by default).
Until then, the exchange and every operation that waits for it are blocked.
To shorten the wait, lower `IgniteConfiguration.clientFailureDetectionTimeout`; clients then also get removed sooner
when they lose connection to the cluster for other reasons.

`TransactionConfiguration.txTimeoutOnPartitionMapExchange` does not help in this case: it rolls back only the
transactions started on the nodes that run the exchange, not a transaction started by a client on the other side.
See link:key-value-api/transactions#timeout-on-partition-map-exchange[Timeout on Partition Map Exchange].

While the exchange waits, the nodes log `Failed to wait for partition map exchange` or
`Failed to wait for partition release future` with the pending transactions.
Each message first appears after twice `IgniteConfiguration.networkTimeout` and then repeats.
Raising `IgniteConfiguration.networkTimeout` postpones the messages but does not shorten the wait.

=== Failure Detection

`IgniteConfiguration.failureDetectionTimeout` (10 seconds by default for server nodes) bounds how long discovery and
communication wait for a node that stopped responding before the node is removed from the cluster.
A message between data centers takes one round trip, so the default covers the latency itself.
In a cluster that spans data centers, the timeout decides how long a network outage between data centers can last:
a short outage, well below the timeout, does not split the cluster, while a longer one can.
Choose the value from the outages the cluster should survive without a split.

The failure detection timeout applies to an SPI only while none of its own timeouts is set explicitly.
Setting `socketTimeout`, `ackTimeout`, `maxAckTimeout` or `reconnectCount` on `TcpDiscoverySpi`, or `connectTimeout`,
`maxConnectTimeout` or `reconnectCount` on `TcpCommunicationSpi`, turns it off for that SPI, which then uses those values.
If `IgniteConfiguration.failureDetectionTimeout` is set as well, the node logs
`Failure detection timeout will be ignored (one of SPI parameters has been set explicitly)`.
A cluster that spans data centers usually does not need these values: set `IgniteConfiguration.failureDetectionTimeout` instead.
See link:clustering/network-configuration#connection-timeouts[Connection Timeouts].

== Cache Topology Validation

Use `MdcTopologyValidator` to protect caches from updates when the visible topology does not contain enough data centers.
Expand Down
Loading