Kafka的Request和Response

  • 先了解Reqeust和Response的構(gòu)成, 有助于我們分析各種請求的處理過程;
  • Kafka的Request基本上分為client->server和server->server兩大類;

基礎(chǔ)數(shù)據(jù)結(jié)構(gòu)類:

Type類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/types/Type.java
  • 這是一個abstrace class, 主要是定義了ByteBuffer與各種Object之間的序列化和反序列化;
public abstract void write(ByteBuffer buffer, Object o);
public abstract Object read(ByteBuffer buffer);
public abstract Object validate(Object o);
public abstract int sizeOf(Object o);
public boolean isNullable();
  • 定義了若干Type類的實現(xiàn)類:
public static final Type INT8
public static final Type INT16
public static final Type INT32
public static final Type INT64
public static final Type STRING
public static final Type BYTES
public static final Type NULLABLE_BYTES

ArrayOf類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/types/ArrayOf.java
  • Type類的具體實現(xiàn), 是Type對象的數(shù)組類型;

Field類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/types/Field.java
  • 定義了在這個schema中的一個字段;
  • 成員:
    final int index;
    public final String name;
    public final Type type;
    public final Object defaultValue;
    public final String doc;
    final Schema schema;

Schema類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/types/schema.java
  • Schema類本身實現(xiàn)了Type類, 又包含了一個Field類對象的數(shù)組, 構(gòu)成了記錄的Schema;

Sturct類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/types/struct.java
  • 包括了一個Schema對象; 一個Object[] values數(shù)組,用于存放Schema描述的所有Field對應(yīng)的值;
    private final Schema schema;
    private final Object[] values;
  • 定義了一系列getXXX方法, 用來獲取schema中某個Field對應(yīng)的值;
  • 定義了set方法, 用來設(shè)置schema中某個Field對應(yīng)的值;
  • writeTo 用來將Stuct對象序列華到ByteBuffer;
  • Schema就是模板,Struct負(fù)責(zé)特化這個模板,向模板里添數(shù)據(jù),構(gòu)造出具體的request對象, 并可以將這個對象與ByteBuffer互相轉(zhuǎn)化;

協(xié)議相關(guān)類型:

Protocol類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/Protocol.java
  • 定義了各種Schema:
public static final Schema REQUEST_HEADER = new Schema(new Field("api_key", INT16, "The id of the request type."),
                                                           new Field("api_version", INT16, "The version of the API."),
                                                           new Field("correlation_id",
                                                                     INT32,
                                                                     "A user-supplied integer value that will be passed back with the response"),
                                                           new Field("client_id",
                                                                     STRING,
                                                                     "A user specified identifier for the client making the request."));
...

ApiKeys類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java
  • 定義了所有Kafka Api 的ID和名字
  • 如下:
    PRODUCE(0, "Produce"),
    FETCH(1, "Fetch"),
    LIST_OFFSETS(2, "Offsets"),
    METADATA(3, "Metadata"),
    LEADER_AND_ISR(4, "LeaderAndIsr"),
    STOP_REPLICA(5, "StopReplica"),
    UPDATE_METADATA_KEY(6, "UpdateMetadata"),
    CONTROLLED_SHUTDOWN_KEY(7, "ControlledShutdown"),
    OFFSET_COMMIT(8, "OffsetCommit"),
    OFFSET_FETCH(9, "OffsetFetch"),
    GROUP_COORDINATOR(10, "GroupCoordinator"),
    JOIN_GROUP(11, "JoinGroup"),
    HEARTBEAT(12, "Heartbeat"),
    LEAVE_GROUP(13, "LeaveGroup"),
    SYNC_GROUP(14, "SyncGroup"),
    DESCRIBE_GROUPS(15, "DescribeGroups"),
    LIST_GROUPS(16, "ListGroups");

Request和Response相關(guān)類型

每個Request和Response都由RequestHeader(ResponseHeader) + 具體的消費體構(gòu)成;

AbstractRequestResponse類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/request/AbstractRequestResponse.java
  • 所有Request和Response的抽象基類
  • 主要數(shù)據(jù)成員: protected final Struct struct
  • 主要接口:
public int sizeOf()
public void writeTo(ByteBuffer buffer)
public String toString()
public int hashCode()
public boolean equals(Object obj)

AbstractRequest類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/request/AbstractRequest.java
  • 繼承自AbstractReqeustResponse類, 增加了接口:
    public abstract AbstractRequestResponse getErrorResponse(int versionId, Throwable e)
  • 最重要的是它提供了一個工廠方法用于從ByteBuffer來產(chǎn)生不同類型的具體的Request;
public static AbstractRequest getRequest(int requestId, int versionId, ByteBuffer buffer) {
        switch (ApiKeys.forId(requestId)) {
            case PRODUCE:
                return ProduceRequest.parse(buffer, versionId);
            case FETCH:
                return FetchRequest.parse(buffer, versionId);
            case LIST_OFFSETS:
                return ListOffsetRequest.parse(buffer, versionId);
            case METADATA:
                return MetadataRequest.parse(buffer, versionId);
            case OFFSET_COMMIT:
                return OffsetCommitRequest.parse(buffer, versionId);
            case OFFSET_FETCH:
                return OffsetFetchRequest.parse(buffer, versionId);
            case GROUP_COORDINATOR:
                return GroupCoordinatorRequest.parse(buffer, versionId);
            case JOIN_GROUP:
                return JoinGroupRequest.parse(buffer, versionId);
            case HEARTBEAT:
                return HeartbeatRequest.parse(buffer, versionId);
            case LEAVE_GROUP:
                return LeaveGroupRequest.parse(buffer, versionId);
            case SYNC_GROUP:
                return SyncGroupRequest.parse(buffer, versionId);
            case STOP_REPLICA:
                return StopReplicaRequest.parse(buffer, versionId);
            case CONTROLLED_SHUTDOWN_KEY:
                return ControlledShutdownRequest.parse(buffer, versionId);
            case UPDATE_METADATA_KEY:
                return UpdateMetadataRequest.parse(buffer, versionId);
            case LEADER_AND_ISR:
                return LeaderAndIsrRequest.parse(buffer, versionId);
            case DESCRIBE_GROUPS:
                return DescribeGroupsRequest.parse(buffer, versionId);
            case LIST_GROUPS:
                return ListGroupsRequest.parse(buffer, versionId);
            default:
                return null;
        }
    }

實現(xiàn)上是調(diào)用各個具體Request對象的parse方法根據(jù)bytebuffer和versionid來產(chǎn)生具體的Request對象;

ProduceRequest類:

  • 我們找其中一個ProduceRqeust類來分析一下, 這個類是客戶端提交消息到broker時使用的請求;
  • 所在文件: clients/src/main/java/org/apache/kafka/common/request/ProduceRequest.java
  • 一個ProduceRequest包括下列字段:
    private final short acks;
    private final int timeout;
    private final Map<TopicPartition, ByteBuffer> partitionRecords;
  • 構(gòu)造函數(shù)public ProduceRequest(Struct struct), 利用Struct里定義的Schame來從ByteBuffer反序列化出ProduceRequest對象;
public ProduceRequest(Struct struct) {
        super(struct);
        partitionRecords = new HashMap<TopicPartition, ByteBuffer>();
        for (Object topicDataObj : struct.getArray(TOPIC_DATA_KEY_NAME)) {
            Struct topicData = (Struct) topicDataObj;
            String topic = topicData.getString(TOPIC_KEY_NAME);
            for (Object partitionResponseObj : topicData.getArray(PARTITION_DATA_KEY_NAME)) {
                Struct partitionResponse = (Struct) partitionResponseObj;
                int partition = partitionResponse.getInt(PARTITION_KEY_NAME);
                ByteBuffer records = partitionResponse.getBytes(RECORD_SET_KEY_NAME);
                partitionRecords.put(new TopicPartition(topic, partition), records);
            }
        }
        acks = struct.getShort(ACKS_KEY_NAME);
        timeout = struct.getInt(TIMEOUT_KEY_NAME);
    }

RequestHeader類:

  • 所在文件: clients/src/main/java/org/apache/kafka/common/request/RequestHeader.java
  • Request的消息頭
  • 主要成員:
    private static final Field API_KEY_FIELD = REQUEST_HEADER.get("api_key");
    private static final Field API_VERSION_FIELD = REQUEST_HEADER.get("api_version");
    private static final Field CLIENT_ID_FIELD = REQUEST_HEADER.get("client_id");
    private static final Field CORRELATION_ID_FIELD = REQUEST_HEADER.get("correlation_id");

關(guān)系圖:

request_response.png

實際上在 core/src/main/scala/kafka/api下也定義了各種Request和Response:

  • 代碼中的注釋:

NOTE: this map only includes the server-side request/response handlers. Newer
request types should only use the client-side versions which are parsed with
o.a.k.common.requests.AbstractRequest.getRequest()

     val keyToNameAndDeserializerMap: Map[Short, (String, (ByteBuffer) =>       
     RequestOrResponse)]=
     Map(ProduceKey -> ("Produce", ProducerRequest.readFrom),
        FetchKey -> ("Fetch", FetchRequest.readFrom),
        OffsetsKey -> ("Offsets", OffsetRequest.readFrom),
        MetadataKey -> ("Metadata", TopicMetadataRequest.readFrom),
        LeaderAndIsrKey -> ("LeaderAndIsr", LeaderAndIsrRequest.readFrom),
        StopReplicaKey -> ("StopReplica", StopReplicaRequest.readFrom),
        UpdateMetadataKey -> ("UpdateMetadata", UpdateMetadataRequest.readFrom),
        ControlledShutdownKey -> ("ControlledShutdown",   
        ControlledShutdownRequest.readFrom),
        OffsetCommitKey -> ("OffsetCommit", OffsetCommitRequest.readFrom),
        OffsetFetchKey -> ("OffsetFetch", OffsetFetchRequest.readFrom)
  • 這部分作解析, 沒有采用schema的形式, 是采用的直接讀取方式:
def readFrom(buffer: ByteBuffer): ProducerRequest = {
    val versionId: Short = buffer.getShort
    val correlationId: Int = buffer.getInt
    val clientId: String = readShortString(buffer)
    val requiredAcks: Short = buffer.getShort
    val ackTimeoutMs: Int = buffer.getInt
    //build the topic structure
    val topicCount = buffer.getInt
    val partitionDataPairs = (1 to topicCount).flatMap(_ => {
      // process topic
      val topic = readShortString(buffer)
      val partitionCount = buffer.getInt
      (1 to partitionCount).map(_ => {
        val partition = buffer.getInt
        val messageSetSize = buffer.getInt
        val messageSetBuffer = new Array[Byte](messageSetSize)
        buffer.get(messageSetBuffer,0,messageSetSize)
        (TopicAndPartition(topic, partition), new ByteBufferMessageSet(ByteBuffer.wrap(messageSetBuffer)))
      })
    })

請求生成與保存

  • 所有進(jìn)來的請求最終會轉(zhuǎn)換成 RequestChannel::Request, 保存在RequestChannelArrayBlockingQueue[RequestChannel.Request]中, 這個前面章節(jié)已經(jīng)講過;

Kafka協(xié)議官網(wǎng)地址

下一篇Kafka初始化流程與請求處理

Kafka源碼分析-匯總
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
  • 序言:七十年代末万矾,一起剝皮案震驚了整個濱河市晓勇,隨后出現(xiàn)的幾起案子癌刽,更是在濱河造成了極大的恐慌洞斯,老刑警劉巖,帶你破解...
    沈念sama閱讀 221,576評論 6 515
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件蔚约,死亡現(xiàn)場離奇詭異奄妨,居然都是意外死亡,警方通過查閱死者的電腦和手機(jī)苹祟,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 94,515評論 3 399
  • 文/潘曉璐 我一進(jìn)店門展蒂,熙熙樓的掌柜王于貴愁眉苦臉地迎上來,“玉大人苔咪,你說我怎么就攤上這事锰悼。” “怎么了团赏?”我有些...
    開封第一講書人閱讀 168,017評論 0 360
  • 文/不壞的土叔 我叫張陵箕般,是天一觀的道長。 經(jīng)常有香客問我舔清,道長丝里,這世上最難降的妖魔是什么? 我笑而不...
    開封第一講書人閱讀 59,626評論 1 296
  • 正文 為了忘掉前任体谒,我火速辦了婚禮杯聚,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘抒痒。我一直安慰自己幌绍,他們只是感情好,可當(dāng)我...
    茶點故事閱讀 68,625評論 6 397
  • 文/花漫 我一把揭開白布。 她就那樣靜靜地躺著傀广,像睡著了一般颁独。 火紅的嫁衣襯著肌膚如雪。 梳的紋絲不亂的頭發(fā)上伪冰,一...
    開封第一講書人閱讀 52,255評論 1 308
  • 那天誓酒,我揣著相機(jī)與錄音,去河邊找鬼贮聂。 笑死靠柑,一個胖子當(dāng)著我的面吹牛,可吹牛的內(nèi)容都是我干的吓懈。 我是一名探鬼主播歼冰,決...
    沈念sama閱讀 40,825評論 3 421
  • 文/蒼蘭香墨 我猛地睜開眼,長吁一口氣:“原來是場噩夢啊……” “哼骄瓣!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起耍攘,我...
    開封第一講書人閱讀 39,729評論 0 276
  • 序言:老撾萬榮一對情侶失蹤榕栏,失蹤者是張志新(化名)和其女友劉穎,沒想到半個月后蕾各,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體扒磁,經(jīng)...
    沈念sama閱讀 46,271評論 1 320
  • 正文 獨居荒郊野嶺守林人離奇死亡,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點故事閱讀 38,363評論 3 340
  • 正文 我和宋清朗相戀三年式曲,在試婚紗的時候發(fā)現(xiàn)自己被綠了妨托。 大學(xué)時的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片。...
    茶點故事閱讀 40,498評論 1 352
  • 序言:一個原本活蹦亂跳的男人離奇死亡吝羞,死狀恐怖兰伤,靈堂內(nèi)的尸體忽然破棺而出,到底是詐尸還是另有隱情钧排,我是刑警寧澤敦腔,帶...
    沈念sama閱讀 36,183評論 5 350
  • 正文 年R本政府宣布,位于F島的核電站恨溜,受9級特大地震影響符衔,放射性物質(zhì)發(fā)生泄漏。R本人自食惡果不足惜糟袁,卻給世界環(huán)境...
    茶點故事閱讀 41,867評論 3 333
  • 文/蒙蒙 一判族、第九天 我趴在偏房一處隱蔽的房頂上張望。 院中可真熱鬧项戴,春花似錦形帮、人聲如沸。這莊子的主人今日做“春日...
    開封第一講書人閱讀 32,338評論 0 24
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽躯枢。三九已至,卻和暖如春槐臀,著一層夾襖步出監(jiān)牢的瞬間锄蹂,已是汗流浹背。 一陣腳步聲響...
    開封第一講書人閱讀 33,458評論 1 272
  • 我被黑心中介騙來泰國打工水慨, 沒想到剛下飛機(jī)就差點兒被人妖公主榨干…… 1. 我叫王不留得糜,地道東北人。 一個月前我還...
    沈念sama閱讀 48,906評論 3 376
  • 正文 我出身青樓晰洒,卻偏偏與公主長得像朝抖,于是被迫代替她去往敵國和親牙瓢。 傳聞我的和親對象是個殘疾皇子付枫,可洞房花燭夜當(dāng)晚...
    茶點故事閱讀 45,507評論 2 359

推薦閱讀更多精彩內(nèi)容