Skip to content

Kafka商业常用场景

Kafka 在商业系统里最常见的不是“博客实战”,而是承接大规模事件流。它通常位于业务系统、数据平台、搜索、风控、监控之间,作为事件管道和缓冲层。

场景一:订单事件总线

订单创建后,很多系统都关心这件事:

  • 库存系统扣减或锁定库存。
  • 搜索系统同步订单索引。
  • 风控系统判断异常交易。
  • 积分系统发放积分。
  • 数据平台沉淀交易数据。
mermaid
flowchart TD
    A["订单服务"] --> B["Topic: order.events"]
    B --> C["库存消费组"]
    B --> D["风控消费组"]
    B --> E["搜索消费组"]
    B --> F["数据仓库消费组"]

为什么适合 Kafka:

  • 多个系统可以使用不同 groupId 独立消费。
  • 新增下游系统时不用改订单服务。
  • 消息保留期内可以重新消费补数据。

注意点:

  • 事件命名用过去式,例如 ORDER_CREATEDORDER_PAID
  • 同一订单用 orderId 作为 key,保证订单维度顺序。
  • 消费端必须做幂等和状态机判断。

场景二:日志采集和可观测性

应用日志、访问日志、审计日志可以先写入 Kafka,再由不同系统消费。

mermaid
flowchart TD
    A["应用服务"] --> B["日志采集 Agent"]
    B --> C["Kafka: app.logs"]
    C --> D["Elasticsearch"]
    C --> E["对象存储"]
    C --> F["告警系统"]

为什么不是直接写 Elasticsearch?

  • ES 短暂不可用时,Kafka 可以缓冲日志。
  • 多个下游可以复用同一份日志。
  • 高峰期 Kafka 承接写入,下游按能力消费。

如果没有缓冲层,ES 写入压力过大时,应用日志可能丢失,甚至反过来拖慢业务线程。

场景三:用户行为埋点

用户点击、浏览、搜索、加购、支付等行为通常量很大,而且有多个下游使用。

mermaid
flowchart TD
    A["前端或服务端埋点"] --> B["网关或采集服务"]
    B --> C["Kafka: user.behavior"]
    C --> D["实时推荐"]
    C --> E["实时数仓"]
    C --> F["用户画像"]
    C --> G["反作弊"]

设计建议:

  • Topic 按事件大类划分,不要每个小事件都建 Topic。
  • key 可以用 userId,便于用户维度顺序分析。
  • value 使用统一事件结构,包含 eventId、userId、eventType、timestamp、traceId。

示例事件:

json
{
  "eventId": "evt-10001",
  "userId": 7,
  "eventType": "PRODUCT_VIEW",
  "occurredAt": 1782920000000,
  "properties": {
    "productId": 1001,
    "source": "search"
  }
}

场景四:CDC 数据同步

CDC 是 Change Data Capture,意思是捕获数据库变更。常见做法是读取 MySQL binlog,把变更写入 Kafka,再由下游消费。

mermaid
flowchart TD
    A["MySQL binlog"] --> B["CDC Connector"]
    B --> C["Kafka Topic"]
    C --> D["搜索索引同步"]
    C --> E["数据仓库"]
    C --> F["缓存刷新"]

商业价值:

  • 业务代码不需要到处写同步逻辑。
  • 数据变更可以被多个系统复用。
  • 可以回放 Topic 重建下游数据。

注意点:

  • CDC 不是万能事务,仍要处理重复事件和乱序。
  • 下游要用主键做幂等更新。
  • 删除事件要明确表达,不要只同步新增和修改。

场景五:搜索索引同步

商品、订单、用户数据经常需要同步到 Elasticsearch。同步方式可以是业务事件,也可以是 CDC。

mermaid
flowchart TD
    A["商品服务更新商品"] --> B["发送 product.changed"]
    B --> C["Kafka"]
    C --> D["搜索同步服务"]
    D --> E["查询数据库补全数据"]
    E --> F["写入 Elasticsearch"]

为什么消费者收到事件后还要查数据库?

事件里通常只放变更标识和必要字段,避免消息过大。搜索同步服务收到事件后查询主库,组装完整索引文档。这样可以减少事件结构和搜索索引结构的强耦合。

风险:

  • 如果消费者不幂等,重复消息会造成重复写或覆盖旧数据。
  • 如果事件乱序,旧数据可能覆盖新数据。
  • 如果没有重试和死信,索引会长期缺数据。

解决:

  • 文档使用业务主键作为 ES _id
  • 事件带版本号或更新时间,只允许新版本覆盖旧版本。
  • 失败进入 DLT,监控消费延迟和死信数量。

场景六:实时风控

支付、登录、提现、下单等事件可以进入 Kafka,风控系统实时消费并计算规则。

mermaid
flowchart TD
    A["交易事件"] --> B["Kafka"]
    B --> C["实时规则引擎"]
    C --> D{"是否风险高?"}
    D -->|"是"| E["冻结或人工审核"]
    D -->|"否"| F["正常通过"]

为什么适合 Kafka:

  • 风控规则需要多个业务域事件。
  • 事件量大,需要高吞吐。
  • 风控系统升级时,可以从 Kafka 回放历史数据验证规则。

注意:强实时拦截不能完全依赖异步 Kafka。如果必须在支付前阻断,仍需要同步风控接口;Kafka 更适合异步监控、补充分析、模型训练和后置处置。

场景七:削峰填谷

秒杀和活动期间,订单事件可能瞬间暴涨。Kafka 可以先承接流量,下游按固定速率消费。

mermaid
flowchart TD
    A["活动流量峰值"] --> B["Kafka 堆积"]
    B --> C["消费者限速处理"]
    C --> D["数据库压力平稳"]

但削峰不是无成本的:

  • 用户看到的结果可能延迟。
  • 消费延迟需要监控和告警。
  • 堆积过大可能导致磁盘压力。
  • 下游处理必须可重试、可补偿。

Topic 设计规范

设计点建议
命名业务域.事件名,例如 order.created
粒度按业务域和事件类型划分,不要过细也不要过粗
Key用业务实体 ID,例如 orderId、userId、deviceId
Value使用稳定事件结构,带 eventId、occurredAt、version
Headers放 traceId、source、schemaVersion 等元信息
保留时间按回放需要和磁盘成本设置

容量规划关注点

上线前至少要估算:

  • 每秒消息条数。
  • 单条消息平均大小和峰值大小。
  • 保留时间。
  • 副本数。
  • 消费组数量。
  • 峰值堆积时需要多少磁盘。

粗略估算:

text
每日数据量 = 每秒消息数 * 单条消息大小 * 86400
实际磁盘 = 每日数据量 * 保留天数 * 副本数

例如每秒 5000 条,每条 1 KB,保留 3 天,3 副本:

text
5000 * 1KB * 86400 * 3 * 3 ≈ 3.9 TB

这还没算索引、日志、压缩效果和预留空间。生产环境一定要留余量。

监控指标

指标说明异常含义
消费延迟 Lag消费者落后多少消息下游处理慢或消费者异常
Broker 磁盘使用率Kafka 数据文件占用保留时间过长或堆积过大
ISR 变化同步副本是否稳定Broker、网络或磁盘异常
请求延迟Produce/Fetch 延迟集群压力或网络问题
DLT 数量死信消息数量业务异常或数据格式错误

小结

Kafka 的商业价值在于事件复用、解耦、缓冲和回放。真正落地时,重点不是“能不能发消息”,而是 Topic 如何设计、key 如何选择、消费者如何幂等、失败如何重试、堆积如何监控、历史数据如何回放。