DCF
Description
Concept
Distributed consensus framework (DCF) is a self-developed independent basic library with high performance, high maturity and reliability, easy to expand, and easy to use. Other systems can easily integrate DCF components through APIs. DCF is implemented based on the Paxos protocol to resolve the distributed consistency issue and improve the high reliability and availability of clusters.
Functionalities and benefits
DCF is often used in database HA scenarios to ensure the consistency of multi-copy log storage and provide high availability for external systems. In addition, DCF can automatically elect the leader node, achieving a small RTO when a single point of failure occurs.
Comparison between DCF and Quorum
- Quorum: a mechanism used in distributed systems to ensure data consistency and reliability, mainly used in multi-copy storage systems. It does not support automatic leader election and arbitration. If log forking occurs, build is required for recovery.
- DCF: a distributed consensus framework developed based on the Paxos protocol. It has the logic of automatic leader election and arbitration. When the leader node is faulty, the follower nodes automatically detect the fault and quickly elect the new leader node with a small RTO. Consistency control is performed before operations such as transaction commit and replay. Log forking does not occur when a node is faulty or a switchover is performed, avoiding unnecessary build.
Principles
Basic Principles of DCF Automatic Leader Election
As shown in Figure 1, DCF automatic leader election has three basic roles: follower, candidate, and leader.
Timeout-driven: The leader sends a heartbeat message every heartbeat timeout. If followers do not receive heartbeat messages from the leader within the election timeout period, the followers initiate an election and become candidates.
Random start: A random value is added to the election timeout period to reduce the probability of election collision.
Election action: The term increases, the voting message is broadcast, and followers become candidates. The system determines whether to vote based on logs and elects the leader if a candidate obtains more than half of the votes.
Correctness: Only one leader can be formed for the same term. Strict mathematical proof is used to prevent split-brain.
A term is equivalent to the logical clock of a distributed system. It is used to handle issues such as leadership contention between nodes, log consistency verification, and expired information identification.
Basic Principles of DCF Log Replication
As shown in Figure 2, the leader node generates logs and replicates them to the follower nodes. After the majority of logs are consistent, they are committed. Finally, the logs on each node are the same, ensuring strong consistency.
- Each log contains log term and log index.
- If a value is written to the logs on most nodes and in the current term, the logs can be committed. (For the concept of term, see the preceding description.)
- Committed logs cannot be modified and can be applied to upper-layer callers.
When the logs on the leader and follower nodes do not match, DCF automatically rectifies the fault. Figure 3 shows three nodes, a leader node for the term of 7, and two follower nodes with different log states.
For follower(a), logs are missing after the index of 4 and from indexes 1 to 4, logs are the same as those on the leader node during its term. The leader node synchronizes logs to follower(a) from the initial nextIndex of 10. However, follower(a) does not have the logs at this location and returns a message indicating the logs do not match. After receiving the response, the leader node decreases nextIndex by 1 (to 9) and continues the synchronization process. When detecting that nextIndex is 5, the leader node synchronizes subsequent logs to follower(a).
The logs on follower(b) conflict with those on the leader node. The leader node continuously synchronizes logs to follower(b) from nextIndex of 10. Until nextIndex is 3, the logs on this follower node are the same as those on the leader node during its term. At this time, the leader node synchronizes its logs to follower(b) from nextIndex of 4 and overwrites the conflicting logs after index of 4 on follower(b).
The process can be summarized as follows:
- The leader finds the start point of the follower log mismatch.
- The mismatched node logs are rewritten.
- Finally, the logs on all nodes are consistent.
Examples
- The DCF profile log stores the following information about election and replication: Figure 4 Example of DCF election and replication information
stream_id: log stream ID. The default value is 1.
node_id: node ID
pause: specifies whether flow control is enabled for a node. The value 0 indicates that flow control is disabled, and a non-zero value indicates that flow control is enabled.
role: 1: leader, 2: follower, 3: logger, 4: passive, 5: pre_candidate, 6: candidate, 7: cascade_follower.
term: term of the current node.
commit_index: consistency index of the current node.
match_index: index of the DCF log received by the current node.
next_index: index of the DCF log that the current node expects the leader node to send.
- The DCF profile log stores the following information. Figure 5 Example of the information stored by DCF
stm id: log stream ID. The default value is 1.
stm first: location of the first log in the log stream. The value is updated when dcf_truncate is called.
stm last: location of the last log in the log stream.
stg first: first log stored in the DCF data directory. The value is updated when dcf_truncate is called.
stg last: last log stored in the DCF data directory.
applied: location of the log where the DCF notifies the DN to apply (location where the DN updates the consensus LSN).
cache begin: first log point in the log message cache array of the storage module. After logs are applied or the number of logs reaches the threshold, the logs are recycled to the cache. Generally, the value of cache begin is greater than that of stm first. stm last – cache begin indicates the number of logs stored in the cache array.
See Also
"Database GUC Parameters > GUC Parameters > DCF Parameter Settings" in Reference.
What is your overall rating for this page?
Thank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot


