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

Kafka新手咨询:生产者多主题发布、分区创建及消费实现

Hey there! Let's break this down step by step since you're already up and running with basic Kafka producers/consumers—great start! Here's how you can handle multiple topics, partitioned messages, and the corresponding consumers in Java:

1. Creating Multiple Topics & Configuring Partitions

You have two straightforward ways to set up topics with partitions:

Command Line (Quick & Simple)

Use the kafka-topics.sh script (or .bat on Windows) from your Kafka bin directory. For example, to create two topics:

  • A topic user-events with 3 partitions and 1 replica
  • A topic order-events with 2 partitions and 1 replica
# Create user-events topic
./kafka-topics.sh --create --topic user-events --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1

# Create order-events topic
./kafka-topics.sh --create --topic order-events --bootstrap-server localhost:9092 --partitions 2 --replication-factor 1

Java Code (Programmatic Creation)

If you need to create topics dynamically in your app, use the AdminClient API:

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.NewTopic;

import java.util.Collections;
import java.util.Properties;
import java.util.concurrent.ExecutionException;

public class TopicCreator {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        Properties config = new Properties();
        config.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

        try (AdminClient adminClient = AdminClient.create(config)) {
            // Define topics with partition counts
            NewTopic userEventsTopic = new NewTopic("user-events", 3, (short) 1);
            NewTopic orderEventsTopic = new NewTopic("order-events", 2, (short) 1);

            // Create topics
            adminClient.createTopics(Collections.singletonList(userEventsTopic)).all().get();
            adminClient.createTopics(Collections.singletonList(orderEventsTopic)).all().get();

            System.out.println("Topics created successfully!");
        }
    }
}
2. Producer: Sending to Multiple Topics & Specific Partitions

Sending to Multiple Topics

It's as simple as creating ProducerRecord instances with different topic names. Here's a complete example:

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.Properties;

public class MultiTopicProducer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", StringSerializer.class.getName());
        props.put("value.serializer", StringSerializer.class.getName());

        try (Producer<String, String> producer = new KafkaProducer<>(props)) {
            // Send to user-events topic
            producer.send(new ProducerRecord<>("user-events", "user1", "Logged in"));
            producer.send(new ProducerRecord<>("user-events", "user2", "Updated profile"));

            // Send to order-events topic
            producer.send(new ProducerRecord<>("order-events", "order101", "Shipped"));
            producer.send(new ProducerRecord<>("order-events", "order102", "Delivered"));

            producer.flush();
            System.out.println("Messages sent to multiple topics!");
        }
    }
}

Sending to Specific Partitions

You can either specify the partition directly in ProducerRecord, or use a custom partitioner for rule-based assignment.

Option 1: Direct Partition Specification

// Send to partition 1 of user-events topic
producer.send(new ProducerRecord<>("user-events", 1, "user3", "Deleted account"));

Option 2: Custom Partitioner (e.g., Hash Key to Partition)

First, create the partitioner class:

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 CustomPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        int numPartitions = cluster.partitionCountForTopic(topic);
        // Hash the key to get a partition index
        return Math.abs(Utils.murmur2(keyBytes)) % numPartitions;
    }

    @Override
    public void close() {}

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

Then use it in your producer config:

props.put("partitioner.class", "com.yourpackage.CustomPartitioner");

// Now messages with the same key will go to the same partition
producer.send(new ProducerRecord<>("user-events", "user1", "Logged out"));
3. Consumer: Consuming from Multiple Topics & Specific Partitions

Consuming from Multiple Topics

Use the subscribe() method with a list of topic names:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class MultiTopicConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "multi-topic-group");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        props.put("auto.offset.reset", "earliest");

        try (Consumer<String, String> consumer = new KafkaConsumer<>(props)) {
            // Subscribe to both topics
            consumer.subscribe(Collections.singletonList("user-events,order-events".split(",")));

            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> {
                    System.out.printf("Topic: %s, Partition: %d, Key: %s, Value: %s%n",
                            record.topic(), record.partition(), record.key(), record.value());
                });
            }
        }
    }
}

Consuming from Specific Partitions

Use assign() to manually specify which partitions to consume:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

public class SpecificPartitionConsumer {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("group.id", "specific-partition-group");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        props.put("auto.offset.reset", "earliest");

        try (Consumer<String, String> consumer = new KafkaConsumer<>(props)) {
            // Assign partition 1 of user-events and partition 0 of order-events
            TopicPartition userPartition1 = new TopicPartition("user-events", 1);
            TopicPartition orderPartition0 = new TopicPartition("order-events", 0);
            consumer.assign(Collections.singletonList(userPartition1, orderPartition0));

            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                records.forEach(record -> {
                    System.out.printf("Topic: %s, Partition: %d, Key: %s, Value: %s%n",
                            record.topic(), record.partition(), record.key(), record.value());
                });
            }
        }
    }
}

A quick note: When using manual partition assignment, you lose consumer group features like rebalancing—so only use this if you need strict control over which partitions your consumer handles.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:21:27