Kafka Streams Topology Description Plugin:让 Broker 记录并对外暴露 Streams Group 处理拓扑

发布时间:2026/9/10 13:44:44
Kafka Streams Topology Description Plugin:让 Broker 记录并对外暴露 Streams Group 处理拓扑 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-1331Streams 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 coordinatorcoordinator 将其交给一个可插拔的、Broker 侧的后端存储插件运维人员随后通过AdminAPI 或bin/kafka-streams-groups.shCLI无需访问应用源码、也不需要运行中的实例就能查看任意 streams group 的拓扑。读完本文你将掌握该功能的启用配置、客户端侧开关、插件接口的完整实现规范、通过 Admin API / CLI 读取拓扑的方法以及状态机status与故障处理语义。Overview功能定位与适用范围该特性仅适用于使用Streams Rebalance Protocolgroup.protocolstreamsKIP-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.enabledtrue时向 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查询调用方通过StreamsGroupDescribeversion 1 及以上IncludeTopologyDescriptiontrue请求拓扑描述时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.classStreamsGroupTopologyDescriptionPlugin实现的完全限定类名。默认未设置此时功能整体禁用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.classorg.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 { CompletableFutureVoid setTopology(String groupId, int topologyEpoch, StreamsGroupTopologyDescription description); CompletableFutureVoid deleteTopology(String groupId); CompletableFutureStreamsGroupTopologyDescription 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(); OptionalStreamsGroupTopologyDescription topology description.topologyDescription(); }返回的StreamsGroupTopologyDescription镜像了org.apache.kafka.streams.TopologyDescription含 source / processor / sink 节点的 subtopologies以及 global stores但不需要依赖kafka-streams库即可在 Broker 侧解析。若目标 Broker 不支持该特性版本低于 4.4请求会以UnsupportedVersionException失败。通过 CLIbin/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 未配置拓扑描述插件或客户端尚未推送。ERRORBroker 从插件获取拓扑描述失败详见 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_FAILEDBroker 不会对 group 打 tombstone周期清理路径对失败的处理一致。可观测性Broker 指标Broker 为每次插件交互暴露指标MBean 组为kafka.server:typegroup-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-errorsetTopology调用结果由客户端推送驱动。streams-group-topology-description-get-success/get-errorgetTopology调用结果由 describe 请求驱动。streams-group-topology-description-delete-success/delete-errordeleteTopology调用结果由组删除与清理驱动。streams-group-topology-description-cleanup-cyclecoordinator 已运行的周期清理轮数。streams-group-topology-description-cleanup-eligible清理扫描发现的可删除插件状态eligible的 group 数。排查时优先关注-error指标get-error上升解释了 describe 响应中的ERRORset-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.enabledfalse。如果 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 小时的指数退避。若推送成员已被 fencingfenced或 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),仅供参考