消息存储机制

Kafka将消息存储在Topic的特定Partition中。每个Topic可以包含多个Partition。消息在单个Partition内是有序的,但跨Partition或跨Topic时则无法保证顺序。

Partition内消息有序的原因

生产者向Partition发送消息时,消息会按顺序追加到该Partition的日志文件中,并分配唯一的offset。消费者从Partition消费时,会从最早offset开始逐个读取,从而保证消息顺序性。

实现顺序消费的方法

  1. 单Partition方案:为Topic仅创建一个Partition,所有消息将按顺序存储在同一Partition中
  2. 指定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();

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐