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:
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-eventswith 3 partitions and 1 replica - A topic
order-eventswith 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!"); } } }
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"));
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

