Kafka Streams Topology Description Plugin:让 Broker 记录并对外暴露 Streams Group 处理拓扑
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
本文基于 Apache Kafka 4.4 引入的 KIP-1331(Streams Group Topology Description Plugin)能力编写,核心内容取自仓库 docs/streams/developer-guide/topology-description-plugin.md。文中涉及的配置参数、接口签名与工作流程,均可在本仓库对应源码与测试中找到实现依据。
导读
从 Apache Kafka 4.4 开始,Broker 可以为每个streams group记录一份人类可读的处理拓扑描述(processing topology description)。Kafka Streams 客户端会把自己Topology#describe()得到的拓扑信息推送给 group coordinator,coordinator 将其交给一个可插拔的、Broker 侧的后端存储插件;运维人员随后通过AdminAPI 或bin/kafka-streams-groups.shCLI,无需访问应用源码、也不需要运行中的实例,就能查看任意 streams group 的拓扑。读完本文,你将掌握:该功能的启用配置、客户端侧开关、插件接口的完整实现规范、通过 Admin API / CLI 读取拓扑的方法,以及状态机(status)与故障处理语义。
Overview:功能定位与适用范围
该特性仅适用于使用Streams Rebalance Protocol(group.protocol=streams,KIP-1071)的 streams group。默认情况下它是关闭的:只有当 Broker 配置group.streams.topology.description.plugin.class指向一个StreamsGroupTopologyDescriptionPlugin实现时,Broker 才会去索要(solicit)、存储并提供拓扑描述。
当功能开启后:
- Kafka Streams 客户端会在 Broker 请求时自动推送拓扑描述,应用代码无需任何改动;可通过 Streams 配置
topology.description.push.enabled按客户端关闭推送。 - 描述以 group 的topology epoch作为版本依据,Broker 始终能判断已存描述是否与 group 当前运行的拓扑一致。
- 已存储的描述可通过
Admin#describeStreamsGroups(配合DescribeStreamsGroupsOptions#includeTopologyDescription(true))或kafka-streams-groups.sh --describe --topology获取。
仓库中docs/streams/developer-guide/streams-rebalance-protocol.md对该协议做了完整介绍;group.coordinator.rebalance.protocols自 4.3 起已标记弃用,streams 协议在 GroupCoordinatorConfig.java 中始终可用。
工作原理:push / describe 完整循环
该特性新增了一个 RPCStreamsGroupTopologyDescriptionUpdate,并扩展了既有的StreamsGroupHeartbeat与StreamsGroupDescribe两个 RPC。整个循环分六步:
1. Solicitation(索要)
当 group coordinator 尚未记录到该 group 当前 topology epoch 的成功推送时(例如新 group,或拓扑变更导致 epoch 增加),它会在StreamsGroupHeartbeat响应中置位TopologyDescriptionRequired标志。
2. Push(推送)
客户端看到该标志、且自身topology.description.push.enabled=true时,向 coordinator 发送StreamsGroupTopologyDescriptionUpdate请求,内容包含:group ID、member ID、topology epoch,以及完整的拓扑描述(subtopologies、sources、processors、sinks、state stores、global stores)。
Broker 侧的数据模型与转换逻辑集中在 StreamsGroupTopologyDescriptionConverter.java,它把线上的TopologyDescription结构转换为插件 API 使用的StreamsGroupTopologyDescription。
3. Store(存储)
Broker 校验发送者是该 group 的已知成员后,调用插件的setTopology(groupId, topologyEpoch, description)方法。成功后 Broker 记录已存储的 topology epoch 并停止索要。同一 epoch 可能多个成员并发推送,推送数据完全一致,插件必须幂等地处理并发。
4. Failure handling(失败处理)
若插件存储失败,Broker 区分两种情况:
- 永久失败(
StreamsTopologyDescriptionPermanentFailureException):表示该描述在当前 topology epoch 永远不会被接受(例如体积过大或语义被拒绝)。Broker 记录失败的 epoch,停止索要,直到 topology epoch 前进。 - 瞬时失败(
StreamsTopologyDescriptionTransientFailureException或任何其他异常):Broker 对该 group 启动指数退避(30 秒起步,上限 1 小时),并在之后的某次心跳中重新索要。
两种情况下,推送客户端收到的错误码都是STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED。客户端不会自行重试——重试完全由 Broker 通过心跳索要驱动。退避计时器的实现见 StreamsGroupTopologyDescriptionBackoff.java 的armOrExtend/clear/clearGroup方法。
5. Describe(查询)
调用方通过StreamsGroupDescribe(version 1 及以上,IncludeTopologyDescription=true)请求拓扑描述时,Broker 调用插件的getTopology(groupId, topologyEpoch),并把描述连同状态字段一起附到响应中(状态语义见下文「Interpreting the topology description status」)。
6. Deletion(删除)
streams group 被删除(DeleteGroups)或过期时,Broker 调用插件的deleteTopology(groupId)以让插件清理已存数据。失败语义见下文「Group deletion and GROUP_DELETION_FAILED」。
Broker 侧协调这些行为的核心是 StreamsGroupTopologyDescriptionManager.java,它提供maybeSetTopologyDescriptionRequired、completeEpochWrite、armBackoff、startCleanupCycle等方法,并通过PluginOutcome.success() / permanent() / transientFailure()三个工厂方法把插件调用结果映射为 Broker 内部状态;group 的 epoch 状态(StoredDescriptionTopologyEpoch/FailedDescriptionTopologyEpoch)则由 StreamsGroup.java 中的setStoredDescriptionTopologyEpoch/setFailedDescriptionTopologyEpoch维护,并经由 StreamsCoordinatorRecordHelpers.java 持久化到 coordinator 日志。
Broker 配置
| 配置项 | 说明 |
|---|---|
group.streams.topology.description.plugin.class | StreamsGroupTopologyDescriptionPlugin实现的完全限定类名。默认未设置,此时功能整体禁用:Broker 从不索要拓扑描述,describe 请求返回状态NOT_STORED。 |
该配置在 GroupCoordinatorConfig.java 中定义(STREAMS_GROUP_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS_CONFIG,类型为CLASS,默认值null),并注册进 coordinator 的CONFIG_DEF(同文件 L537),因此你可以像下面这样把它写进 Broker 的server.properties:
# 启用 streams group 拓扑描述插件 group.streams.topology.description.plugin.class=org.apache.kafka.server.streams.InMemoryTopologyDescriptionPluginApache Kafka 自带一个参考实现org.apache.kafka.server.streams.InMemoryTopologyDescriptionPlugin,源码位于 InMemoryTopologyDescriptionPlugin.java,用内存ConcurrentHashMap为每个 group 存一份描述。它面向测试、并作为真实实现的起点,不适合生产环境:Broker 重启后状态即丢失,且数据不跨 Broker 共享。
从该参考实现可以看到插件接口各方法的推荐实现方式(详见下一节):setTopology直接覆盖写入;deleteTopology按 groupId 移除;getTopology仅当请求的 topology epoch 与存储的 epoch 一致时返回描述,否则返回null——这与 Broker 侧的NOT_STORED语义精确对应。
客户端配置
| 配置项 | 说明 |
|---|---|
topology.description.push.enabled | 控制 Kafka Streams 客户端是否在 Broker 请求时发送拓扑描述。设为false时客户端不会准备或推送拓扑描述。默认开启。 |
该配置在 StreamsConfig.java 中定义为TOPOLOGY_DESCRIPTION_PUSH_ENABLED_CONFIG,并由 StreamThread.java 在心跳链路中消费。
需要注意:该配置只控制客户端是否响应 Broker 的索要。如果 Broker 侧没有配置插件,客户端永远不会被要求推送,无论此项设为什么。
实现一个插件
插件实现group-coordinator-api模块中的StreamsGroupTopologyDescriptionPlugin接口:
public interface StreamsGroupTopologyDescriptionPlugin extends Configurable, AutoCloseable { CompletableFuture<Void> setTopology(String groupId, int topologyEpoch, StreamsGroupTopologyDescription description); CompletableFuture<Void> deleteTopology(String groupId); CompletableFuture<StreamsGroupTopologyDescription> getTopology(String groupId, int topologyEpoch); }该接口标注为@InterfaceAudience.Public/@InterfaceStability.Evolving,即公共且仍在演进中的 API。接口的 Javadoc(同文件 L25-L80)逐条给出了实现契约,与官方文档互相印证,是实现时最重要的权威参考:
- 线程安全。
setTopology可能被同一 group 的多个成员并发调用;(groupId, topologyEpoch)相同的调用携带完全一致的数据,必须幂等。deleteTopology对同一 group 可能被调用多次(包括本就无数据时)。 - 异步完成 future,绝不同步抛出异常。失败必须通过让返回的 future 异常完成来传达;
setTopology的同步抛出会被视为一次带通用客户端可见错误消息的永久失败。 - 对失败分类:
setTopology的 future 用StreamsTopologyDescriptionPermanentFailureException完成表示该描述在当前 epoch 永远不会被接受;用StreamsTopologyDescriptionTransientFailureException(或任何其他异常)表示可重试的后端故障。永久/瞬时之分是 Broker 内部状态,推送客户端统一只看到STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED+ 异常消息。 - 按 group 与 epoch 为键存取。
getTopology(groupId, topologyEpoch)仅在请求 epoch 与存储一致时才应返回描述;当插件已无该数据(如后端被清空)时以null完成,Broker 随即上报状态NOT_STORED;若 future 异常完成,Broker 为该 group 上报读错误(状态ERROR)。注意:Broker 只会在尚未记录当前 topology epoch 的成功推送时才重新索要;插件丢失已存数据后不会自动被重新填充,直到 topology epoch 前进。因此生产实现必须使用持久化存储。 - 生命周期。插件在每个 Broker 上实例化一次,通过
Configurable#configure(Map)传入 Broker 配置完成初始化,在 Broker 关闭时通过AutoCloseable#close()释放资源。
作为对照,InMemoryTopologyDescriptionPlugin的configure为空实现、setTopology直接store.put(...)并返回completedFuture(null)、getTopology做 epoch 匹配、close清空 map——它演示了"如何满足接口契约"的最小可行形态,而生产插件应当在此之上替换为持久化后端并妥善分类失败。
数据结构速览:插件 API 中的StreamsGroupTopologyDescription是一个 record,由subtopologies与globalStores组成;节点模型与org.apache.kafka.streams.TopologyDescription形状一致(Source、Processor、Sink三种 sealed 的Node),但它位于group-coordinator-api模块,因此插件实现无需依赖kafka-streams库。线上 schema 只携带后继关系(successor relation),需要前驱方向的插件可在一次遍历中自行重建。
读取拓扑描述
通过 Admin API
向Admin#describeStreamsGroups传入DescribeStreamsGroupsOptions#includeTopologyDescription(true):
try (Admin admin = Admin.create(props)) { DescribeStreamsGroupsResult result = admin.describeStreamsGroups( List.of("my-streams-app"), new DescribeStreamsGroupsOptions().includeTopologyDescription(true)); StreamsGroupDescription description = result.describedGroups().get("my-streams-app").get(); StreamsGroupTopologyDescriptionStatus status = description.topologyDescriptionStatus(); Optional<StreamsGroupTopologyDescription> topology = description.topologyDescription(); }返回的StreamsGroupTopologyDescription镜像了org.apache.kafka.streams.TopologyDescription(含 source / processor / sink 节点的 subtopologies,以及 global stores),但不需要依赖kafka-streams库即可在 Broker 侧解析。若目标 Broker 不支持该特性(版本低于 4.4),请求会以UnsupportedVersionException失败。
通过 CLI
bin/kafka-streams-groups.sh的--describe配合--topology选项:
kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --topology当描述可用时,输出格式与Topology#describe()完全一致:
Topologies: Sub-topology: 0 Source: KSTREAM-SOURCE-0000000000 (topics: [streams-plaintext-input]) --> KSTREAM-FLATMAPVALUES-0000000001 Processor: KSTREAM-FLATMAPVALUES-0000000001 (stores: []) --> KSTREAM-AGGREGATE-0000000002 <-- KSTREAM-SOURCE-0000000000 Processor: KSTREAM-AGGREGATE-0000000002 (stores: [counts-store]) --> KSTREAM-SINK-0000000003 <-- KSTREAM-FLATMAPVALUES-0000000001 Sink: KSTREAM-SINK-0000000003 (topic: streams-wordcount-output) <-- KSTREAM-AGGREGATE-0000000002若无描述可用,工具会打印说明信息并以非零退出码结束。CLI 的完整参考见 kafka-streams-groups.sh 文档。
解读拓扑描述状态(status)
每个请求了拓扑描述的 describe 响应都携带一个StreamsGroupTopologyDescriptionStatus。描述本身仅在状态为AVAILABLE时出现:
| 状态 | 含义 |
|---|---|
NOT_REQUESTED | 调用方未请求拓扑描述(未设置includeTopologyDescription(true))。 |
NOT_STORED | 该 group 未记录任何拓扑描述——例如 Broker 未配置拓扑描述插件,或客户端尚未推送。 |
ERROR | Broker 从插件获取拓扑描述失败,详见 Broker 日志。 |
AVAILABLE | 拓扑描述可用并已随响应返回。 |
组删除与 GROUP_DELETION_FAILED
当启用了拓扑描述插件的 streams group 被删除时,Broker 会在移除 group 前调用插件的deleteTopology方法。若插件删除失败:
DeleteGroups请求针对该 group 返回错误码GROUP_DELETION_FAILED,插件异常消息放在 per-group 的ErrorMessage字段(DeleteGroupsversion 3 及以上可用),且Broker 不会删除该 group。- 重试删除会幂等地再次调用
deleteTopology。 - 因周期清理而过期的 group 处理方式相同——其删除会推迟到后续清理周期,直到插件删除成功为止。
接口层面,StreamsGroupTopologyDescriptionPlugin#deleteTopology的 Javadoc 明确:future 异常完成即向DeleteGroups调用方上报GROUP_DELETION_FAILED,Broker 不会对 group 打 tombstone;周期清理路径对失败的处理一致。
可观测性:Broker 指标
Broker 为每次插件交互暴露指标,MBean 组为kafka.server:type=group-coordinator-metrics,完整列表见 group coordinator 监控参考。每个 sensor 同时发布-rate(每秒)与-count(累计)两个指标,例如streams-group-topology-description-set-success会展开为streams-group-topology-description-set-success-rate与streams-group-topology-description-set-success-count。
| 指标 | 含义 |
|---|---|
streams-group-topology-description-set-success/set-error | setTopology调用结果,由客户端推送驱动。 |
streams-group-topology-description-get-success/get-error | getTopology调用结果,由 describe 请求驱动。 |
streams-group-topology-description-delete-success/delete-error | deleteTopology调用结果,由组删除与清理驱动。 |
streams-group-topology-description-cleanup-cycle | coordinator 已运行的周期清理轮数。 |
streams-group-topology-description-cleanup-eligible | 清理扫描发现的可删除插件状态(eligible)的 group 数。 |
排查时优先关注-error指标:get-error上升解释了 describe 响应中的ERROR,set-error上升解释了描述迟迟不出现,delete-error上升解释了GROUP_DELETION_FAILED。
Troubleshooting 故障排查
读 Broker 日志之前,先查看上文 Observability 中的streams-group-topology-description-*-error指标——它们能直接定位是set、get还是delete哪个插件调用在失败。
--topology报告 "No topology description is stored"(状态NOT_STORED)
- 确认
group.streams.topology.description.plugin.class已设置在所有托管该 group coordinator 的 Broker 上。缺少该配置则功能整体禁用。 - 确认应用未设置
topology.description.push.enabled=false。 - 如果 group(或其 topology epoch)是新建的,客户端可能只是还没推送——Broker 通过心跳索要,描述通常会在几个心跳间隔内出现。
- 若描述仍不出现,检查 Broker 日志中失败的
setTopology调用。发生永久失败(例如插件拒绝该描述)后,Broker 会停止索要,直到 topology epoch 前进。 - 若描述以前有、现在消失了,可能是插件丢失了已存数据。当插件的
getTopology返回null时,Broker 以NOT_STORED呈现该状态(记录WARN日志)并在后续 describe 中持续返回NOT_STORED。由于 Broker 只在 topology epoch 前进时才重新索要,不提升 epoch 就重启应用无法恢复描述——需要推进 topology epoch 或清空插件状态以触发一次全新推送。
describe 时状态为ERROR
- 插件在 Broker 上的
getTopology调用失败。查看 group coordinator 所在 Broker 的日志,寻找底层异常。
推送投递问题
- 失败的推送在客户端表现为
STREAMS_TOPOLOGY_DESCRIPTION_UPDATE_FAILED,并由 Streams 客户端记录日志;客户端不会自行重试。Broker 通过心跳重新索要——瞬时插件失败带 30 秒至 1 小时的指数退避。 - 若推送成员已被 fencing(fenced)或 group 已被删除,推送会以
UNKNOWN_MEMBER_ID失败,客户端重新加入 group;这是预期行为,可自愈。
DeleteGroups报GROUP_DELETION_FAILED
- 插件删除已存描述失败。per-group 错误消息包含插件失败原因,Broker 日志包含完整异常。解决插件/后端问题后重试删除即可——
deleteTopology会被幂等地再次调用。
请求拓扑描述时抛UnsupportedVersionException
- 目标 Broker 版本低于 Apache Kafka 4.4,不支持
StreamsGroupDescribeversion 1。升级 Broker,或在不请求拓扑描述的情况下 describe group。
相关阅读
- Streams Rebalance Protocol 文档:streams group 协议基础,本特性的前置条件。
- kafka-streams-groups.sh 文档:CLI 完整参数参考。
- Streams 配置文档:
topology.description.push.enabled等客户端配置说明。 - group coordinator 监控参考:本节指标的完整清单。
- 升级指南 与 Streams 升级指南:4.4 版本升级注意事项。
- 集成测试参考:GroupCoordinatorServiceTopologyDescriptionTest.java 覆盖了无插件场景;streams_topology_description_plugin_test.py 覆盖了端到端插件行为。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考