Skip to content

Latest commit

 

History

History

Folders and files

NameName
Last commit message
Last commit date

parent directory

..
 
 
 
 

README.md

TOLERANT Match Cluster Guide

This guide explains how to configure and run TOLERANT Match in a cluster for production use.

It is written for administrators and operators.

0. Quickstart

Use this as the shortest production-oriented flow:

  1. Prepare at least 2 nodes (3 preferred), synchronized time (NTP), and node-to-node network access.
  2. Configure one shared cluster identity (clusterName) and consistent cluster ports on all nodes.
  3. Use one discovery method across all nodes (Kubernetes service discovery, static TCP list, or AWS discovery).
  4. Ensure each node has unique node identity and separate local data storage.
  5. Start nodes and verify all members join the same cluster and report expected clusterSize.
  6. Route traffic only to healthy nodes and only when project state is SYNCHRONIZED.
  7. Monitor synchronization and cluster metrics; investigate any ERROR state before full traffic.

Path selection:

  • Kubernetes/Helm operators: start with sections 4.5, 4.7, 7, and 8.
  • On-prem/custom operators: start with sections 4.4, 4.8, 4.9, 17.1, and 17.2.

1. What a Match cluster does

A Match cluster runs two or more Match service nodes for:

  • higher availability
  • better request throughput
  • controlled data synchronization between nodes

In a typical setup, clients call a load balancer, and the load balancer forwards requests to cluster nodes.

1.1 Example topology

Cluster topology

2. Recommended architecture

Use this baseline architecture:

  1. At least 2 Match nodes (3 preferred for production).
  2. One load balancer in front of Match nodes.
  3. Separate local data storage per node.
  4. Stable node identity and stable network naming.
  5. Time synchronization (NTP) on all nodes.

Important: cluster nodes must not share the same local data directory or database schema tables.

2.1 Do and don't

Do:

  • keep Match versions identical across all cluster nodes
  • keep cluster/discovery configuration consistent across all nodes
  • keep local node data directories separated per node
  • route production traffic only to projects in SYNCHRONIZED

Don't:

  • do not share one local data directory across nodes
  • do not mix discovery methods unintentionally in the same cluster
  • do not route full production traffic to SYNCHRONIZING or ERROR project states
  • do not change cluster identity/ports on only a subset of nodes
  • do not keep Savepoint Settings on all nodes the same if copying the same config files between servers

3. Prerequisites checklist

Before enabling cluster mode, confirm:

  • all nodes run the same Match version
  • all nodes have compatible platform/runtime setup
  • required node-to-node ports are open
  • load balancer can reach all Match nodes
  • all nodes can resolve each other by DNS/IP
  • project configuration is aligned across nodes

4. Cluster communication and configuration files

Match uses Hazelcast configuration files for cluster discovery and communication.

Quick path selection:

  • Kubernetes/Helm: use chart-generated hazelcast.yaml (recommended)
  • On-prem/custom: use external Hazelcast config files, or matchRuntime fallback if no Hazelcast config file exists

4.1 Supported file names and lookup order

Place one of these files in the Match config directory:

  1. hazelcast.yaml
  2. hazelcast.yml
  3. hazelcast.xml
  4. clusterconfig.yaml
  5. clusterconfig.yml
  6. clusterconfig.xml

The first file found in this order is used.

4.2 YAML or XML

Both YAML and XML are supported.
Use one format consistently across all nodes.

4.3 Discovery modes

Common discovery modes:

  • Kubernetes service discovery
  • static TCP member list
  • AWS discovery (for AWS profile deployments)

4.4 Suggested static TCP/IP configuration

Use static member discovery when multicast or platform discovery is not suitable.

Example (hazelcast.yaml):

hazelcast:
  cluster-name: match-cluster
  network:
    port:
      port: 5701
    join:
      multicast:
        enabled: false
      tcp-ip:
        enabled: true
        member-list:
          - 192.168.10.11:5701
          - 192.168.10.12:5701
          - 192.168.10.13:5701

Guidelines:

  • disable multicast when using TCP/IP member list
  • list all expected cluster members
  • keep cluster name and port consistent on all nodes
  • ensure firewall rules allow member-to-member traffic

4.5 Suggested Kubernetes configuration

For Kubernetes, prefer chart-managed Hazelcast config and StatefulSet identity.

Example Helm values:

global:
  db:
    type: rocksdb
    rocksdbCluster:
      enabled: true
  match:
    environmentProfile: k8s
    clusterSyncServicePort: 10100
    hazelcast:
      enabled: true
      port: 5701
      clusterName: match-cluster
      joinMode: kubernetes

Notes:

  • use stable pod identity (StatefulSet)
  • keep per-pod persistent storage
  • keep readiness checks enabled and route traffic only to ready pods

4.6 Which configuration is used? (precedence)

Important:

  • Match checks Hazelcast files in the lookup order from section 4.1 (hazelcast.yaml, hazelcast.yml, hazelcast.xml, clusterconfig.yaml, clusterconfig.yml, clusterconfig.xml)
  • if a file is found, that file is used for Hazelcast cluster configuration
  • matchRuntime cluster attributes are used only when no Hazelcast configuration file was found

4.7 Kubernetes/Helm

In Helm/Kubernetes deployments:

  • when global.match.hazelcast.enabled=true, the chart generates and mounts hazelcast.yaml
  • this generated file is the effective runtime cluster configuration
  • use Helm values (for example global.match.hazelcast.joinMode, global.match.hazelcast.tcpMembers, global.match.hazelcast.aws.*, ports, and cluster name) instead of manually editing external clusterconfig.* files

4.8 On-prem activation

For on-prem installations, cluster mode is activated by setting cluster identity in matchRuntime.

Minimum activation:

  • set clusterName in <matchRuntime ...>

All service instances with the same clusterName are treated as one cluster.

4.9 Fallback: matchRuntime cluster attributes

These attributes are used as fallback baseline cluster settings only if no Hazelcast config file is found:

  • clusterName: shared cluster identity across nodes
  • clusterPort: port for cluster member communication
  • clusterSyncServicePort: port for synchronization data exchange service
  • clusterSyncServiceEncrKey (optional): enables encrypted sync payload exchange when set

Example:

<matchRuntime clusterName="cl1" clusterPort="5701" clusterSyncServicePort="10100" />

For advanced network/member settings, use Hazelcast configuration files.

5. Load balancer guidance

You can use software or hardware load balancers.
A reverse proxy setup (for example, Apache/Nginx) is common.

Recommendations:

  • route client traffic only to healthy Match nodes
  • keep health checks active
  • use session stickiness only if your architecture requires it
  • for clustered projects, route full traffic only when project state is SYNCHRONIZED

6. Cluster-related Match settings

Configure consistent cluster identity and ports across nodes:

  • cluster name (shared by all nodes in one cluster)
  • node name / service identity (unique per node)
  • cluster sync service port
  • Hazelcast/member communication port

If values differ unintentionally between nodes, synchronization issues may occur.

7. Kubernetes / Helm quick start

For Helm-based deployments:

  • keep global.match.environmentProfile: k8s (default)
  • enable Hazelcast for clustered synchronization
  • run multi-replica setups as StatefulSet for stable identity
  • keep per-node persistent storage

For RocksDB multi-replica deployments:

  • enable RocksDB cluster mode
  • enable Hazelcast
  • keep per-pod PVCs

For SQL deployments (Postgres/MariaDB):

  • use node-specific service IDs/prefixes
  • keep SQL connection settings valid for all replicas

8. Node and project status checks

Check cluster status from the Match admin interface/API/CLI.

Typical project cluster states:

  • SYNCHRONIZING: node is catching up and should not receive production traffic
  • SYNCHRONIZED: node is ready for normal traffic
  • ERROR: synchronization/configuration problem; investigate before routing traffic

Operational rule: only route full traffic to nodes in SYNCHRONIZED.

8.1 Status check example

The following example shows the readiness of a cluster with four members:

~ $ service.sh backend --endpoint health --function readiness

  _____ ___  _    ___ ___    _   _  _ _____   __  __      _      _    
 |_   _/ _ \| |  | __| _ \  /_\ | \| |_   _| |  \/  |__ _| |_ __| |_  
   | || (_) | |__| _||   / / _ \| .` | | |   | |\/| / _` |  _/ _| ' \ 
   |_| \___/|____|___|_|_\/_/ \_\_|\_| |_|   |_|  |_\__,_|\__\__|_||_|
                                                                      

 Version 12.1.58634, Copyright (c) 2026 TOLERANT Software GmbH & Co KG

{
 "name": "match",
 "status": "UP",
 "details": {
  "projects": {
   "name": "match",
   "status": "UP",
   "details": {
    "matchProject-1": {
     "name": "match",
     "status": "UP",
     "details": {
      "errorMessage": "",
      "cluster-information": {
       "cluster-project-state": "SYNCHRONIZED",
       "cluster-mode": "MASTER_SLAVE",
       "cluster-size": 4,
       "cluster-node-name": "tolerant-match-smoke-match-1",
       "cluster-uuid": "153d789f-3dde-4cc9-a1cc-24d92a901241",
       "cluster-name": "match-cluster"
      },
      "componentStatus": "RUNNING",
      "active": true,
      "configurationTime": "2026-05-22T11:34:43.033157072"
     }
    }
   }
  },
  "compositeDiscoveryClient()": {
   "name": "match",
   "status": "UP",
   "details": {
    "services": {}
   }
  },
  "readiness": {
   "name": "match",
   "status": "UP",
   "details": {
    "runtimeState": "RUNNING",
    "projectCount": 1,
    "allProjectsRunning": true,
    "allClusterProjectsSynchronized": true,
    "unsynchronizedClusterProjects": []
   }
  },
  "service": {
   "name": "match",
   "status": "UP"
  },
  "diskSpace": {
   "name": "match",
   "status": "UP",
   "details": {
    "total": 1081101176832,
    "free": 1003819864064,
    "threshold": 10485760
   }
  },
  "runtime": {
   "name": "match",
   "status": "UP",
   "details": {
    "cluster-information": {
     "cluster-name": "match-cluster",
     "cluster-member-state": "CONNECTED",
     "cluster-node-name": "tolerant-match-smoke-match-1",
     "cluster-UUID": "153d789f-3dde-4cc9-a1cc-24d92a901241"
    },
    "errorMessage": "",
    "componentStatus": "RUNNING",
    "configurationTime": "2026-05-22T11:34:38.426613761"
   }
  },
  "gracefulShutdown": {
   "name": "match",
   "status": "UP",
   "details": {
    "activeTasks": 1
   }
  }
 }
}
Processing successful (RC: 0, Elapsed time: 1s)

The same result can be achieved by using the following http request:

 wget -q -O- http://<Servername>:<Serverport>/health/readiness

The following metrics give you a quick view over the cluster state:

  • cluster.member.state
  • cluster.project.member-count
  • cluster.project.state

State code mapping (cluster.member.state / ClusterMemberState):

  • 0: not available / cluster disabled
  • 1: UNCONFIGURED, CONFIGURED, CONNECTED, SPLIT_BRAIN
  • 4: ERROR

State code mapping (cluster.project.state / ClusterProjectState):

  • 0: not available / cluster disabled
  • 1: UNCONFIGURED, CONFIGURED_ACTIVE, CONFIGURED_INACTIVE
  • 2: SYNCHRONIZING, SYNCHRONIZING_WAITING
  • 3: SYNCHRONIZED
  • 4: ERROR

Example:

~ $ service.sh backend --endpoint metrics --function cluster.project.member-count

  _____ ___  _    ___ ___    _   _  _ _____   __  __      _      _    
 |_   _/ _ \| |  | __| _ \  /_\ | \| |_   _| |  \/  |__ _| |_ __| |_  
   | || (_) | |__| _||   / / _ \| .` | | |   | |\/| / _` |  _/ _| ' \ 
   |_| \___/|____|___|_|_\/_/ \_\_|\_| |_|   |_|  |_\__,_|\__\__|_||_|
                                                                      

 Version 12.1.58634, Copyright (c) 2026 TOLERANT Software GmbH & Co KG

{
 "name": "cluster.project.member-count",
 "measurements": [{
  "statistic": "VALUE",
  "value": 4
 }],
 "availableTags": [{
  "tag": "projectId",
  "values": ["matchProject-1"]
 }],
 "description": "Cluster member count for this project (0 when project is not clustered)"
}
Processing successful (RC: 0, Elapsed time: 1s)

9. How synchronization behaves (practical view)

When a node restarts or rejoins:

  1. it finds a suitable peer
  2. it synchronizes missing data
  3. it becomes SYNCHRONIZED

Synchronization can use backlog/history when available, which is usually faster than full transfer.

If history is insufficient, Match can use larger recovery transfer paths, which take longer.

10. Backlog/history retention recommendation

Keep history retention long enough to cover expected outage windows.

Practical recommendation:

  • set retention so normal node restarts can recover from backlog/history
  • avoid very short retention values that force frequent full re-sync
  • configure retention with project attribute maxHistoryDuration (seconds)

11. Initial load and re-sync operations

If you run an initial load on one node while other nodes are active:

  • expect synchronization work after startup/rejoin
  • monitor all affected projects until they return to SYNCHRONIZED

Do not treat the cluster as fully available until synchronization completes.

12. Monitoring and operations

Track at least:

  • project cluster state transitions
  • synchronization duration and failures
  • repeated retries/timeouts
  • node availability and readiness endpoints (/health)

Use alerts for:

  • nodes stuck in SYNCHRONIZING
  • any ERROR state
  • recurring sync failures after restart events

13. Common problems and fixes

Problem: node does not join cluster

Check:

  • same cluster name on all nodes
  • matching discovery method/config
  • network/firewall access to cluster ports
  • config file is present and valid (YAML/XML syntax)

Problem: node remains in synchronizing state too long

Check:

  • donor node health
  • backlog/history retention settings
  • storage and network throughput
  • very large data changes since node went offline

Problem: configuration mismatch warnings

Check:

  • project config consistency across nodes
  • runtime config consistency across nodes
  • accidental drift after manual changes

14. Security and hardening basics

  • allow cluster ports only within trusted networks
  • do not expose internal cluster sync ports to the public internet
  • protect credentials/secrets via your platform secret manager
  • use least-privilege firewall and security group rules

15. Change management recommendations

For safe production changes:

  1. apply config changes in staging first
  2. validate node join and synchronization
  3. roll out progressively in production
  4. verify all projects return to SYNCHRONIZED

16. Tolerant Match Version Upgrade

  • Always verify functionality on testing environment before production
  • Check whether a full initial load is required after a product update
  • Verify that all nodes use the same Tolerant Match major and minor version
  • Follow steps of playbook to ensure full functionality of services
  • Perform UAT or Testcases to ensure response integrity
  • After successful setup on testing environment move to production

17. Playbook (on-prem)

17.1 Initial setup (on-prem)

  1. Prepare infrastructure

    • use at least two nodes (three recommended)
    • ensure node-to-node connectivity for clusterPort and clusterSyncServicePort
    • ensure load balancer can reach all nodes
    • ensure time sync (NTP) on all nodes
  2. Configure cluster identity and ports

    • choose one cluster name for all nodes
    • configure either:
      • Hazelcast config file (hazelcast.yaml / hazelcast.xml / clusterconfig.*), or
      • matchRuntime fallback attributes (clusterName, clusterPort, clusterSyncServicePort)
    • ensure each node has unique service identity/node identity
  3. Start services

    • start node A
    • start node B (and remaining nodes)
    • wait until each node is reachable and healthy
  4. Configure and enable load balancer routing

    • add all nodes as targets
    • route only to healthy/ready nodes
  5. Verify cluster join

    • check that all nodes report the same clusterName
    • check expected clusterSize
    • check project states move to SYNCHRONIZED

17.2 Check everything is working

Use this checklist after setup and after operational changes.

Node health checks:

  • backend process is running on every node
  • /health endpoints respond successfully
  • required ports are open/listening

Cluster checks:

  • all nodes show same cluster name
  • clusterSize equals expected node count
  • project state is SYNCHRONIZED on all active nodes
  • no persistent cluster errors in logs

Load balancer checks:

  • requests through LB are distributed to active nodes
  • unhealthy nodes are not receiving production traffic

Synchronization checks:

  • restart one node and verify it returns to SYNCHRONIZED
  • verify other nodes stay available during that restart
  • verify data changes are visible across nodes after synchronization

Optional resilience test:

  • stop one node intentionally
  • verify service remains available via remaining node(s)
  • start node again and confirm successful rejoin/synchronization

Acceptance checklist (pass/fail):

  • all nodes healthy
  • cluster joined with expected size
  • all projects synchronized
  • LB routing healthy
  • restart/rejoin synchronization successful

17.3 Restart and recovery (single node)

Use this recovery flow when one restarted node does not return cleanly:

  1. Remove the affected node from load balancer routing.
  2. Stop the affected Match service instance.
  3. Verify the old process is fully terminated and releases file handles.
  4. Start the node again and wait for health/readiness.
  5. Verify cluster member count and project state reaches SYNCHRONIZED.
  6. Re-enable load balancer routing for that node.

17.4 Loading new data

Use this flow when new project data must be loaded into an existing cluster.

  1. Select one source node for the load.

    • Prefer a healthy node that is currently SYNCHRONIZED.
    • Remove the source node from load balancer routing if the load may affect response latency or project availability.
  2. Run the initial load on the selected source node.

    • Stop the affected project on that node before loading.
    • Set --data-epoch if the data is to be rolled out to all other cluster members.
    • Keep the other nodes running so the service remains available through synchronized nodes.

    Linux shell example:

    service.sh backend --endpoint operations --function stop.project --parameter "projectId=matchProject-1"
    matchInitialLoad.sh -delete-backlog config/matchserviceconfig.xml matchProject-1 --data-epoch 1713700000000
    service.sh backend --endpoint operations --function start.project --parameter "projectId=matchProject-1"

    Windows cmd example:

    service.exe backend --endpoint operations --function stop.project --parameter "projectId=matchProject-1"
    matchInitialLoad.bat -delete-backlog config\matchserviceconfig.xml matchProject-1 --data-epoch 1713700000000
    service.exe backend --endpoint operations --function start.project --parameter "projectId=matchProject-1"
  3. Verify that the affected project is running again on the source node.

  4. Monitor cluster synchronization.

    • Expect other nodes to enter synchronization while they catch up.
    • Do not route full production traffic to nodes whose affected project is SYNCHRONIZING or ERROR.
  5. Verify completion.

    • all nodes show the expected clusterSize
    • all affected projects return to SYNCHRONIZED
    • smoke-test representative requests and response integrity
  6. Re-enable normal load balancer routing after validation.

18. Scope and related documents

  • This file is the customer-facing cluster manual.
  • Helm deployment details remain in helm/README.md.

19. Failure recovery scenarios

These scenarios summarize expected cluster failure recovery behavior:

  1. One node fails, savepoints available.
  2. One node fails, no savepoint for restarting node.
  3. One node fails, no usable history.
  4. One node fails, no usable history and limited savepoint coverage.
  5. One node fails, no savepoints available.
  6. One node fails, no savepoint on preferred master node.
  7. One node fails, no savepoints on either node.
  8. Two nodes fail at different times, savepoints available.
  9. Two nodes fail at different times, no usable history.
  10. Cluster communication outage, savepoint available.
  11. Cluster communication outage, no usable history.

Operational interpretation:

  • If backlog/history is sufficient, recovery is usually faster.
  • If backlog/history is insufficient, full savepoint/database/index transfer paths may be required.
  • A higher maxHistoryDuration typically reduces expensive full recovery operations.

19.1 Recovery scenario descriptions

The descriptions below use a consistent scenario format and should be read together with each scenario diagram.

Common context for all scenario images:

  • two nodes (A, B)
  • F: represents an error
  • S: start
  • SP: full savepoint
  • H: history start time
  • T(x): timestamp of event x

Timing expressions such as T(S) > T(F) mean event S happened after event F.

Important note:

  • importing an incremental savepoint alone from another cluster node is not possible
  • incremental savepoints always depend on the latest full savepoint
  • therefore, only full savepoints are exchanged directly
  • for database-based synchronization, Match transfers the full savepoint and all incremental savepoints written after that full savepoint

Scenario 1: One node fails, savepoints available
Context: Node B fails and is restarted.
Timing relation: T(S) > T(SP A2) > T(F) > T(H A)
Recovery: use data from SP B1, Backlog B, and History A > T(F) / Backlog A. Scenario 1

Scenario 2: One node fails, no savepoint available for restarting node
Context: Node B fails and is restarted; no savepoint exists for B.
Timing relation: T(S) > T(SP A2) > T(F) > T(H A), no SP for B
Recovery: use data from Backlog B, History A > T(F), and Backlog A. Scenario 2

Scenario 3: One node fails, no history available
Context: Node B fails and is restarted; required history is not available.
Timing relation: T(S) > T(SP A2) > T(H A) > T(F)
Recovery: use data from SP A2 (full savepoint + increments); Backlog B is discarded. Scenario 3

Scenario 4: One node fails, no history and no local savepoint available
Context: Node B fails and is restarted; history is missing and local savepoint is missing.
Timing relation: T(S) > T(SP A2) > T(H A) > T(F), no SP for B
Recovery: use data from SP A2 (full savepoint + increments); Backlog B is discarded. Scenario 4

Scenario 5: One node fails, no usable savepoints in current path
Context: Node B fails and is restarted.
Timing relation: T(S) > T(SP B1) > T(SP A2)
Recovery: use data from SP B1, plus Backlog B and Backlog A > T(F). Scenario 5

Scenario 6: One node fails, no savepoint for master
Context: Node B fails and is restarted; no savepoint exists for A, savepoint exists for B.
Timing relation: no SP for A, SP for B exists
Recovery: use data from SP B1, Backlog B, and Backlog A > T(F). Scenario 6

Scenario 7: One node fails, no savepoints available
Context: Node B fails and is restarted; no savepoints exist for A or B.
Timing relation: no SP for A and B
Recovery: use data from Backlog B and Backlog A > T(F). Scenario 7

Scenario 8: Two nodes fail at different times, savepoints available
Context: Nodes A and B fail in sequence.
Timing relation: T(S A) > T(S B) > T(F A) > T(F B)
Recovery:
A uses Backlog A with Backlog B > T(F A).
B at T(S B): no action.
B at T(S A): use Backlog A > T(F B) for keys without newer updates. Scenario 8

Scenario 9: Two nodes fail at different times, no history available
Context: Nodes A and B fail in sequence; history coverage is insufficient.
Timing relation: T(S A) > T(S B) > T(F A) > T(H A) > T(F B)
Recovery:
A uses Backlog A with Backlog B > T(F A).
B at T(S B): no action.
B at T(S A): use SP A1 (full savepoint + increments), after A synchronization. Scenario 9

Scenario 10: Cluster communication fails, savepoint available
Context: Cluster communication is interrupted temporarily.
Timing relation: T(CS) > max(T(SP A1), T(SP B1))
Recovery: A uses Backlog B > T(CS) and B uses Backlog A > T(CS). Scenario 10

Scenario 11: Cluster communication fails, no history available
Context: Cluster communication is interrupted and history coverage is insufficient.
Timing relation: T(CR) > T(SP B1) > T(H B) > T(CS) or T(H A) > T(CS)
Recovery: broader export/import reconciliation between A and B; in key conflicts, master-node entries win. Scenario 11