🤖
AI审核中

分区不再限制消费者数量:Kafka Share Groups如何把事件流变成任务队列

Java 22分钟 121浏览 0评论

长期以来,Kafka 最擅长的是事件流,而不是传统意义上的任务队列。

当业务需要订单异步处理、图片转码、邮件发送、Webhook 投递、AI 推理任务时,很多团队会发现一个问题:Kafka 明明拥有很高的吞吐量,但消费者扩容能力却始终受分区数量限制。一个只有 8 个分区的 Topic,即使启动 20 个消费者实例,也只有 8 个实例能够真正工作。

为了提升峰值处理能力,团队不得不提前创建大量分区,随之而来的却是更多文件句柄、更高的元数据管理成本以及更复杂的分区规划。

Kafka 4.0 首次以 Early Access 形式引入 KIP-932,Kafka 4.2 将 Share Groups 推向生产可用,Kafka 4.3 又增加了更多 Share Group 配置和协调器优化。本文以 Kafka 4.3.x 为基础,分析这套新消费模型究竟解决了什么问题。(Apache Kafka)

一、Kafka 过去为什么不适合做任务队列

传统 Kafka Consumer Group 的核心规则是:

一个分区在同一个消费者组中,同一时间只能由一个消费者负责。

假设一个 Topic 有 3 个分区:

flowchart LR
    P0["分区 0"] --> C1["消费者 A"]
    P1["分区 1"] --> C2["消费者 B"]
    P2["分区 2"] --> C3["消费者 C"]
    C4["消费者 D"] --> I["没有分区可分配"]

当消费者 D 加入后,由于已经没有空闲分区,它只能保持空闲。

这套模型对日志采集、数据库变更订阅、事件驱动系统非常合理,因为它能够保证同一分区内的记录由一个消费者按顺序处理。

但对于任务队列,业务更关心的是:

  • 一条任务只交给一个 Worker;
  • Worker 数量可以随流量快速增加;
  • 每条任务可以单独确认;
  • 临时失败的任务可以重新投递;
  • 永久失败的任务可以停止重试;
  • 某个 Worker 崩溃后,任务能够自动交给其他 Worker。

传统消费者组虽然能够通过手动提交 Offset、Retry Topic、死信 Topic 等方式实现类似效果,但整个重试体系基本都需要业务自己搭建。

更关键的是,消费者数量始终无法突破分区数量。很多团队只能通过“过度分区”为未来峰值预留并行度。KIP-932 正是为了解除消费者数量与分区数量之间的强绑定。(Apache 维基)

二、Share Groups 改变的不是 Topic,而是消费协议

Share Groups 并没有为 Kafka 新增一种名为 Queue 的存储资源。

生产者依旧向普通 Kafka Topic 写入消息,消息依旧保存在分区日志中。真正发生变化的是消费者如何获取和确认记录。

传统消费者组分配的是整个分区,而 Share Group 分配的重点变成了分区中的具体记录。

flowchart LR
    P["同一个 Kafka 分区"] --> L["Broker 按记录获取锁"]
    L --> R1["记录 1 交给 Worker A"]
    L --> R2["记录 2 交给 Worker B"]
    L --> R3["记录 3 交给 Worker C"]

同一个分区可以同时分配给多个 Share Consumer,但同一条已经被获取的记录,在锁有效期间不会再交给同一 Share Group 中的其他消费者。

因此,即使 Topic 只有一个分区,也可以启动多个消费者并发处理不同记录。

Kafka 官方将 Share Group 定义为一种新的 Group 类型,与传统的 classicconsumer Group 并列。多个 Share Group 可以独立订阅同一个 Topic,每个 Share Group 都维护自己的处理进度和记录状态。(Apache 维基)

Consumer Group 与 Share Group 的核心区别

对比维度 Consumer Group Share Group
最小分配单位 分区 记录
分区归属 一个分区独占分配给一个消费者 一个分区可以分配给多个消费者
并行度上限 通常受分区数量限制 消费者数量可以超过分区数量
进度模型 提交 Offset 维护记录级状态和滑动窗口
失败处理 通常依赖业务重试机制 支持 RELEASE、REJECT 等确认结果
顺序保证 同一分区内顺序较强 跨批次可能乱序
典型场景 事件流、CDC、顺序处理 独立任务、弹性 Worker、异步作业

Share Groups 不是对 Consumer Group 的升级替换,两者解决的是不同问题。

需要严格分区顺序时,Consumer Group 仍然更合适;需要动态增加 Worker、逐条确认和失败重投时,Share Group 更自然。

三、一条记录会经历怎样的状态变化

Share Group 为每条正在处理的记录维护状态。

一条记录主要会经历四种状态:

flowchart LR
    A["Available 可获取"] --> B["Acquired 已加锁"]
    B -->|ACCEPT| C["Acknowledged 已确认"]
    B -->|RELEASE或锁超时| A
    B -->|REJECT| D["Archived 已归档"]
    B -->|达到投递上限| D
    C -->|窗口前移| D

Available:可以被消费者获取

记录尚未被任何消费者处理,可以分配给 Share Group 中的任意消费者。

Acquired:已经被某个消费者锁定

消费者获取记录时,Broker 会为记录创建一个有时限的获取锁。

Kafka 4.3 默认的记录锁时间是 30 秒。在锁有效期间,这条记录不会再交给同一个 Share Group 中的其他消费者。锁时间可以通过 Group 配置 share.record.lock.duration.ms 调整。(Apache Kafka)

Acknowledged:已经成功处理

消费者返回 ACCEPT 后,记录进入已确认状态,不再参与后续投递。

Archived:不再投递

记录可能因为以下原因进入 Archived 状态:

  • 消费者明确返回 REJECT
  • 记录达到最大投递次数;
  • Share Partition 的处理窗口已经越过该记录;
  • Topic 保留策略删除了对应日志。

需要注意,Archived 只是该 Share Group 对记录的消费状态,并不代表 Kafka 立即从日志中物理删除这条消息。

Kafka Topic 仍然按照自己的时间、大小或压缩策略管理日志。(Apache 维基)

四、四种确认结果决定任务的命运

KafkaShareConsumer 支持四种 AcknowledgeType

确认类型 含义 后续行为
ACCEPT 处理成功 不再投递
RELEASE 本次失败,但可以重试 重新变为可获取状态
REJECT 永久失败,不应继续重试 进入 Archived 状态
RENEW 任务仍在处理中 延长当前获取锁

官方 Java API 将 RELEASE 定义为允许再次投递,将 REJECT 定义为不再参与后续投递,而 RENEW 用于处理时间超过锁期限的长任务。(Apache Kafka)

ACCEPT:任务真正完成

例如订单同步成功、邮件发送成功或者图片转码完成,可以返回:

consumer.acknowledge(record, AcknowledgeType.ACCEPT);

RELEASE:临时错误,稍后重试

数据库短暂不可用、第三方接口超时、网络抖动等问题,通常适合返回:

consumer.acknowledge(record, AcknowledgeType.RELEASE);

记录会重新变为可获取状态,之后可能被当前消费者重新获取,也可能交给另一个消费者。

REJECT:永久错误,停止重试

参数格式错误、业务对象不存在、签名校验失败等无法通过重试解决的问题,可以返回:

consumer.acknowledge(record, AcknowledgeType.REJECT);

但生产系统不应直接丢弃失败原因。更合理的做法是先将原始消息、异常原因和处理上下文写入失败表或死信 Topic,再执行 REJECT

RENEW:任务还没完成,延长锁

如果视频转码、模型推理或者大文件处理需要几分钟,而记录锁只有 30 秒,就需要周期性续锁:

consumer.acknowledge(record, AcknowledgeType.RENEW);

RENEW 不是最终确认。任务完成后仍然需要返回 ACCEPTRELEASEREJECT

五、Kafka 如何避免毒消息无限重试

传统 Kafka Consumer 遇到一条始终处理失败的消息时,如果没有完善的重试与死信机制,很容易形成无限消费、无限报错。

Share Group 在 Broker 端维护记录的投递次数。

Kafka 4.3 中,share.delivery.count.limit 默认值为 5。当记录反复被获取、释放或者因为锁超时重新投递,并达到投递上限后,它会进入 Archived 状态,不再继续投递。(Apache Kafka)

不过,这个投递次数并不是严格的业务计数。

Kafka 官方说明,相关状态更新并不具备精确一次语义,因此投递次数主要用于防止毒消息无限循环,不能当作准确的业务重试次数。Share Group 整体仍然提供至少一次投递语义。(Apache 维基)

因此,业务代码仍然需要保证幂等。

例如一个扣减库存任务重复执行时,不能仅依赖消费者“理论上只执行一次”,而应通过业务流水号、唯一索引或状态机阻止重复扣减。

六、用 Java 编写一个 Share Consumer

下面使用 Kafka 4.3.1 客户端演示逐条确认。4.3.1 是 Kafka 4.3 系列的修复版本。(Apache Kafka)

Maven 依赖

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>4.3.1</version>
</dependency>

完整消费者示例

package com.example.kafka;

import org.apache.kafka.clients.consumer.AcknowledgeType;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaShareConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

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

public final class OrderShareWorker {

    private static final Duration POLL_TIMEOUT = Duration.ofSeconds(1);

    private OrderShareWorker() {
    }

    public static void main(String[] args) {
        Properties properties = new Properties();

        properties.put("bootstrap.servers", "localhost:9092");
        properties.put("group.id", "sg-order-workers-v1");
        properties.put(
                "key.deserializer",
                StringDeserializer.class.getName()
        );
        properties.put(
                "value.deserializer",
                StringDeserializer.class.getName()
        );

        // 使用逐条显式确认
        properties.put("share.acknowledgement.mode", "explicit");

        // 严格限制单次 poll 获取的记录数量
        properties.put("share.acquire.mode", "record_limit");
        properties.put("max.poll.records", "20");

        try (KafkaShareConsumer<String, String> consumer =
                     new KafkaShareConsumer<>(properties)) {

            consumer.subscribe(List.of("order-tasks"));

            while (!Thread.currentThread().isInterrupted()) {
                ConsumerRecords<String, String> records =
                        consumer.poll(POLL_TIMEOUT);

                for (ConsumerRecord<String, String> record : records) {
                    AcknowledgeType acknowledgeType = handle(record);
                    consumer.acknowledge(record, acknowledgeType);
                }

                if (!records.isEmpty()) {
                    var commitResult = consumer.commitSync();

                    commitResult.forEach((partition, error) ->
                            error.ifPresent(exception ->
                                    System.err.printf(
                                            "确认提交失败,partition=%s,error=%s%n",
                                            partition,
                                            exception.getMessage()
                                    )
                            )
                    );
                }
            }
        }
    }

    private static AcknowledgeType handle(
            ConsumerRecord<String, String> record
    ) {
        try {
            process(record.value());

            System.out.printf(
                    "处理成功,topic=%s,partition=%d,offset=%d%n",
                    record.topic(),
                    record.partition(),
                    record.offset()
            );

            return AcknowledgeType.ACCEPT;
        } catch (RetryableTaskException exception) {
            System.err.printf(
                    "临时失败,等待重试,offset=%d,error=%s%n",
                    record.offset(),
                    exception.getMessage()
            );

            return AcknowledgeType.RELEASE;
        } catch (IllegalArgumentException exception) {
            System.err.printf(
                    "永久失败,停止重试,offset=%d,error=%s%n",
                    record.offset(),
                    exception.getMessage()
            );

            // 实际项目中应先保存失败上下文或写入死信 Topic
            return AcknowledgeType.REJECT;
        } catch (Exception exception) {
            System.err.printf(
                    "未知异常,暂时按可重试处理,offset=%d,error=%s%n",
                    record.offset(),
                    exception.getMessage()
            );

            return AcknowledgeType.RELEASE;
        }
    }

    private static void process(String payload) {
        if (payload == null || payload.isBlank()) {
            throw new IllegalArgumentException("任务内容为空");
        }

        if (payload.startsWith("TEMP_FAIL")) {
            throw new RetryableTaskException("下游服务暂时不可用");
        }

        System.out.println("正在处理任务:" + payload);
    }

    private static final class RetryableTaskException
            extends RuntimeException {

        private RetryableTaskException(String message) {
            super(message);
        }
    }
}

在显式确认模式下,acknowledge() 首先只是更新消费者本地的确认状态。后续执行 commitSync()commitAsync() 或下一次 poll() 时,确认结果才会提交给 Kafka。

同时,显式模式要求上一次 poll() 返回的记录全部完成确认,否则直接进行下一次 poll() 会抛出 IllegalStateException。(Apache Kafka)

七、两个容易被忽略的客户端模式

1. implicit 与 explicit

通过 share.acknowledgement.mode 可以选择确认模式。

模式 特点 适用场景
implicit 下一次 poll 或 commit 时,将上一批记录视为成功 整批处理、失败逻辑简单
explicit 每条记录必须明确返回确认类型 独立任务、逐条重试、错误分类

默认值是 implicit

隐式模式代码简单,但只要开始下一次 poll(),上一批记录就可能被视为处理成功。因此,只要业务中存在“部分成功、部分失败”,就应优先使用 explicit。(Apache Kafka)

2. batch_optimized 与 record_limit

通过 share.acquire.mode 可以控制记录获取方式。

batch_optimized 是默认模式。它会尽量按照 Kafka 原始 Record Batch 边界获取数据,吞吐量更高,但一次 poll() 返回的记录数量可能超过 max.poll.records

record_limit 会严格限制记录数量,并关闭预获取。它的吞吐量可能稍低,但能更准确地控制并发任务数和锁开始计时的时机。(Apache Kafka)

对于单条处理耗时较长的任务,建议使用:

share.acquire.mode=record_limit
max.poll.records=20

对于处理逻辑非常轻、追求批量吞吐的任务,可以保留:

share.acquire.mode=batch_optimized

八、长任务不能只靠调大锁时间

假设一个 AI 视频生成任务平均需要 2 分钟,而获取锁只有 30 秒。

如果消费者在 30 秒内没有返回确认,记录锁会自动释放,任务可能被另一个消费者重新获取。此时两个 Worker 可能同时处理同一个业务任务。

最简单的方案是调大 share.record.lock.duration.ms,但锁时间并非越长越好。

锁设置为 10 分钟后,如果 Worker 在处理开始时直接宕机,这条记录可能需要较长时间才能重新投递,故障恢复速度也会降低。

更合理的方式是:

  1. 设置能够覆盖大部分普通任务的锁时间;
  2. 对超过锁时间的长任务周期性返回 RENEW
  3. 任务完成后再返回最终确认结果;
  4. 无论如何都保证业务幂等。

KafkaShareConsumer 本身不是线程安全的。对于异步 Worker Pool,不能让多个工作线程直接并发操作同一个 Consumer。官方建议所有 poll()acknowledge() 和提交操作由受控线程完成。(Apache Kafka)

一种更安全的线程模型如下:

flowchart LR
    P["Poll 与 ACK 线程"] --> Q["待处理任务队列"]
    Q --> W1["工作线程 A"]
    Q --> W2["工作线程 B"]
    Q --> W3["工作线程 C"]
    W1 --> R["处理结果队列"]
    W2 --> R
    W3 --> R
    R --> P

Poll 线程负责:

  • 调用 poll()
  • 将任务提交给线程池;
  • 接收处理结果;
  • 执行 ACCEPTRELEASEREJECTRENEW
  • 提交确认状态。

工作线程只执行业务逻辑,不直接操作 Kafka Consumer。

九、Share Groups 不等于 RabbitMQ

Share Groups 让 Kafka 更适合队列型工作负载,但它并没有把 Kafka 变成传统消息队列。

1. 确认成功不代表消息被删除

RabbitMQ 中,消息确认后通常会从队列中移除。

Kafka 中,消息仍然保存在 Topic 日志中,只是当前 Share Group 将其标记为已确认。另一个 Consumer Group 或 Share Group 仍然可以独立读取相同消息。

这意味着 Kafka 依旧保留了事件回放、多组独立消费和按时间恢复的能力。

2. Topic 保留策略仍然生效

所有 Share Group 共享同一份 Topic 日志。

如果 Topic 使用按大小保留,并且磁盘达到上限,Kafka 可能删除尚未被某个 Share Group 处理完成的旧日志。该 Share Group 的起始位置会随日志起始 Offset 向前移动,相应记录也不再可用。(Apache 维基)

因此,任务队列场景不能只关注消费速度,还要确保:

Topic 保留时间大于系统能够容忍的最长积压时间。

例如业务最多允许停机 3 天,就不应只保留 24 小时消息。

3. 顺序保证会明显减弱

Share Group 允许多个消费者并发处理同一分区,因此不能再依赖传统消费者组的严格分区顺序。

Kafka 只保证同一次返回的批次中,同一分区的记录按 Offset 递增。不同批次之间,特别是出现锁超时和重新投递后,Offset 可能倒退。

例如:

  1. Worker A 获取 Offset 100~109 后崩溃;
  2. Worker B 成功处理 Offset 110~119;
  3. 100~109 的锁超时;
  4. Worker B 再次获取 100~109。

最终业务看到的完成顺序可能是先处理 110~119,再处理 100~109。(Apache 维基)

4. Archived 不是自动死信转发

达到投递上限后,记录会进入 Archived 状态,但 Kafka 不会自动把它复制到业务定义的死信 Topic。

如果系统需要人工补偿、失败查询和重新驱动,就应在执行 REJECT 前主动记录失败信息。

5. 不适合复杂消息路由场景

Share Groups 的重点是记录共享、获取锁和逐条确认。

当业务强依赖交换机路由、消息优先级、原生延迟队列、灵活的死信转发时,RabbitMQ 等传统消息队列仍然更直接。

十、生产落地必须解决的五个问题

1. 用业务幂等对抗至少一次投递

Share Groups 是至少一次投递模型。

以下情况都可能导致重复执行:

  • Worker 处理成功,但确认提交失败;
  • Worker 处理成功后进程崩溃;
  • 获取锁在处理完成前超时;
  • Broker 或网络发生切换;
  • 消费者执行 RELEASE

常用幂等方案包括:

  • 使用任务 ID 建立唯一索引;
  • 写入业务处理流水表;
  • 通过状态机限制重复状态变更;
  • 调用第三方接口时传递幂等键;
  • 使用 Inbox 或 Outbox 模式管理跨系统副作用。

不要把 Kafka 的确认结果当成业务事务提交。

2. 明确失败分类

异常不能全部 RELEASE,否则格式错误的毒消息会反复占用资源。

也不能全部 REJECT,否则一次网络抖动就可能造成永久任务丢失。

推荐建立统一分类:

异常类型 处理方式
业务执行成功 ACCEPT
网络超时、限流、临时不可用 RELEASE
数据格式错误、对象不存在 记录失败后 REJECT
仍在执行的长任务 RENEW
无法判断的未知异常 通常先 RELEASE,同时告警

3. 不要直接照搬默认配置

Kafka 4.3 的部分 Share Group 默认配置如下:(Apache Kafka)

配置 默认值 作用
share.record.lock.duration.ms 30000 获取锁持续时间
share.delivery.count.limit 5 最大投递次数
share.partition.max.record.locks 2000 单个 Share Partition 最大记录锁数量
share.auto.offset.reset latest 新 Share Group 初始位置
share.isolation.level read_uncommitted 事务消息读取级别
share.renew.acknowledge.enable true 是否允许 RENEW

锁时间应参考任务处理耗时的 P95 或 P99,而不是只看平均值。

投递次数应结合第三方故障恢复时间、重试成本以及任务价值决定。

4. 为 Archived 记录建立业务补偿通道

Broker 的 Archived 状态主要解决无限重试问题,并不能代替完整的失败治理系统。

生产系统至少应保存:

  • Task ID;
  • 原始消息;
  • Topic、Partition 和 Offset;
  • 失败类型;
  • 异常堆栈摘要;
  • 首次失败时间;
  • 最后失败时间;
  • 人工处理状态;
  • 是否允许重新投递。

这样才能实现失败任务查询、人工修复和重新驱动。

5. 同时监控积压与处理质量

Kafka 4.2 为 Share Groups 增加了 Share Partition Lag 持久化和相关查询能力,用于观察消费进度和后续自动扩缩容。(Apache Kafka)

除了 Lag,还应监控:

  • 单位时间获取记录数;
  • ACCEPT 数量;
  • RELEASE 数量;
  • REJECT 数量;
  • 锁超时数量;
  • 确认提交失败数量;
  • 单任务处理耗时;
  • Worker 活跃数量;
  • 下游接口错误率;
  • Archived 任务增长速度。

如果 Lag 不高但 RELEASE 比例持续升高,系统可能已经进入“不断取出、不断失败”的无效循环。

十一、从传统 Consumer Group 迁移到 Share Group

Share Group 使用新的消费者类型和状态模型,不能简单地把原来的 KafkaConsumer 配置改一下就直接替换。

更稳妥的迁移过程如下。

第一步:创建新的 Group ID

Consumer Group 和 Share Group 位于同一个 Group ID 命名空间中。

如果某个 Group ID 已经被创建为传统 Consumer Group,就不能再用同一个 ID 创建 Share Group,反之亦然。官方建议提前制定命名规范。(Apache 维基)

可以使用:

  • cg-order-events-v1
  • sg-order-workers-v1

分别标识普通消费者组和共享消费者组。

第二步:确认任务不依赖严格顺序

检查业务是否依赖:

  • 同一用户操作必须串行;
  • 同一订单状态必须按 Offset 顺序变化;
  • 后一条消息依赖前一条消息的处理结果;
  • 通过消息 Key 保证实体顺序。

只要存在这些要求,就不能直接使用 Share Group 并发处理同一分区。

第三步:保证所有副作用幂等

在真正切换前,主动制造以下故障:

  • 业务处理成功后终止消费者;
  • 处理过程中断网;
  • 处理时间超过记录锁;
  • 确认提交失败;
  • Broker Leader 切换。

观察重复投递是否会导致重复扣款、重复发券、重复发送通知或状态回退。

第四步:初始化 Share Group 位置

新 Share Group 默认从 latest 开始。

需要读取历史任务时,可以在 Share Group 没有活跃成员的情况下执行 Offset 重置:

bin/kafka-share-groups.sh \
  --bootstrap-server localhost:9092 \
  --group sg-order-workers-v1 \
  --topic order-tasks \
  --reset-offsets \
  --to-earliest \
  --execute

Share Group 的起始位置只能在组为空、没有活跃成员时安全重置。重置操作会丢弃现有的在途状态和投递计数。(Apache 维基)

第五步:灰度运行

可以为同一个 Topic 创建新的 Share Group,让它以影子模式处理消息。

灰度期间重点比较:

  • 传统消费者与 Share Consumer 的处理结果;
  • 重复执行比例;
  • 顺序变化对业务的影响;
  • 高峰期扩容速度;
  • RELEASE 和 REJECT 分布;
  • Worker 宕机后的恢复时间。

确认无误后,再逐步切换正式流量。

十二、哪些场景最适合 Share Groups

场景 是否推荐 原因
图片压缩、视频转码 推荐 任务独立,可水平增加 Worker
邮件、短信异步发送 推荐 适合逐条确认和失败重试
Webhook 投递 推荐 可区分临时失败与永久失败
AI 推理和文档解析 推荐 处理时间差异大,可使用 RENEW
订单创建事件广播 谨慎 更像事件流,普通 Consumer Group 更自然
同一账户顺序记账 不推荐 不能依赖严格分区顺序
数据库 CDC 同步 通常不推荐 强依赖 Offset 和顺序推进
复杂优先级与延迟队列 谨慎 专业消息队列能力更完整
需要随流量快速增加 Worker 推荐 消费者数量可超过分区数量

一个简单的判断标准是:

消息代表“发生过什么”,优先考虑 Consumer Group;消息代表“需要完成什么”,可以重点评估 Share Group。

十三、总结

Share Groups 最重要的变化,不只是让 Kafka 多了几个确认 API,而是改变了 Kafka 长期以来的分区独占消费模型。

过去,分区既是存储并行度,也是消费并行度。想增加消费者,就必须提前增加分区。

现在,Share Group 将并行处理进一步下沉到记录级别。Broker 使用获取锁避免同一记录被同时交给多个消费者,并通过 ACCEPTRELEASEREJECTRENEW 表达每条任务的处理结果。

它让 Kafka 能够更自然地承载任务队列、弹性 Worker 和逐条重试场景,同时继续保留日志存储、多组独立消费和历史回放能力。

但代价同样明确:跨批次顺序不再可靠、系统仍是至少一次投递、确认不会立即删除消息、Archived 也不等于完整的死信治理。

因此,Share Groups 并不是“Kafka 终于取代所有消息队列”,而是让 Kafka 在事件流之外,获得了一套更符合任务处理场景的消费协议。

真正决定它能否进入生产环境的,不是能不能启动 KafkaShareConsumer,而是业务是否已经准备好处理幂等、重试、锁超时、失败归档、日志保留和顺序变化。

0 条评论
如果你觉得文章对你有帮助,那就请作者喝杯咖啡吧☕
微信
支付宝
  0 条评论