You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用SmallRye-Kafka按DatabaseID控制Kafka主题分区?

如何在Quarkus的SmallRye Kafka生产者中按DatabaseID控制消息分区?

我有一个Quarkus项目,负责创建名为databases的Kafka主题(默认100个分区)并发送数据库列表消息。现在需要实现按DatabaseID分配分区,确保相同ID的消息进入同一个分区。以下是当前的配置和代码,请教如何实现分区控制?

当前配置(application.properties)

kafka.bootstrap.servers=localhost:9092

mp.messaging.outgoing.databases.connector=smallrye-kafka
mp.messaging.outgoing.databases.key.serializer=org.apache.kafka.common.serialization.StringSerializer
mp.messaging.outgoing.databases.value.serializer=com.secupi.dbmanager.model.DatabaseSerializer
quarkus.class-loading.removed-resources."org.slf4j\:slf4j-api"=org.slf4j.Logger.class

当前代码

@Outgoing("databases")
public <T> Multi<Record<String, Database>> generate() {
    Collection<Database> dbs = databaseDao.getAll(new FilterParams());
    return Multi.createFrom().items(dbs.stream().map(db -> Record.of(db.getName(), db)));                   
}

解决方案

核心思路

Kafka生产者默认通过消息Key的哈希值计算分区,只要相同DatabaseID的消息使用一致的Key,就能保证它们进入同一个分区。只需调整消息Key的取值,并根据需求选择默认或自定义分区逻辑即可。

步骤1:将消息Key替换为DatabaseID

当前代码用db.getName()作为消息Key,无法保证相同ID的消息Key一致(比如名称可能重复或变更)。把Key改为db.getId()即可(假设Database类提供了getId()方法,返回String类型):

@Outgoing("databases")
public Multi<Record<String, Database>> generate() {
    Collection<Database> dbs = databaseDao.getAll(new FilterParams());
    // 用DatabaseID作为消息Key,确保相同ID的消息Key唯一且稳定
    return Multi.createFrom().items(dbs.stream()
            .map(db -> Record.of(db.getId(), db)));
}

步骤2:验证默认分区器行为(可选)

Kafka的DefaultPartitioner会对Key进行哈希计算,再映射到主题分区上。只要主题分区数(当前是100)保持不变,相同Key就会始终被分配到同一个分区,完全满足需求。

如果需要自定义分区逻辑(比如指定某些ID固定映射到特定分区),可以添加自定义分区器:

  1. 在application.properties中配置自定义分区器类:
mp.messaging.outgoing.databases.partitioner.class=com.secupi.dbmanager.producer.CustomDatabasePartitioner
  1. 实现自定义分区器:
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.utils.Utils;

import java.util.Map;

public class CustomDatabasePartitioner implements Partitioner {

    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        int totalPartitions = cluster.partitionCountForTopic(topic);
        // 基于DatabaseID的哈希值计算分区,确保相同ID映射到同一分区
        return Math.abs(Utils.murmur2(keyBytes)) % totalPartitions;
    }

    @Override
    public void close() {}

    @Override
    public void configure(Map<String, ?> configs) {}
}

注:自定义分区器逻辑和默认分区器类似,除非有特殊业务需求,否则无需额外实现。

步骤3:保持主题分区数稳定

分区计算依赖于主题总分区数,如果后续修改databases主题的分区数量,相同Key的消息可能会被分配到不同分区。因此,若要长期保证相同ID的消息在同一分区,需避免随意调整分区数。


内容的提问来源于stack exchange,提问作者user5260143

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 06:22:34