{"record":{"id":"c0639ee66b19f488","repo":"apache/kafka","slug":"unknown-topology-node-type-nodetype","errorCode":null,"errorMessage":"Unknown topology node type: {nodeType}","messagePattern":"Unknown topology node type: (.+?)","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"clients/src/main/java/org/apache/kafka/clients/admin/internals/DescribeStreamsGroupsHandler.java","lineNumber":304,"sourceCode":"    }\n\n    private StreamsGroupTopologyDescription.Node convertTopologyNode(\n            final StreamsGroupDescribeResponseData.TopologyDescriptionNode node,\n            final Map<String, Set<String>> predecessors) {\n        final Set<String> successors = Set.copyOf(node.successors());\n        final Set<String> nodePredecessors = predecessors.getOrDefault(node.name(), Set.of());\n        switch (node.nodeType()) {\n            case NODE_TYPE_SOURCE:\n                return new StreamsGroupTopologyDescription.Source(\n                    node.name(), Set.copyOf(node.sourceTopics()), successors, nodePredecessors);\n            case NODE_TYPE_PROCESSOR:\n                return new StreamsGroupTopologyDescription.Processor(\n                    node.name(), Set.copyOf(node.stores()), successors, nodePredecessors);\n            case NODE_TYPE_SINK:\n                return new StreamsGroupTopologyDescription.Sink(\n                    node.name(), Optional.ofNullable(node.sinkTopic()), successors, nodePredecessors);\n            default:\n                throw new IllegalStateException(\"Unknown topology node type: \" + node.nodeType());\n        }\n    }\n\n    private StreamsGroupMemberAssignment.TaskIds convertTaskIds(final StreamsGroupDescribeResponseData.TaskIds taskIds) {\n        return new StreamsGroupMemberAssignment.TaskIds(\n            taskIds.subtopologyId(),\n            taskIds.partitions()\n        );\n    }\n\n    private StreamsGroupMemberAssignment convertAssignment(final StreamsGroupDescribeResponseData.Assignment assignment) {\n        return new StreamsGroupMemberAssignment(\n            assignment.activeTasks().stream().map(this::convertTaskIds).collect(Collectors.toList()),\n            assignment.standbyTasks().stream().map(this::convertTaskIds).collect(Collectors.toList()),\n            assignment.warmupTasks().stream().map(this::convertTaskIds).collect(Collectors.toList())\n        );\n    }\n","sourceCodeStart":286,"sourceCodeEnd":322,"githubUrl":"https://github.com/apache/kafka/blob/996fb4585aa1bcc8980b0e1b8d6b168b986cd979/clients/src/main/java/org/apache/kafka/clients/admin/internals/DescribeStreamsGroupsHandler.java#L286-L322","documentation":"Thrown by DescribeStreamsGroupsHandler.convertTopologyNode when the node.nodeType() byte does not match NODE_TYPE_SOURCE, NODE_TYPE_PROCESSOR, or NODE_TYPE_SINK. The client received a topology node type it does not recognize, indicating either a newer broker adding a new node kind or a corrupt payload.","triggerScenarios":"Describing a Streams group whose topology contains a node type byte outside the three known values; common when the broker is newer than the client.","commonSituations":"Client older than broker (rolling upgrade); experimental Streams features; corrupt response.","solutions":["Upgrade the kafka-clients dependency to match the broker version.","Restrict the client to API versions supported by both ends.","Catch IllegalStateException around describeStreamsGroups and surface a clear 'unsupported version' error."],"exampleFix":"// before\nMap<String, StreamsGroupDescription> r =\n    admin.describeStreamsGroups(groups).all().get();\n\n// after\ntry {\n    Map<String, StreamsGroupDescription> r =\n        admin.describeStreamsGroups(groups).all().get();\n} catch (ExecutionException e) {\n    if (e.getCause() instanceof IllegalStateException\n        && e.getCause().getMessage().contains(\"Unknown topology node type\")) {\n        throw new IllegalStateException(\"Upgrade kafka-clients to match broker\", e);\n    }\n    throw e;\n}","handlingStrategy":"try-catch","validationCode":"// Cannot pre-validate unknown node types; ensure client >= broker version\nassert clientVersionIsAtLeast(brokerVersion);","typeGuard":null,"tryCatchPattern":"try {\n    return admin.describeStreamsGroups(groups).all().get();\n} catch (ExecutionException e) {\n    if (e.getCause() instanceof IllegalStateException\n        && e.getCause().getMessage().contains(\"Unknown topology node type\")) {\n        throw new IllegalStateException(\"kafka-clients out of date; upgrade to match broker\", e);\n    }\n    throw e;\n}","preventionTips":["Upgrade kafka-clients in lockstep with broker upgrades.","Restrict negotiate api versions to those the client understands.","Add version-mismatch detection to your client health checks."],"tags":["admin-client","streams","wire-protocol","illegal-state","version-mismatch"],"backgroundTag":null,"analyzedSha":"996fb4585aa1bcc8980b0e1b8d6b168b986cd979","analyzedAt":"2026-08-11T22:03:28.655Z","contentChangedAt":null,"schemaVersion":2},"datasetVersion":"2026-09-14T00:17:10.932Z"}