You are viewing an old version of this page. View the current version.
Compare with Current
View Page History
« Previous
Version 39
Next »
Motivation
Currently Communication protocol doesn't support message exchange with node of another version. For making Rolling upgrade feature possible we must make the protocol 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 of Communication protocol:
- All peers hold predefined
- 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.
GridCacheMessage holds additional field depInfo.
Discovery protocol
Current implementation of Discovery protocol:
- Peers use JdkMarshaller (java serialization) to serialize and deserialize a message.
Proposed changes
Message DTO
MessageWriter, MessageReader logic is depends on a remote IgniteProductVersion.Message#writeTo, Message#readFrom is auto-generated and stored separately from Message DTO classes.- Add
Message#writeJavaObject - for serializing Java functions and user objects (ComputeJob) with JdkMarshaller. Message fields contain:- primitive classes
- known collections
- POJO - that are other
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
- Add
@Since, annotations for Message classes and fields.
Version check
- Ignite CI must notify for IF-clauses with condition based on IgniteVersion older than
(curVer - 1). - Ignite CI must forbid code changes if Message DTO changed without corresponding @Since, @Until annotations.
- Ignite CI must forbid code changes if Message DTO contains byte[] fields.
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:
Message#writeTo consumes MessageWriter that stores IgniteProduceVersion of destination node and use it for serializing data for this version (mostly, for ignoring some fields).Message#readFrom consumes MessageReader that stores IgniteProductVersion of source node and use it for deserializing data (mostly, for setting default values of new fields).
public interface Message {
public boolean writeTo(ByteBuffer buf, MessageWriter writer);
public boolean readFrom(ByteBuffer buf, MessageReader reader);
public short directType();
}
- annotations
@Since and @Until for Message classes and Message fields, to use it for generating code for Message#writeTo and Message#readFrom:
// Package that store all schemas.
package org.apache.ignite.internal.messages.schema;
// package private class.
@Since(version = "2.19.0")
class MyMessage {
private int id;
/** Remove field. */
@Until(version = "2.20.0")
private String rmFld;
/** New field. */
@Since(version = "2.20.0")
private String newFld;
}
// Generated code from the schema ^.
package org.apache.ignite.internal.messages;
// public class.
public class MyMessage {
private int id;
/**
* Remove field.
* @deprecated since 2.20.0.
*/
@Deprecated
private String rmFld;
/**
* New field.
* @since 2.20.0
*/
private String newFld;
@Override public 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)) {
if (!writer.writeString(rmFld))
return false;
}
if (destVer.greaterThanEqual(2, 20, 0)) {
if (!writer.writeString(newFld))
return false;
}
return true;
}
@Override public boolean readFrom(ByteBuffer buf, MessageReader reader, IgniteProductVersion srcVer) {
id = reader.readString();
if (srcVer.greaterThanEqual(2, 20, 0))
newFld = reader.readString();
else
newFld = null;
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 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.
Removing annotated entities is allowed after current version is greater than (@Until + 1) or (@Since + 1).
Transport compatibility
Communication handshake
Settings that affects both communicating nodes must be same:
- usePairedConnections
- connectionsPerNode
Marshaller compatibility
- JdkMarshaller - 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.