Motivation
Ignite use a few protocols of inter-node message exchange:
- Communication protocol
- Discovery protocol
All protocols have own serialization mechanisms and doesn't support message exchange with node of another version. For making Rolling upgrade feature possible we must make these protocols compatible between Ignite versions.
Prerequisites
- Communication protocol is peer-to-peer protocol: every message is sent between one node-sender and one node-receiver.
- Communication protocol is an internal protocol used for system messages exchange between Ignite nodes:
- There is a limited amount of peers, and all peers are aware of versions of each other.
- All peers are aware of all possible messages and their schemas.
- At most time peers are of same version. In short period of time (during RU) messags schemas might differ, but not much (few messages might differ with few fields only)
...
- .
- It is proposed to
...
- provide compatibility between versions that differ by 1 minor version, for example between 2.20.X
...
- and 2.19.X (but not between 2.20.X and 2.18.X).
- Messages can be pretty big, contains cache entries (putAll, historic rebalance).
- Users applications deployed on cluster (services, etc) must support compatibility by itself.
Current implementation
Communication protocol
- All peers hold predefined messages schemas (described with java classes).
- Messages follow schema in strict way - no fields are skipped and order of fields is guaranteed.
- Peers send messages as byte stream serializing fields one by one, no delimiters are used between fields and messages.
- For serializing POJO fields different techniques are used:
- Message consumes already serialized byte array (e.g. JdkMarshaller for IgniteDiagnosticMessage, SchemaOperationStatusMessage... BinaryMarshaller for GridJobExecuteRequest#jobBytes).
- User's BinaryObjects, Binarylizable are marshalled to byte array prior to serializing a message (the array is stored in cache structures).
- Some messages contain POJO field annotated with GridDirectTransient, and additional byte array field fill with GridCacheMessage#prepareMarshal (e.g. BinaryMarshaller for GridChangeGlobalStateMessageResponse).
- Some messages contains POJO that implements Message interface.
Discovery protocol
- All peers hold predefined messages schemas (described with java classes).
- Peers use JdkMarshaller (java serialization) to serialize and deserialize full message. Borders between messages are controlled by java serialization.
- Responses with status code or ResponseMessage.
Management command protocol
There is an additional serialization protocol between control.sh (thin client) and server nodes
- Management operates compute tasks for sending commands.
- Compute args are serialized with IgniteDataTransferObject (that is actually simple compacted Externalizable - fields without java staff).
This protocol is out of scope of this proposal.
Proposed changes
Validate joining node version
The joining node version must be checked versus actual Ignite versions in a cluster (within OnDiscoveryNodeValidationProcessor):
- If a cluster contains nodes of only version, then the joining node can be greater/less than this version by up to 1 minor version.
- If a cluster contains nodes of two versions, then joining node version must be one of them.
Message serialization framework
There should be one serialization framework for communication and discovery protocol.
MessageWriter, MessageReader logic is depends on a remote IgniteProductVersion.Message#writeTo, Message#readFrom is auto-generated and stored separately from Message classes.- Add
Message#writeObject - for serializing Objects (java functions, and user objects) with BinaryMarshaller. Message fields may contain:- primitive classes.
- known collections.
- others
Message. - Users classes and java functions (e.g. ComputeJob).
- byte[] fields (objects serialized externally) must be avoided as much as possible, because their compatibility can't be guaranteed.
- Order of fields in Message is fixed with
@Order annotation. - Adding and removing fields must be go along with
@Since, @Until annotations for Message classes and fields.
Code checks
- Ignite CI notifies for IF-clauses with condition based on IgniteVersion older than
(curVer - 1). - Ignite CI forbids code changes if Message fields changed without corresponding @Since, @Until annotations.
- Ignite CI forbids code changes if Message contains byte[] fields.
- Ignite CI forbids code changes if Message field changes type.
- Ignite CI checks @Order annotation of fields - starts with 0, no lags.
Nice to have (for later research)
These possible improvements can be implemented later. Ignite message code generator will support this features:
- Optional tagged fields. Such fields is not part of Message schema. The fields can be attached to any message. Cases: securityId, traceId, incremenalIndex, sessionAttributes, depInfo, etc.
Current approach - is creating a new message-wrapper (IncrementalSnapshotAwareMessage, TransactionAttributesAwareRequest) that wraps original Message with extra data. - Lazy deserialization/unmarshalling of specific fields. Cases: skip deserializing optional fields, transfer cache entries as byte arrays. It can be achieved by storing these fields as byte array.
- Compact length of varlen/collections - currently we use 4 bytes (int) to write length of collection or varlen type. In most cases this too much and the length can be encoded with less data (using 1-2 bytes instead).
Communication protocol
Communication protocol consist of 2 parts:
- Data - set of declared
Message classes, including ser/des algorithm. - Transport - algorithm of transport the
Messages to remote node.
Data compatibility
Message - is base class for all messages transported between nodes. Proposed changes:
- Deprecate
Message#writeTo and Message#readFrom in favor generated MessageSerializer. MessageSerializer#writeTo consumes MessageWriter that stores IgniteProductVersion of a receiver Message#writeTo consumes IgniteProductVersion of destination node and use it for serializing data for this version (mostly, for ignoring some fields).Message#readFrom MessageSerializer#readFrom consumes MessageReader that stores IgniteProductVersion of source node and use it for deserializing data (mostly, for setting default values of new fields).
| Code Block |
|---|
|
public interface MessageMessageSerializer {
public boolean writeTo(Message msg, ByteBuffer buf, MessageWriter writer, IgniteProductVersion destVer);
public boolean readFrom(Message msg, ByteBuffer buf, MessageReader reader,);
}
public interface MessageWriter {
public IgniteProductVersion srcVerreceiverVersion();
}
public interface MessageReader {
public shortIgniteProductVersion directTypesenderVersion();
} |
- Introduce annotations
@Since and @Until for Message classes and Message fields, to use it for generating code for Message#writeTo and Message#readFrom:
| Code Block |
|---|
| language | java |
|---|
| title | MyMessage |
|---|
|
// Package where the Message is defined.
package org.apache.ignite.internal.my.message;
@Since(version = "2.19.0")
public class MyMessage implements Message {
@Order(0)
private int id;
/** Remove field. */
@Until(version = "2.20.0")
@Order(1)
private String rmFld;
/** New field. */
@Since(version = "2.20.0")
@Order(2)
private String newFld;
// Message @Overridemust publichave boolean writeTo(ByteBuffer buf, MessageWriter writer, IgniteProductVersion destVer) {
if (destVer.lessThan(2, 19, 0))
throw new IgniteException("Must not send the message to destination node");
if (!writer.writeString(id))
return false;
if (destVer.lessThan(2, 20, 0)) {
setters/getters for all @Ordered fields. Methods names are equal to a corresponding field name.
public void id(int id) {
this.id = id;
}
public int id() {
return id;
}
public void newFld(String newFld) {
this.newFld = newFld;
}
public String newFld() {
return newFld;
}
}
// Generated code from the message ^ for Ignite version 2.20.0.
// Use the same package as corresponding message.
package org.apache.ignite.internal.my.message;
class MyMessageSerializer {
public boolean writeTo(Message m, ByteBuffer buf, MessageWriter writer) {
MyMessage msg = (MyMessage)m;
IgniteProductVersion rcvVer = writer.receiverVersion();
if (!writer.writeString(msg.id(rmFld));
if (rcvVer.lessThan(2, 20, 0))
return false;
}writer.writeString(msg.rmFld());
if (destVerrcvVer.greaterThanEqual(2, 20, 0)) {
if (!writer.writeString(msg.newFld())
return false;
}
return true;
}
@Overridepublic publicstatic boolean readFrom(Message m, ByteBuffer buf, MessageReader reader, IgniteProductVersion srcVer) {
) {
MyMessage msg = (MyMessage)m;
IgniteProductVersion idsrcVer = reader.readStringsenderVersion();
msg.id(reader.readString());
if (srcVer.greaterThanEquallessThan(2, 20, 0))
newFld = msg.rmFld(reader.readString());
if (srcVer.greaterThanEqual(2, 20, else0))
newFld = nullmsg.newFld(reader.readString());
return true;
}
} |
Rules to describe Message (must be automated and validated):
- Do not remove
Message class or Message fields, but annotate it with @Until. - Do not change types or order
@Order of fields. - New fields must be annotated with
@Since. All such fields must be optional, default value is null. Handling the nulls is care of Message consumer on the reader side.
...
- Setters and getters must follow name of the field.
Transport compatibility
FeatureTable
In case of transport protocol should be changed, Ignite must support -1 Ignite version.
- It's proposed to add
FeatureTable class that used to track changes in transport protocol. Each entry must be annotated with @Since version. - Ignite release must disable a communication feature with Ignite node with
(@Since - 1) version. - On release Ignite should check the
@Since version and notify release manager to drop support of old versions.
Communication handshake
Handshake algorithm is extended on new step - validating TcpCommunicationConfiguration consistency. Settings that affects both communicating nodes must be same:
...
Marshaller compatibility
- JdkMarshaller - current version RU does not support work with different Java versions. TODO: hide jdkMarshaller?BinaryMarshaller - ?- Jdk serialization is compatible between JDK versions if serialVersionUID is specified. Still some messages use it (IgniteDiagnosticMessage - should replace it is much as possible).
- BinaryMarshaller - backward compatibility is guaranteed. Require API for getting marshaller for specific version.
Other implementations
Protobuf
https://protobuf.dev/programming-guides/encoding/
- Field numbers are serialized with data:
- Do not send null fields, ignore unknown fields.
- Order of fields isn't guaranteed: message can be concatenations of same fields in different order. For optimizations (compression)?
- Make possible easily change schema (but user must preserve field ids).
- Serializer must know len of varlen fields (Messages).
- It can be worked in streaming way using CodedOutputStream. To customize serialization, but it still requires len for Messages be written before.
- It allows deserialize only required fields using CodedInputStream. It requires writing a code for iterating over tags and skipping fields.
FlatBuffers
https://flatbuffers.dev/internals/
- Offsets are serialized with data:
- Order of fields is not guaranteed - for optimization like compaction.
- First 4bytes - offset to the root of vtable, that stores offsets to other fields. Vtable can be anywhere relative to fields.
- Size of data must be know before serializing to prepare the vtable.
- It starts serializing from nested objects, calculate it sizes and fill tables, and then write root object.
- No streaming is possible.
Avro
https://avro.apache.org/docs/
- Avro relies that schema is known on both sides. And avoid writing field numbers, offsets. Only data. And it can be done in streaming way.
- Schema resolution is based on field names.
Bson
- Writes field names like json. Suppose that received doesn't have a schema.
Kafka
https://kafka.apache.org/protocol.html
KIP-482: The Kafka Protocol should Support Optional Tagged Fields
- All messages are size delimited
- Fields order is preserved in serialization.
- Clients and brokers are aware of versions of each other and send messages in the form for specific versions known by each others.
- There are optional tagged fields beyond a message schema, that can be attached to messages.