You are viewing an old version of this page. View the current version.

Compare with Current View Page History

« Previous Version 54 Current »

Motivation

Ignite use a few protocols of inter-node message exchange:

  1. Communication protocol
  2. 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

  1. Communication protocol is peer-to-peer protocol: every message is sent between one node-sender and one node-receiver.
  2. Communication protocol is an internal protocol used for system messages exchange between Ignite nodes:
    1. There is a limited amount of peers, and all peers are aware of versions of each other.
    2. All peers are aware of all possible messages and their schemas.
    3. 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).
  3. 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).
  4. Messages can be pretty big, contains cache entries (putAll, historic rebalance).
  5. 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.

  1. MessageWriter, MessageReader logic is depends on a remote IgniteProductVersion.
  2. Message#writeTo, Message#readFrom is auto-generated and stored separately from Message classes.
  3. Add Message#writeObject - for serializing Objects (java functions, and user objects):
    1. MessageWriter for communication protocol use BinaryMarshaller
    2. MessageWriter for discovery protocol use JdkMarshaller
  4. Message fields contain:
    1. primitive classes.
    2. known collections.
    3. POJO - that are other Message.
    4. Users classes and java functions (e.g. ComputeJob).
    5. byte[] fields (objects serialized externally) must be avoided as much as possible, because their compatibility can't be guaranteed.
  5. Order of fields in Message is fixed with @Order annotation.
  6. Adding and removing fields must be go along with @Since, @Until annotations for Message classes and fields.

Code checks

  1. Ignite CI notifies for IF-clauses with condition based on IgniteVersion older than (curVer - 1).
  2. Ignite CI forbids code changes if Message fields changed without corresponding @Since, @Until annotations.
  3. Ignite CI forbids code changes if Message contains byte[] fields.
  4. Ignite CI forbids code changes if Message field changes type.
  5. 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:

  1. 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.
  2. 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.
  3. 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:

  1. Data - set of declared Message classes, including ser/des algorithm.
  2. Transport - algorithm of transport the Messages to remote node.

Data compatibility

Message - is base class for all messages transported between nodes. Proposed changes:

  1. Message#writeTo consumes MessageWriter that stores IgniteProductVersion of a receiver node and use it for serializing data for this version (mostly, for ignoring some fields).
  2. Message#readFrom consumes MessageReader that stores IgniteProductVersion of source node and use it for deserializing data (mostly, for setting default values of new fields).

    Message
    public interface MessageWriter {
         public IgniteProductVersion receiverVersion();
    }
    
    public interface MessageReader {
         public IgniteProductVersion senderVersion();
    }
  3. Introduce annotations @Since and @Until for Message classes and Message fields, to use it for generating code for Message#writeTo and Message#readFrom:

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 must have 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;
	}

    // Delegates #writeTo method call to generated serializer.
    @Override public boolean writeTo(ByteBuffer buf, MessageWriter writer) {
        return MyMessageSerializer.writeTo(this, buf, writer);
    }

 	// Delegates #readFrom method call to generated serializer.
    @Override public boolean readFrom(ByteBuffer buf, MessageReader reader) {
         return MyMessageSerializer.readFrom(this, buf, reader);
    }
}

// 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 static boolean writeTo(MyMessage msg, MByteBuffer buf, MessageWriter writer) {
        IgniteProductVersion rcvVer = writer.receiverVersion();

        writer.writeString(msg.id());

        if (rcvVer.lessThan(2, 20, 0))
            writer.writeString(msg.rmFld());

        if (rcvVer.greaterThanEqual(2, 20, 0))
            writer.writeString(msg.newFld());
    
        return true;
    }

    public static boolean readFrom(MyMessage msg, ByteBuffer buf, MessageReader reader) {
        IgniteProductVersion srcVer = reader.senderVersion();

        msg.id(reader.readString());

		if (srcVer.lessThan(2, 20, 0))
			msg.rmFld(reader.readString());

        if (srcVer.greaterThanEqual(2, 20, 0))
            msg.newFld(reader.readString());

        return true;
    }
}


Rules to describe Message (must be automated and validated):

  1. Do not remove Message class or Message fields, but annotate it with @Until
  2. Do not change types or @Order of fields.
  3. 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.
  4. Setters and getters must follow name of the field.

Transport compatibility 

Communication handshake

Handshake algorithm is extended on new step - validating TcpCommunicationConfiguration consistency. Settings that affects both communicating nodes must be same:

  1. usePairedConnections
  2. connectionsPerNode

Marshaller compatibility

  1. 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).
  2. BinaryMarshaller - backward compatibility is guaranteed. Require API for getting marshaller for specific version.


Other implementations

Protobuf

https://protobuf.dev/programming-guides/encoding/

  1. Field numbers are serialized with data:
    1. Do not send null fields, ignore unknown fields.
    2. Order of fields isn't guaranteed: message can be concatenations of same fields in different order. For optimizations (compression)?
    3. Make possible easily change schema (but user must preserve field ids).
    4. Serializer must know len of varlen fields (Messages). 
  2. It can be worked in streaming way using CodedOutputStream. To customize serialization, but it still requires len for Messages be written before.
  3. 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/

  1. Offsets are serialized with data:
    1. Order of fields is not guaranteed - for optimization like compaction.
    2. First 4bytes - offset to the root of vtable, that stores offsets to other fields. Vtable can be anywhere relative to fields.
    3. Size of data must be know before serializing to prepare the vtable.
    4. It starts serializing from nested objects, calculate it sizes and fill tables, and then write root object.
  2. No streaming is possible.

Avro

https://avro.apache.org/docs/

  1. 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.
  2. Schema resolution is based on field names.

Bson

  1. 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

  1. All messages are size delimited
  2. Fields order is preserved in serialization.
  3. Clients and brokers are aware of versions of each other and send messages in the form for specific versions known by each others.
  4. There are optional tagged fields beyond a message schema, that can be attached to messages.


  • No labels