Kafka에서 메시지 소비
Kafka 토픽의 데이터 스트림을 읽고 처리하며 오프셋과 컨슈머 그룹을 다루는 컨슈머를 구축하는 방법을 배웁니다.
Kafka에서 메시지 소비은(는) CoddyKit의 무료 Apache Kafka & Stream Processing Fundamentals 강의입니다. 이것은 4개 중 2번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 Apache Kafka & Stream Processing Fundamentals 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. Apache Kafka & Stream Processing Fundamentals 강의에는 총 4개의 강의가 포함되어 있습니다.
이 강의의 일부는 아직 번역되지 않았으며 영어로 표시됩니다.
What are Kafka Consumers?
Kafka Consumers are applications that read data from Kafka topics. They subscribe to one or more topics and process the messages as they arrive.
- Consumers are the 'listeners' in your Kafka ecosystem.
- They pull data from Kafka brokers, rather than having data pushed to them.
- Essential for building real-time data pipelines and analytics systems.
Working with Consumer Groups
Consumers often work together in a consumer group. This is a collection of consumer instances that share a common group.id.
- Scalability: Messages from a topic's partitions are distributed among consumers in the same group, allowing parallel processing.
- Fault Tolerance: If a consumer instance fails, its assigned partitions are automatically reassigned to other active consumers in the group.
- Each message in a partition is delivered to only one consumer within a group.
How Consumers Read Data (Polling)
Kafka consumers don't just 'receive' messages. Instead, they actively poll Kafka brokers for new data.
- The consumer repeatedly asks Kafka, "Do you have any new records for me?"
- This polling mechanism gives the consumer control over its processing rate.
- If no new records are available, the
poll()method will simply return an empty list after waiting for a specified duration.
Basic Java Consumer Setup
To create a Kafka consumer in Java, you use the KafkaConsumer class. First, you need to set up some basic properties like the Kafka server address and deserializers.
Here's the start of a simple consumer program:
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Properties;
import java.util.Collections;
public class SimpleConsumerSetup {
public static void main(String[] args) {
// 1. Define consumer properties
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-first-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// 2. Create the consumer
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
// 3. Subscribe to a topic
consumer.subscribe(Collections.singletonList("my-topic"));
// Consumer loop and closing will go here
System.out.println("Consumer setup complete. Subscribed to 'my-topic'.");
consumer.close();
}
}Essential Consumer Properties
Let's break down the key properties we set for our consumer:
BOOTSTRAP_SERVERS_CONFIG: A list of Kafka broker addresses (e.g.,localhost:9092) that the consumer will connect to.GROUP_ID_CONFIG: A unique identifier for the consumer group this consumer belongs to (e.g.,my-first-group).KEY_DESERIALIZER_CLASS_CONFIG: Specifies how to convert message keys from bytes (received from Kafka) into Java objects (e.g.,StringDeserializerfor text keys).VALUE_DESERIALIZER_CLASS_CONFIG: Specifies how to convert message values from bytes into Java objects (e.g.,StringDeserializerfor text values).
The Consumer Loop: Polling for Records
Consumers continuously poll Kafka for new data. This is typically done in an infinite loop. The poll(Duration timeout) method returns a batch of ConsumerRecords or an empty list if no new records are available within the specified timeout.
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Properties;
import java.util.Collections;
public class PollingConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-polling-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("my-topic"));
try {
while (true) { // Infinite loop to continuously poll
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) {
System.out.println("No records received. Waiting...");
} else {
System.out.println("Received " + records.count() + " records.");
// Processing of records would happen here
}
}
} finally {
consumer.close();
System.out.println("Consumer closed.");
}
}
}Processing Received Messages
Once you receive ConsumerRecords from poll(), you can iterate over them to process each message individually. Each ConsumerRecord contains the topic, partition, offset, key, and value of the message.
This is where your application's specific logic for handling the data lives.
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Properties;
import java.util.Collections;
public class ProcessingConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-processing-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("my-topic"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("Topic: %s, Partition: %d, Offset: %d, Key: %s, Value: %s%n",
record.topic(), record.partition(), record.offset(), record.key(), record.value());
// Add your custom processing logic here, e.g., store in database, perform analytics
}
}
} finally {
consumer.close();
System.out.println("Consumer closed.");
}
}
}Tracking Progress with Offsets
An offset is a unique, sequential ID assigned to each message within a Kafka topic partition. It acts like a pointer, indicating the position of a message.
- Consumers use offsets to track which messages they have already processed.
- Kafka stores the latest committed offset for each consumer group and partition.
- This ensures that if a consumer restarts or a partition is reassigned, it knows exactly where to resume reading to avoid reprocessing or missing messages.
Offsets are crucial for reliable message processing.
Automatic Offset Committing
By default, Kafka consumers are configured to auto-commit offsets periodically. This is convenient for simple applications but might lead to message loss or duplicates in certain failure scenarios.
- Set
enable.auto.committotrue(this is the default setting). - The
auto.commit.interval.msproperty determines how often offsets are committed (e.g., every 5 seconds). - This approach is generally suitable for 'at-least-once' processing where some duplicates are acceptable.
Manual Offset Control
For more precise control over message processing and stronger guarantees (like 'exactly-once' processing), you can manually commit offsets after your application has successfully processed a batch of messages.
- Set
enable.auto.committofalseto disable automatic commits. - Use
consumer.commitSync(): This method blocks until the commit operation is successful. - Use
consumer.commitAsync(): This method is non-blocking and faster, but requires a callback function for handling potential errors.
Manual commits ensure that records are processed *before* their offsets are marked as 'done', preventing data loss on consumer failure.
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Properties;
import java.util.Collections;
public class ManualCommitConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-manual-commit-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // Disable auto-commit
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("my-topic"));
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("Processed: Partition %d, Offset %d, Value: %s%n",
record.partition(), record.offset(), record.value());
// Simulate processing work that must complete before committing
}
if (!records.isEmpty()) {
consumer.commitSync(); // Commit offsets after processing the batch
System.out.println("Offsets committed manually.");
}
}
} finally {
consumer.close();
System.out.println("Consumer closed.");
}
}
}Check Your Understanding
Consider a Kafka topic with 3 partitions. A consumer group with 2 consumer instances is subscribed to this topic. If one consumer instance fails, what happens to the partitions it was processing, and how does Kafka ensure no data is lost?
Lesson Recap: Consuming Data
Great job! In this lesson, you learned how Kafka consumers read data:
- Consumers subscribe to topics and actively poll Kafka for new messages.
- Consumer groups enable scalable and fault-tolerant message processing across multiple consumer instances.
- Offsets are sequential IDs that track a consumer's progress within a topic partition.
- You can choose between automatic (simpler, default) or manual (more control, better guarantees) offset committing for reliability.
Next, we'll dive deeper into partitions and offsets to understand how Kafka achieves high throughput and message ordering.
자주 묻는 질문
“Kafka에서 메시지 소비” 강의는 무료인가요?
네 — “Kafka에서 메시지 소비” 전체 내용을 이 웹사이트에서 무료로 읽을 수 있습니다. 인터랙티브하게 실습하려면(내장 코드 에디터와 24/7 AI 튜터), CoddyKit PRO로 업그레이드하면 Apache Kafka & Stream Processing Fundamentals 강의 전체를 잠금 해제할 수 있습니다. Apache Kafka & Stream Processing Fundamentals 강의에는 총 4개의 강의가 포함되어 있습니다.
“Kafka에서 메시지 소비”에서 뭘 배우나요?
Kafka 토픽의 데이터 스트림을 읽고 처리하며 오프셋과 컨슈머 그룹을 다루는 컨슈머를 구축하는 방법을 배웁니다. 브라우저에서 직접 실행하는 실습 코드로 Apache Kafka & Stream Processing Fundamentals을(를) 배우며, 24/7 AI 튜터가 강의를 진행하면서 질문에 답변해줍니다.
Apache Kafka & Stream Processing Fundamentals을(를) 시작하는 데 경험이 필요한가요?
사전 경험은 필요하지 않습니다. CoddyKit의 Apache Kafka & Stream Processing Fundamentals은(는) 초급자부터 고급 학습자까지를 위해 구성되어 있으므로, 여기서 시작하거나 처음부터 시작할 수 있으며 자신의 속도대로 진행할 수 있습니다. 이것은 4개 중 2번째 강의입니다.
“Kafka에서 메시지 소비” 강의는 얼마나 걸리나요?
대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.
이 Apache Kafka & Stream Processing Fundamentals 강의에서 코드를 작성하고 실행할 수 있나요?
네. 모든 Apache Kafka & Stream Processing Fundamentals 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.
이 강의의 모든 강의
- Kafka로 메시지 생성
- Kafka에서 메시지 소비
- 파티션 및 오프셋 이해
- 메시지 키와 파티셔닝 전략