Kafka 是分布式系统的「中央神经系统」——订单、日志、事件溯源、CDC 全部汇入其中。Scala 开发者面对它时有两个成熟选择:FS2-Kafka(基于 fs2 + Cats Effect 的纯函数式客户端)与 Alpakka Kafka(Akka Streams 的 Reactive Streams 连接器)。两者都建立在同一个官方 Java 客户端之上,区别在于抽象层:前者把消费与生产建模成 Stream[F, *],天然享受 fiber 并发、资源管理与 effect 组合;后者沿用 Actor 模型的 Materializer 与图 DSL。本文不止讲 API,更聚焦那些决定系统是否可运维的细节:commit 时机、exactly-once 的代价、rebalance 期间的重复消费、lag 监控与优雅停机。
目录
- 1. Kafka 核心概念回顾
- 2. FS2-Kafka 入门
- 3. 消费语义与交付保证
- 4. 流式管道设计
- 5. Alpakka Kafka 对比
- 6. 序列化与 Schema Registry
- 7. 状态与有状态处理
- 8. 可观测性
- 9. 生产实践
- 10. 速查表与一句话记忆
- 延伸阅读
1. Kafka 核心概念回顾
Kafka 的可靠性来自一组简单但互相牵制的抽象。Topic 是逻辑日志,Partition 是它的物理分片与并行单元,Offset 是分区内单调递增的位置。Consumer Group 把多个消费者实例组织成一个逻辑订阅者,每个分区在同一时刻只被组内一个实例消费——这既是水平扩展的来源,也是顺序保证的边界:Kafka 只保证分区内有序,跨分区无序。Rebalance 是组成员变化时重新分配分区的过程,它会暂停消费(stop-the-world),是延迟抖动与重复消费的主要来源之一。
- topic:逻辑主题,可配置 partitions / replication.factor / retention.ms
- partition:有序不可变日志,offset 单调递增,是并行与顺序的最小单元
- offset:消费者位点,__consumer_offsets 内部主题持久化提交结果
- consumer group:组内分区独占,实例数 > 分区数时多出的实例空转
- rebalance:分区再分配;触发条件=成员加入/离开、订阅变化、分区数变化
- 分区键:相同 key 必落同一分区,是「保序」的唯一手段
- ISR:in-sync replicas,acks=all 时需 ISR 全员确认才算写入成功
import org.apache.kafka.clients.admin.{AdminClient, NewTopic}
import org.apache.kafka.common.config.TopicConfig
import java.util.Properties
val props = new Properties()
props.put("bootstrap.servers", "localhost:9092")
val admin = AdminClient.create(props)
val topic = new NewTopic("orders", 12, 3.toShort)
.configs(Map(
TopicConfig.RETENTION_MS_CONFIG -> "604800000",
TopicConfig.CLEANUP_POLICY_CONFIG -> "delete",
TopicConfig.MIN_INSUFFICIENT_REPLICAS_CONFIG -> "2"
).asJava)
admin.createTopics(List(topic).asJava).all().get()
工程要点:分区数一旦确定只能增不能减,且增加分区会破坏按 key 的既有分区映射(同 key 的旧数据在旧分区、新数据可能落新分区)——先估算峰值吞吐(单分区消费约几 MB/s 量级)再定分区数。acks=all + min.insync.replicas=2 是生产底线,宁可写失败也不要静默丢数据。
2. FS2-Kafka 入门
FS2-Kafka 把 Kafka 客户端包装成 Cats Effect 的 Resource:KafkaConsumer.resource 负责连接的获取与关闭,KafkaConsumer.stream 返回一个 Stream[F, CommittableConsumerRecord],整个消费循环就是一个可以 map / filter / parEvalMap 的纯流。生产侧用 KafkaProducer.pipe(settings),它返回一个 Pipe[F, ProducerRecord, ProducerResult],直接嵌进流管道即可。
- 依赖:org.typelevel %% fs2-kafka % 3.x(对应 Kafka 3.x 客户端)
- KafkaConsumerSettings:bootstrapServers + groupId + 反序列化器
- ConsumerSettings[F, K, V].withAutoOffsetReset(First/Last) 决定无位点时的起点
- KafkaConsumer.stream(settings).subscribe(topic):返回 CommittableConsumerRecord 流
- 提交:stream.commitBatchWithin(maxBatch, maxInterval) 或手动 commit
- 生产:KafkaProducer.pipe(producerSettings) 得到 Pipe
import cats.effect.{IO, IOApp}
import fs2.kafka.*
import scala.concurrent.duration.*
object ConsumerApp extends IOApp.Simple {
val consumerSettings: ConsumerSettings[IO, String, String] =
ConsumerSettings[IO, String, String]
.withAutoOffsetReset(AutoOffsetReset.Earliest)
.withBootstrapServers("localhost:9092")
.withGroupId("order-processor")
.withEnableAutoCommit(false) // 关闭自动提交,手动控制
.withMaxPollRecords(500)
val producerSettings: ProducerSettings[IO, String, String] =
ProducerSettings[IO, String, String]
.withBootstrapServers("localhost:9092")
.withEnableIdempotence(true)
val run: IO[Unit] =
KafkaConsumer
.stream(consumerSettings)
.subscribeTo("orders")
.records
.map(_.value.toUpperCase)
.through(KafkaProducer.pipe(producerSettings))
.compile
.drain
}
工程要点:永远显式关闭 enable.auto.commit。自动提交在后台按时间间隔提交「已 poll 到的最大 offset」,而不是「已处理完的 offset」——崩溃时必然丢消息。FS2-Kafka 的 commitBatchWithin(n, interval) 才是正确的批量提交抽象。
3. 消费语义与交付保证
交付保证有三个层次。At-most-once:先提交再处理,崩溃丢消息。At-least-once:先处理再提交,崩溃重放,需要下游幂等。Exactly-once:通过 Kafka 事务把「消费位点提交」与「生产输出」绑定为一个原子操作,代价是事务开销、transactional.id 管理与对下游的强约束。绝大多数业务用 at-least-once + 幂等下游,比强行上 exactly-once 更简单可靠。
- at-least-once:处理完成后 commit;崩溃重放 → 下游必须幂等
- at-most-once:commit 后处理;丢数据,仅用于可丢弃的监控类场景
- exactly-once(EOS):事务性生产者,消费位点写入同一事务
- 幂等生产者:enable.idempotence=true 保证单分区内不重复(重试不产生副本)
- 事务前提:transactional.id 唯一且稳定,max.in.flight <= 5,acks=all
- 跨系统 EOS:Kafka 事务无法覆盖外部 DB → 用 outbox 或幂等键
- 重复来源:rebalance 未提交位点、poll 超时被踢出、处理超 max.poll.interval.ms
import fs2.kafka.*
val transactional: ProducerSettings[IO, String, String] =
ProducerSettings[IO, String, String]
.withBootstrapServers("localhost:9092")
.withTransactionalId("order-tx-1") // 必须全局唯一且稳定
.withEnableIdempotence(true)
KafkaProducer
.transactional(transactional) // 返回 TransactionalKafkaProducer
.use { producer =>
producer.produce(
ProducerRecords.one(ProducerRecord("orders-out", "k", "v"))
)
}
工程要点:max.poll.interval.ms(默认 5 分钟)是隐藏杀手——单批处理耗时超过它,消费者会被踢出组触发 rebalance,位点回退导致整批重复。要么降低 max.poll.records,要么把耗时处理移出 poll 线程。事务性生产者要求消费者用 readCommitted 隔离级别,否则会读到未提交数据。
4. 流式管道设计
真实管道的骨架是:消费 → 解析 → 校验 → 富化(调用外部服务)→ 转换 → 生产/落库。fs2 的组合子让这条链路可读且可测:evalMap 串行处理、parEvalMap(n) 限并发并行、chunkN 批量聚合、handleErrorWith 兜底、through 复用一个 Pipe。背压由 fs2 的 pull 模型天然提供——下游不拉取,上游就停,不需要任何缓冲调参。
- 串行:evalMap(一次一个 Effect)
- 并行:parEvalMap(n) 保持顺序、parEvalMapUnordered(n) 吞吐优先
- 批量:chunkN(n) / grouped(n) 后 evalMap 一次性写库
- 重试:retry(_.retry) 固定、.delay + retryWithDelay 指数退避
- 超时:.timeout(d) / .timeoutTo(d, fallback)
- 兜底:handleErrorWith 转死信,不要吞异常
- 死信队列:解析/处理失败的消息 produce 到 <topic>.DLQ 并附原始头
- 背压:fs2 天然 pull-based,无需 buffer 调参;buffer(n) 仅平滑抖动
import fs2.{Stream, Pipe}
import scala.concurrent.duration.*
final case class Order(id: String, amount: BigDecimal, raw: String)
def parse(r: CommittableConsumerRecord[IO, String, String]): IO[Either[String, Order]] =
IO(Order(r.record.key, BigDecimal(r.record.value), r.record.value)).attempt
.map(_.left.map(_.getMessage))
def pipeline: Pipe[IO, CommittableConsumerRecord[IO, String, String], Unit] =
_.parEvalMap(8) { rec => // 限并发 8,保序
parse(rec)
.timeout(3.seconds)
.retryWithDelay(3, 500.millis) // 指数退避重试
.flatMap {
case Right(o) => process(o)
case Left(err) => toDeadLetter(rec, err)
}
}
.chunkN(100)
.evalMap(batch => commitBatch(batch))
工程要点:parEvalMapUnordered 吞吐更高但打乱顺序,若下游依赖顺序必须用 parEvalMap。重试务必区分可重试(网络、超时)与不可重试(反序列化失败、业务校验失败)——后者重试一万次也不会成功,直接进死信队列。死信消息要带上原始 offset、异常栈与失败时间,否则事后无从排查。
5. Alpakka Kafka 对比
Alpakka Kafka 是 Akka Streams 的 Kafka 连接器,能力与 FS2-Kafka 高度重叠,差异在抽象风格。Consumer.plainSource 给一个普通消息流(自动提交)、Consumer.committableSource 给带提交能力的流、Producer.flexiFlow 支持在流中透传消息上下文。它更适合已经在 Akka 生态里的系统,或需要图 DSL 复杂拓扑(Fan-in/Fan-out、Partition)的场景。
| 维度 | FS2-Kafka | Alpakka Kafka |
|---|---|---|
| 抽象 | fs2 Stream + Cats Effect | Akka Streams 图 |
| 资源管理 | Resource / Stream | Materializer + KillSwitch |
| 背压 | pull 模型原生 | Reactive Streams 协议 |
| 提交控制 | commitBatchWithin / 手动 | CommittableOffset |
| 测试 | 纯流 compile 断言 | TestKit probe |
| 生态绑定 | Typelevel 栈 | Akka / Pekko 栈 |
import akka.actor.ActorSystem
import akka.kafka.scaladsl.{Committer, Consumer, Producer}
import akka.kafka.{ConsumerSettings, ProducerSettings, Subscriptions}
import akka.stream.scaladsl.Source
import org.apache.kafka.clients.consumer.ConsumerConfig
implicit val system: ActorSystem = ActorSystem("kafka")
val consumerSettings = ConsumerSettings(system, new StringDeserializer, new StringDeserializer)
.withBootstrapServers("localhost:9092")
.withGroupId("order-processor")
.withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
val producerSettings = ProducerSettings(system, new StringSerializer, new StringSerializer)
val done = Consumer
.committableSource(consumerSettings, Subscriptions.topics("orders"))
.mapAsync(8) { msg =>
process(msg.record.value()).map(_ => msg.committableOffset)
}
.groupedWithin(100, 5.seconds)
.mapAsync(1) { offsets =>
// 批量提交位点:把 committableOffset 集合交给 Committer
Source(offsets).runWith(Committer.sink(committerSettings))
}
.run()
工程要点:Alpakka 的 plainSource 默认自动提交,语义是 at-most-once,生产环境请一律用 committableSource + 显式 Committer。注意 Akka 的 BSL 许可变更——新项目建议评估 Pekko 分支(org.apache.pekko 命名空间)。
6. 序列化与 Schema Registry
Kafka 只存字节,序列化格式决定了跨服务契约的健壮性。JSON(circe) 可读、易调试,但无强制 schema,演进靠自律。Avro 体积小、支持 schema 演进与 Confluent Schema Registry 的兼容性检查,是数据管道主流。Protobuf 适合已有 gRPC 契约、需要跨语言强类型的场景。Key 的序列化同样重要——分区键必须与消费端的分区语义一致。
- JSON + circe:Decoder/Encoder 派生,灵活但无 schema 约束
- Avro:紧凑二进制,Schema Registry 强制兼容性(BACKWARD/FORWARD/FULL)
- Protobuf:强类型、跨语言,适合与 gRPC 共享 .proto
- Schema Registry:subject 命名 <topic>-key / <topic>-value
- 兼容策略:新增字段给默认值=向后兼容;删字段/改类型=破坏性
- 序列化异常处理:DeserializationException 不可重试,直接死信
- 错误处理:KafkaConsumerSettings 可配 withDeserializationExceptionHandler
import io.circe.{Decoder, Encoder}
import io.circe.generic.semiauto.*
import fs2.kafka.*
final case class OrderEvent(id: String, amount: BigDecimal, ts: Long)
object OrderEvent {
given Encoder[OrderEvent] = deriveEncoder
given Decoder[OrderEvent] = deriveDecoder
}
val serde: Serde[IO, OrderEvent] =
Serde.instance(
serializer = Serializer.lift[IO, OrderEvent](o => IO.pure(Encoder[OrderEvent](o).noSpaces.getBytes)),
deserializer = Deserializer.lift[IO, OrderEvent](b =>
IO.fromEither(io.circe.parser.decode[OrderEvent](new String(b)).left.map(e => new RuntimeException(e.getMessage))))
)
工程要点:Schema Registry 的兼容性检查是在注册时发生的,一旦允许破坏性变更,老消费者会直接崩。生产上把兼容级别设为 BACKWARD 并禁止手工绕过。反序列化失败不可重试,必须走死信并告警——它通常意味着上游发布了不兼容的 schema。
7. 状态与有状态处理
Kafka 的并行模型决定了状态必须按分区局部化:要聚合,就用分区键把同一实体的消息路由到同一分区,再用该分区内的 offset 顺序做累积。fs2 提供了 mapAccumulate 做无界累积,Stream.fixedRate 或 groupWithin 做窗口,状态可放在内存(须可重建)或外部存储(Redis / RocksDB)。真正的流处理引擎(Flink / Kafka Streams)提供的是状态容错与 changelog,纯客户端方案要自己承担。
- 分区键:key = entityId,保证同一实体消息进同一分区
- 累积:mapAccumulate 无界状态,须考虑内存增长
- 窗口:groupWithin(size, time) 做时间+数量双触发
- 状态存储:内存(可重建)或外部(Redis/RocksDB/Postgres)
- 状态容错:Kafka Streams 用 changelog + 本地 RocksDB 恢复
- 幂等:状态更新配合幂等键,避免重放导致重复累加
- 分区再平衡:状态须随分区迁移或可从源重建
import fs2.Stream
import scala.concurrent.duration.*
def windowedCounts[F[_]](
in: Stream[F, (String, BigDecimal)]
): Stream[F, Map[String, BigDecimal]] =
in.groupWithin(1000, 5.seconds) // 1000 条或 5 秒触发
.map { chunk =>
chunk.groupBy(_._1).view.mapValues(_.map(_._2).sum).toMap
}
工程要点:内存状态在 rebalance 时会随分区迁移而失效——如果状态无法从源重建,就必须落外部存储。别把「本地 Map 缓存」当成生产级状态存储:一次 rebalance 或重启,累计值就归零了。
8. 可观测性
Kafka 消费者最核心的指标是 consumer lag(分区最新 offset 减已提交 offset)。它直接反映「积压多久」,比 CPU/内存更早暴露问题。除此之外需要监控 rebalance 次数、poll 间隔、提交失败率、处理耗时分布。分布式追踪方面,Kafka 消息头(headers)是天然的 trace 载体——生产时注入 traceparent,消费时提取并续接 span。
- consumer lag:kafka_consumergroup_lag(JMX / kafka-exporter / Burrow)
- 分区级 lag:定位是全局慢还是个别分区倾斜
- rebalance 计数:kafka_consumer_rebalance_total 突增说明实例不稳定
- poll/commit:commit 失败率、max.poll.interval 触顶次数
- 处理耗时:p50/p95/p99 直方图,按 topic/partition 维度
- trace 传播:ProducerRecord headers 注入 traceparent,消费侧提取续接
- 告警:lag 持续增长 + rebalance 频繁 = 处理能力不足或处理逻辑变慢
import fs2.kafka.*
import org.apache.kafka.common.header.internals.RecordHeader
val traceId = "0123456789abcdef"
val record = ProducerRecord("orders", "key", "value")
.withHeaders(Headers(new RecordHeader("traceparent", s"00-$traceId-...-01".getBytes)))
// 消费侧提取
val span = record.headers("traceparent") match {
case h :: _ => Some(new String(h.value()))
case Nil => None
}
工程要点:lag 的绝对值不重要,趋势才重要——恒定 10 万 lag 可能只是吞吐匹配,持续攀升才是真积压。用 kafka-exporter 或 Burrow 做分区级 lag 采集,别自己写 offset 差值逻辑。trace 传播一定要在生产与消费两端都做,断链的 trace 等于没有。
9. 生产实践
从能跑到可运维,差距全在细节。优雅停机:收到 SIGTERM 后先停止拉取、处理完在途消息、提交位点、再关闭——否则每次发布都产生一批重复。Rebalance 监听:在分区被撤销前提交当前位点,减少重复。测试:Testcontainers 起真实 Kafka 或 embedded-kafka,别用 mock 掩盖协议行为。调优:fetch.min.bytes、max.poll.records、fetch.max.wait.ms 共同决定吞吐与延迟的权衡。
- 优雅停机:KillSwitch / 中断信号 → 停止拉取 → 处理在途 → 提交 → 关闭
- RebalanceListener:onPartitionsRevoked 前提交位点(同步提交)
- 消费者并发上限 = 分区数,实例数超过分区数纯属浪费
- 吞吐调优:fetch.min.bytes 增大提升吞吐、max.poll.records 控制单批
- 会话与心跳:session.timeout.ms、heartbeat.interval.ms 需协调
- 测试:Testcontainers KafkaContainer 或 embedded-kafka,覆盖 rebalance 场景
- 版本:客户端版本 >= broker 版本兼容,Kafka 3.x 起 KRaft 免 ZooKeeper
import fs2.Stream
import cats.effect.{IO, Deferred, Resource}
def graceful: Resource[IO, Unit] =
for {
stop <- Resource.eval(Deferred[IO, Unit])
_ <- Resource.make(
Stream.eval(stop.get).compile.drain.background
)(_ => stop.complete(()).void)
} yield ()
// 停机信号到达时:停止拉取 → 排空在途 → commit → close
// Kubernetes: terminationGracePeriodSeconds 必须大于最长批次处理时间
工程要点:Kubernetes 的默认 terminationGracePeriodSeconds=30s 常常短于一个批次处理时间,导致容器被 SIGKILL、位点未提交、消息重复。把它调到「最长批次处理时间 × 2」以上,并在应用内实现 SIGTERM 处理逻辑。
10. 速查表与一句话记忆
| 问题 | 一句话答案 |
|---|---|
| 顺序保证 | 仅分区内有序,同 key 落同分区 |
| 默认语义 | at-least-once,下游必须幂等 |
| 提交方式 | 关自动提交,处理完再 commitBatchWithin |
| exactly-once | 事务性生产者 + transactional.id + readCommitted |
| 并行处理 | parEvalMap 保序,parEvalMapUnordered 追吞吐 |
| 批量写库 | chunkN / grouped 后一次性提交 |
| 失败消息 | 不可重试的进死信队列,带原始头与异常 |
| lag 监控 | 看趋势不看绝对值,分区级采集 |
| 优雅停机 | 停拉取、排空在途、提交、关闭 |
一句话记忆:Kafka 集成 = 分区键定顺序 + at-least-once 配幂等下游 + 手动提交控位点 + parEvalMap 控并发 + 死信队列接住不可重试的错 + lag 趋势做告警 + SIGTERM 优雅停机——FS2-Kafka 用 Resource 管连接、Stream 管管道,Alpakka 用 Materializer 管图;选哪套看你的效果栈,选对语义才决定数据丢不丢。
延伸阅读
- fs2 与 Akka Streams 的流处理基础
- Alpakka Kafka 背后的图 DSL 与背压
- FS2-Kafka 依赖的并发与资源模型
- Avro / Protobuf / JSON 的选型与演进
- 指标、日志与 trace 传播体系
- Testcontainers 与流式代码测试
- Kafka 专题 — 分区、副本、事务与运维细节
- 分布式系统专题 — 消息队列在架构中的位置
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。