Kafka商业常用场景
Kafka 在商业系统里最常见的不是“博客实战”,而是承接大规模事件流。它通常位于业务系统、数据平台、搜索、风控、监控之间,作为事件管道和缓冲层。
场景一:订单事件总线
订单创建后,很多系统都关心这件事:
- 库存系统扣减或锁定库存。
- 搜索系统同步订单索引。
- 风控系统判断异常交易。
- 积分系统发放积分。
- 数据平台沉淀交易数据。
flowchart TD
A["订单服务"] --> B["Topic: order.events"]
B --> C["库存消费组"]
B --> D["风控消费组"]
B --> E["搜索消费组"]
B --> F["数据仓库消费组"]为什么适合 Kafka:
- 多个系统可以使用不同 groupId 独立消费。
- 新增下游系统时不用改订单服务。
- 消息保留期内可以重新消费补数据。
注意点:
- 事件命名用过去式,例如
ORDER_CREATED、ORDER_PAID。 - 同一订单用
orderId作为 key,保证订单维度顺序。 - 消费端必须做幂等和状态机判断。
场景二:日志采集和可观测性
应用日志、访问日志、审计日志可以先写入 Kafka,再由不同系统消费。
flowchart TD
A["应用服务"] --> B["日志采集 Agent"]
B --> C["Kafka: app.logs"]
C --> D["Elasticsearch"]
C --> E["对象存储"]
C --> F["告警系统"]为什么不是直接写 Elasticsearch?
- ES 短暂不可用时,Kafka 可以缓冲日志。
- 多个下游可以复用同一份日志。
- 高峰期 Kafka 承接写入,下游按能力消费。
如果没有缓冲层,ES 写入压力过大时,应用日志可能丢失,甚至反过来拖慢业务线程。
场景三:用户行为埋点
用户点击、浏览、搜索、加购、支付等行为通常量很大,而且有多个下游使用。
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。
示例事件:
{
"eventId": "evt-10001",
"userId": 7,
"eventType": "PRODUCT_VIEW",
"occurredAt": 1782920000000,
"properties": {
"productId": 1001,
"source": "search"
}
}场景四:CDC 数据同步
CDC 是 Change Data Capture,意思是捕获数据库变更。常见做法是读取 MySQL binlog,把变更写入 Kafka,再由下游消费。
flowchart TD
A["MySQL binlog"] --> B["CDC Connector"]
B --> C["Kafka Topic"]
C --> D["搜索索引同步"]
C --> E["数据仓库"]
C --> F["缓存刷新"]商业价值:
- 业务代码不需要到处写同步逻辑。
- 数据变更可以被多个系统复用。
- 可以回放 Topic 重建下游数据。
注意点:
- CDC 不是万能事务,仍要处理重复事件和乱序。
- 下游要用主键做幂等更新。
- 删除事件要明确表达,不要只同步新增和修改。
场景五:搜索索引同步
商品、订单、用户数据经常需要同步到 Elasticsearch。同步方式可以是业务事件,也可以是 CDC。
flowchart TD
A["商品服务更新商品"] --> B["发送 product.changed"]
B --> C["Kafka"]
C --> D["搜索同步服务"]
D --> E["查询数据库补全数据"]
E --> F["写入 Elasticsearch"]为什么消费者收到事件后还要查数据库?
事件里通常只放变更标识和必要字段,避免消息过大。搜索同步服务收到事件后查询主库,组装完整索引文档。这样可以减少事件结构和搜索索引结构的强耦合。
风险:
- 如果消费者不幂等,重复消息会造成重复写或覆盖旧数据。
- 如果事件乱序,旧数据可能覆盖新数据。
- 如果没有重试和死信,索引会长期缺数据。
解决:
- 文档使用业务主键作为 ES
_id。 - 事件带版本号或更新时间,只允许新版本覆盖旧版本。
- 失败进入 DLT,监控消费延迟和死信数量。
场景六:实时风控
支付、登录、提现、下单等事件可以进入 Kafka,风控系统实时消费并计算规则。
flowchart TD
A["交易事件"] --> B["Kafka"]
B --> C["实时规则引擎"]
C --> D{"是否风险高?"}
D -->|"是"| E["冻结或人工审核"]
D -->|"否"| F["正常通过"]为什么适合 Kafka:
- 风控规则需要多个业务域事件。
- 事件量大,需要高吞吐。
- 风控系统升级时,可以从 Kafka 回放历史数据验证规则。
注意:强实时拦截不能完全依赖异步 Kafka。如果必须在支付前阻断,仍需要同步风控接口;Kafka 更适合异步监控、补充分析、模型训练和后置处置。
场景七:削峰填谷
秒杀和活动期间,订单事件可能瞬间暴涨。Kafka 可以先承接流量,下游按固定速率消费。
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 等元信息 |
| 保留时间 | 按回放需要和磁盘成本设置 |
容量规划关注点
上线前至少要估算:
- 每秒消息条数。
- 单条消息平均大小和峰值大小。
- 保留时间。
- 副本数。
- 消费组数量。
- 峰值堆积时需要多少磁盘。
粗略估算:
每日数据量 = 每秒消息数 * 单条消息大小 * 86400
实际磁盘 = 每日数据量 * 保留天数 * 副本数例如每秒 5000 条,每条 1 KB,保留 3 天,3 副本:
5000 * 1KB * 86400 * 3 * 3 ≈ 3.9 TB这还没算索引、日志、压缩效果和预留空间。生产环境一定要留余量。
监控指标
| 指标 | 说明 | 异常含义 |
|---|---|---|
| 消费延迟 Lag | 消费者落后多少消息 | 下游处理慢或消费者异常 |
| Broker 磁盘使用率 | Kafka 数据文件占用 | 保留时间过长或堆积过大 |
| ISR 变化 | 同步副本是否稳定 | Broker、网络或磁盘异常 |
| 请求延迟 | Produce/Fetch 延迟 | 集群压力或网络问题 |
| DLT 数量 | 死信消息数量 | 业务异常或数据格式错误 |
小结
Kafka 的商业价值在于事件复用、解耦、缓冲和回放。真正落地时,重点不是“能不能发消息”,而是 Topic 如何设计、key 如何选择、消费者如何幂等、失败如何重试、堆积如何监控、历史数据如何回放。
