0Pricing
Advanced Spring Boot 4: Event-Driven Architecture (Kafka) · 课时

手动提交偏移量

实现对偏移量提交的手动控制,精确控制消息处理保障,避免数据丢失或重复

手动提交偏移量 是 CoddyKit 上的免费 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课时。 这是第 1 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

Why Manual Commits?

In Kafka, an offset marks the last message a consumer group has successfully processed from a topic partition. Committing an offset tells Kafka: "I've handled messages up to this point."

By default, Spring Kafka uses auto-commit, where offsets are committed periodically in the background. While convenient, this can sometimes lead to data loss or duplication if your application crashes mid-processing.

Manual offset committing gives you precise control, allowing you to decide exactly when an offset is marked as processed. This is crucial for ensuring message processing guarantees.

Auto-Commit: A Quick Look

With auto-commit, Kafka automatically commits offsets at a set interval (e.g., every 5 seconds). This means:

  • Messages are processed.
  • Offsets are committed later by Kafka.

If your application processes a message but crashes *before* Kafka's auto-commit interval passes, that message's offset might not be committed. When the application restarts, it will re-read and re-process that message, leading to potential duplicates (at-least-once processing).

Switching to Manual Mode

To take control of offset management, you need to disable auto-commit in your Spring Boot application's Kafka configuration. This is typically done by setting the AckMode.

The AckMode determines when a consumer acknowledges messages. For manual control, we'll use MANUAL_IMMEDIATE or MANUAL.

Here's how you might configure it in application.properties:

spring.kafka.consumer.enable-auto-commit=false
spring.kafka.listener.ack-mode=MANUAL_IMMEDIATE

The Acknowledgment Object

When ack-mode is set to a manual option, your @KafkaListener method can receive an additional parameter: the Acknowledgment object.

This object is your direct interface to signal to Kafka that you have successfully processed a message (or a batch of messages) and its offset can now be committed.

You'll call its acknowledge() method when you're ready.

Basic Manual Commit Example

Let's see a simple example where we manually commit the offset after processing each message. Notice the Acknowledgment acknowledgment parameter.

Run this code, then stop and restart. You'll see messages are not re-processed if committed.

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;

@SpringBootApplication
public class ManualCommitApp implements CommandLineRunner {

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public static void main(String[] args) {
        SpringApplication.run(ManualCommitApp.class, args);
    }

    @Override
    public void run(String... args) throws Exception {
        System.out.println("Sending message...");
        kafkaTemplate.send("my-topic", "Hello CoddyKit!");
        System.out.println("Message sent.");
    }

    @KafkaListener(topics = "my-topic", groupId = "manual-group")
    public void listen(String message, Acknowledgment acknowledgment) {
        System.out.println("Received: " + message);
        // Simulate processing
        try {
            Thread.sleep(500);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        System.out.println("Processed: " + message + ", committing offset.");
        acknowledgment.acknowledge(); // Manual commit
    }
}

Configuring for the Example

For the previous example to work, you'd need a src/main/resources/application.properties file with Kafka broker details and the manual ack mode:

  • spring.kafka.bootstrap-servers=localhost:9092 (or your Kafka broker)
  • spring.kafka.consumer.enable-auto-commit=false
  • spring.kafka.listener.ack-mode=MANUAL_IMMEDIATE
  • spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
  • spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
  • spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
  • spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer

Remember to have a local Kafka running!

When to Acknowledge?

The core principle for manual committing is: commit only after your business logic has successfully completed.

  • After each message: As shown in the previous example. Good for low throughput or critical messages.
  • After a batch: Process multiple messages, then commit once for the entire batch. This is more efficient for high throughput.
  • After external interactions: If you write data to a database, commit the offset *only after* the database transaction is successful.

Choosing the right strategy depends on your application's requirements for performance and data consistency.

Batch Processing & Manual Commit

When your listener consumes a batch of messages (e.g., List<String>), you should commit the offset only after *all* messages in that batch have been successfully processed. The Acknowledgment object still works for the entire batch.

This is often combined with AckMode.BATCH, though MANUAL_IMMEDIATE also works.

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.CommandLineRunner;
import java.util.List;

@SpringBootApplication
public class BatchManualCommitApp implements CommandLineRunner {

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    public static void main(String[] args) {
        SpringApplication.run(BatchManualCommitApp.class, args);
    }

    @Override
    public void run(String... args) throws Exception {
        System.out.println("Sending 3 messages...");
        kafkaTemplate.send("my-batch-topic", "Batch Msg 1");
        kafkaTemplate.send("my-batch-topic", "Batch Msg 2");
        kafkaTemplate.send("my-batch-topic", "Batch Msg 3");
        System.out.println("Messages sent.");
    }

    // Ensure spring.kafka.listener.ack-mode=MANUAL_IMMEDIATE in properties
    @KafkaListener(topics = "my-batch-topic", groupId = "batch-manual-group")
    public void listenBatch(List<String> messages, Acknowledgment acknowledgment) {
        System.out.println("Received batch of " + messages.size() + " messages.");
        for (String msg : messages) {
            System.out.println("  Processing: " + msg);
            // Simulate processing each message in the batch
            try {
                Thread.sleep(100);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
        System.out.println("Finished processing batch. Committing offset.");
        acknowledgment.acknowledge(); // Commit once for the entire batch
    }
}

Handling Errors & Reprocessing

What happens if an error occurs during message processing *before* acknowledgment.acknowledge() is called?

Since the offset was not committed, Kafka considers the message (or batch) as not processed. Upon restart or rebalance, the consumer will re-fetch and re-process those messages. This is the foundation of at-least-once processing semantics.

While this guarantees no data loss, it requires your message processing logic to be idempotent, meaning processing the same message multiple times has the same effect as processing it once.

Trade-offs & Best Practices

Manual offset committing offers control but comes with considerations:

  • Overhead: Committing too frequently can add network and Kafka broker overhead.
  • Reprocessing Scope: Committing too infrequently means more messages might be reprocessed if a failure occurs.
  • Idempotency: Always design your consumers to be idempotent when using manual commits to handle potential duplicates gracefully.
  • Error Handling: Combine manual commits with robust exception handling (e.g., retries, Dead Letter Topics) to manage failures effectively.

Quick Check: Manual Commits

You are using manual offset committing in your Spring Kafka consumer. If an error occurs while processing a message, and the acknowledgment.acknowledge() method is NOT called for that message, what will happen?

Recap: Manual Offset Control

We've explored manual offset committing, a powerful technique for precise control over message processing in Spring Kafka.

  • It disables auto-commit, giving you control.
  • You use the Acknowledgment object to explicitly commit offsets.
  • Committing should happen only after successful business logic execution.
  • It enables at-least-once processing, but requires idempotent consumers.

This fine-grained control is vital for building robust and reliable event-driven applications.

常见问题解答

「手动提交偏移量」课时是免费的吗?

是的 — 「手动提交偏移量」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程的其余内容,请升级到 CoddyKit PRO。 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程共包含 4 节课。

「手动提交偏移量」这节课中我会学到什么?

实现对偏移量提交的手动控制,精确控制消息处理保障,避免数据丢失或重复 你通过在浏览器中直接运行的动手代码来练习 Advanced Spring Boot 4: Event-Driven Architecture (Kafka),全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 1 节课,共 4 节。

「手动提交偏移量」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课中编写并运行代码吗?

能。每节 Advanced Spring Boot 4: Event-Driven Architecture (Kafka) 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 手动提交偏移量
  2. 暂停和恢复消费者
  3. 并发与线程管理
  4. 再平衡监听器与静态成员资格
← 返回 Advanced Spring Boot 4: Event-Driven Architecture (Kafka)