Kafka如何实现顺序消费?
·
消息存储机制
Kafka将消息存储在Topic的特定Partition中。每个Topic可以包含多个Partition。消息在单个Partition内是有序的,但跨Partition或跨Topic时则无法保证顺序。
Partition内消息有序的原因
生产者向Partition发送消息时,消息会按顺序追加到该Partition的日志文件中,并分配唯一的offset。消费者从Partition消费时,会从最早offset开始逐个读取,从而保证消息顺序性。
实现顺序消费的方法
- 单Partition方案:为Topic仅创建一个Partition,所有消息将按顺序存储在同一Partition中
- 指定Partition方案:将需要保证顺序的消息发送到同一Partition
消息定向发送方法
1. 直接指定Partition
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
public class KafkaProducerExample {
public static void main(String[] args) {
Producer<String, String> producer = new KafkaProducer<>(getProperties());
String topic = "hollis_topic";
String message = "Hello World!";
int partition = 0;
ProducerRecord<String, String> record =
new ProducerRecord<>(topic, partition, null, message);
producer.send(record);
producer.close();
}
}
2. 通过Key指定
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
public class KafkaProducerExample {
public static void main(String[] args) {
Producer<String, String> producer = new KafkaProducer<>(getProperties());
String topic = "hollis_topic";
String message = "Hello World!";
String key = "Hollis_key";
ProducerRecord<String, String> record =
new ProducerRecord<>(topic, key, message);
producer.send(record);
producer.close();
}
}
3. 自定义Partitioner
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.record.InvalidRecordException;
public class CustomPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
int numPartitions = partitions.size();
if (keyBytes == null || !(key instanceof String)) {
throw new InvalidRecordException("键不能为空且必须是字符串类型");
}
return Math.abs(((String)key).hashCode()) % numPartitions;
}
@Override public void configure(Map<String, ?> configs) {}
@Override public void close() {}
}
配置使用自定义Partitioner:
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("partitioner.class", "com.hollis.CustomPartitioner");
Producer<String, String> producer = new KafkaProducer<>(props);
ProducerRecord<String, String> record =
new ProducerRecord<>("hollis_topic", "Hollis_key", "Hello World!");
producer.send(record);
producer.close();
更多推荐
所有评论(0)