教育行業(yè)A股IPO第一股(股票代碼 003032)

全國(guó)咨詢/投訴熱線:400-618-4000

kafka自定義攔截器實(shí)例教程[傳智教育]

更新時(shí)間:2019年09月17日15時(shí)32分 來(lái)源:傳智教育 瀏覽次數(shù):

1、攔截器原理

Producer攔截器(interceptor)是在Kafka 0.10版本被引入的,主要用于實(shí)現(xiàn)clients端的定制化控制邏輯。

對(duì)于producer而言,interceptor使得用戶在消息發(fā)送前以及producer回調(diào)邏輯前有機(jī)會(huì)對(duì)消息做一些定制化需求,比如修改消息等。同時(shí),producer允許用戶指定多個(gè)interceptor按序作用于同一條消息從而形成一個(gè)攔截鏈(interceptor chain)。Intercetpor的實(shí)現(xiàn)接口是org.apache.kafka.clients.producer.ProducerInterceptor,其定義的方法包括:

(1)configure(configs)

獲取配置信息和初始化數(shù)據(jù)時(shí)調(diào)用。

(2)onSend(ProducerRecord):

該方法封裝進(jìn)KafkaProducer.send方法中,即它運(yùn)行在用戶主線程中。Producer確保在消息被序列化以及計(jì)算分區(qū)前調(diào)用該方法。用戶可以在該方法中對(duì)消息做任何操作,但最好保證不要修改消息所屬的topic和分區(qū),否則會(huì)影響目標(biāo)分區(qū)的計(jì)算。

(3)onAcknowledgement(RecordMetadata, Exception):

該方法會(huì)在消息從RecordAccumulator成功發(fā)送到Kafka Broker之后,或者在發(fā)送過(guò)程中失敗時(shí)調(diào)用。并且通常都是在producer回調(diào)邏輯觸發(fā)之前。onAcknowledgement運(yùn)行在producer的IO線程中,因此不要在該方法中放入很重的邏輯,否則會(huì)拖慢producer的消息發(fā)送效率。

(4)close:

關(guān)閉interceptor,主要用于執(zhí)行一些資源清理工作

如前所述,interceptor可能被運(yùn)行在多個(gè)線程中,因此在具體實(shí)現(xiàn)時(shí)用戶需要自行確保線程安全。另外倘若指定了多個(gè)interceptor,則producer將按照指定順序調(diào)用它們,并僅僅是捕獲每個(gè)interceptor可能拋出的異常記錄到錯(cuò)誤日志中而非在向上傳遞。這在使用過(guò)程中要特別留意。

2、攔截器案例

1)需求:

實(shí)現(xiàn)一個(gè)簡(jiǎn)單的雙interceptor組成的攔截鏈。第一個(gè)interceptor會(huì)在消息發(fā)送前將時(shí)間戳信息加到消息value的最前部;第二個(gè)interceptor會(huì)在消息發(fā)送后更新成功發(fā)送消息數(shù)或失敗發(fā)送消息數(shù)。

2)案例實(shí)操

(1)增加時(shí)間戳攔截器

package com.heima.kafka.interceptor;
import java.util.Map;
import org.apache.kafka.clients.producer.ProducerInterceptor;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
public class TimeInterceptor implements ProducerInterceptor {
@Override
public void configure(Map configs) {
}
@Override
public ProducerRecord onSend(ProducerRecord record) {
// 創(chuàng)建一個(gè)新的record,把時(shí)間戳寫入消息體的最前部
return new ProducerRecord(record.topic(), record.partition(), record.timestamp(), record.key(),
System.currentTimeMillis() + "," + record.value().toString());
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
}
@Override
public void close() {
}
}

kafka自定義攔截器

(2)統(tǒng)計(jì)發(fā)送消息成功和發(fā)送失敗消息數(shù),并在producer關(guān)閉時(shí)打印這兩個(gè)計(jì)數(shù)器

package com.heima.kafka.interceptor;
import java.util.Map;
import org.apache.kafka.clients.producer.ProducerInterceptor;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
public class CounterInterceptor implements ProducerInterceptor{
    private int errorCounter = 0;
    private int successCounter = 0;
@Override
public void configure(Map configs) {
}
@Override
public ProducerRecord onSend(ProducerRecord record) {
return record;
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
// 統(tǒng)計(jì)成功和失敗的次數(shù)
        if (exception == null) {
            successCounter++;
        } else {
            errorCounter++;
        }
}
@Override
public void close() {
        // 保存結(jié)果
        System.out.println("Successful sent: " + successCounter);
        System.out.println("Failed sent: " + errorCounter);
}
}

(3)producer主程序

package com.heima.kafka.interceptor;

import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
public class InterceptorProducer {
public static void main(String[] args) throws Exception {
// 1 設(shè)置配置信息
Properties props = new Properties();
props.put("bootstrap.servers", "hadoop102:9092");
props.put("acks", "all");
props.put("retries", 0);
props.put("batch.size", 16384);
props.put("linger.ms", 1);
props.put("buffer.memory", 33554432);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 2 構(gòu)建攔截鏈
List interceptors = new ArrayList<>();
interceptors.add("com.heima.kafka.interceptor.TimeInterceptor"); interceptors.add("com.heima.kafka.interceptor.CounterInterceptor"); 
props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
String topic = "first";
Producer producer = new KafkaProducer<>(props);
// 3 發(fā)送消息
for (int i = 0; i < 10; i++) {
    ProducerRecord record = new ProducerRecord<>(topic, "message" + i);
    producer.send(record);
}
// 4 一定要關(guān)閉producer,這樣才會(huì)調(diào)用interceptor的close方法
producer.close();
}
}
3)測(cè)試
(1)在kafka上啟動(dòng)消費(fèi)者,然后運(yùn)行客戶端java程序。

[root@hadoop102 kafka]$ bin/kafka-console-consumer.sh \
--bootstrap-server hadoop102:9092 --from-beginning --topic first
1501904047034,message0
1501904047225,message1
1501904047230,message2
1501904047234,message3
1501904047236,message4
1501904047240,message5
1501904047243,message6
1501904047246,message7
1501904047249,message8
1501904047252,message9
推薦了解:
大數(shù)據(jù)培訓(xùn)
web前端開(kāi)發(fā)
python+人工智能

0 分享到:
和我們?cè)诰€交談!