需求

实现一个简单的双interceptor组成的拦截链。第一个interceptor会在消息发送前将时间戳信息加到消息value的最前部;第二个interceptor会在消息发送后更新成功发送消息数或失败发送消息数。

增加时间戳拦截器


  1. package com.interceptor;
  2. import java.util.Map;
  3. import org.apache.kafka.clients.producer.ProducerInterceptor;
  4. import org.apache.kafka.clients.producer.ProducerRecord;
  5. import org.apache.kafka.clients.producer.RecordMetadata;
  6. public class TimeInterceptor implements ProducerInterceptor<String, String> {
  7. @Override
  8. public void configure(Map<String, ?> configs) {
  9. }
  10. @Override
  11. public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
  12. // 创建一个新的record,把时间戳写入消息体的最前部
  13. return new ProducerRecord(record.topic(), record.partition(), record.timestamp(), record.key(),
  14. System.currentTimeMillis() + "," + record.value().toString());
  15. }
  16. //收到ack以后调用这个方法
  17. @Override
  18. public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
  19. }
  20. //关闭producer时候调用这个方法
  21. @Override
  22. public void close() {
  23. }
  24. }

统计发送消息成功和发送失败消息数,并在producer关闭时打印这两个计数器


  1. package com.interceptor;
  2. import java.util.Map;
  3. import org.apache.kafka.clients.producer.ProducerInterceptor;
  4. import org.apache.kafka.clients.producer.ProducerRecord;
  5. import org.apache.kafka.clients.producer.RecordMetadata;
  6. public class CounterInterceptor implements ProducerInterceptor<String, String>{
  7. private int errorCounter = 0;
  8. private int successCounter = 0;
  9. @Override
  10. public void configure(Map<String, ?> configs) {
  11. }
  12. @Override
  13. public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
  14. return record;
  15. }
  16. @Override
  17. public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
  18. // 统计成功和失败的次数
  19. if (exception == null) {
  20. successCounter++;
  21. } else {
  22. errorCounter++;
  23. }
  24. }
  25. @Override
  26. public void close() {
  27. // 保存结果
  28. System.out.println("Successful sent: " + successCounter);
  29. System.out.println("Failed sent: " + errorCounter);
  30. }
  31. }

生产者主程序

  1. package com.interceptor;
  2. import java.util.ArrayList;
  3. import java.util.List;
  4. import java.util.Properties;
  5. import org.apache.kafka.clients.producer.KafkaProducer;
  6. import org.apache.kafka.clients.producer.Producer;
  7. import org.apache.kafka.clients.producer.ProducerConfig;
  8. import org.apache.kafka.clients.producer.ProducerRecord;
  9. public class InterceptorProducer {
  10. public static void main(String[] args) throws Exception {
  11. // 1 设置配置信息
  12. Properties props = new Properties();
  13. props.put("bootstrap.servers", "zjj101:9092");
  14. props.put("acks", "all");
  15. props.put("retries", 0);
  16. props.put("batch.size", 16384);
  17. props.put("linger.ms", 1);
  18. props.put("buffer.memory", 33554432);
  19. props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
  20. props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
  21. // 2 构建拦截链
  22. List<String> interceptors = new ArrayList<>();
  23. interceptors.add("com.interceptor.TimeInterceptor");
  24. interceptors.add("com.interceptor.CounterInterceptor");
  25. props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
  26. String topic = "first";
  27. Producer<String, String> producer = new KafkaProducer<>(props);
  28. // 3 发送消息
  29. for (int i = 0; i < 10; i++) {
  30. ProducerRecord<String, String> record = new ProducerRecord<>(topic, "message" + i);
  31. producer.send(record);
  32. }
  33. // 4 一定要关闭producer,这样才会调用interceptor的close方法
  34. producer.close();
  35. }
  36. }

消费者主程序

  1. package com.consumer;
  2. import org.apache.kafka.clients.consumer.ConsumerRecord;
  3. import org.apache.kafka.clients.consumer.ConsumerRecords;
  4. import org.apache.kafka.clients.consumer.KafkaConsumer;
  5. import java.util.Arrays;
  6. import java.util.Properties;
  7. public class CustomConsumer {
  8. public static void main(String[] args) {
  9. Properties props = new Properties();
  10. props.put("bootstrap.servers", "zjj101:9092");
  11. props.put("group.id", "test");
  12. //enable.auto.commit:是否开启自动提交offset功能
  13. props.put("enable.auto.commit", "true");
  14. //auto.commit.interval.ms:自动提交offset的时间间隔
  15. props.put("auto.commit.interval.ms", "1000");
  16. props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
  17. props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
  18. KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
  19. //订阅 名字为 first的 topic ,可以订阅好几个topic
  20. consumer.subscribe(Arrays.asList("first"));
  21. while (true) {
  22. ConsumerRecords<String, String> records = consumer.poll(100);
  23. for (ConsumerRecord<String, String> record : records) {
  24. System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
  25. //输出: offset = 400, key = 0, value = 0
  26. }
  27. }
  28. }
  29. }

如何测试

启动一个消费者,然后运行这个生产者 就能看到效果了.