Kafka MCP

io.github.denizgursoyv0.2.1Updated Oct 7, 2026

Debug Kafka from an agent: find messages, measure consumer lag, unblock consumers, compare clusters.

Overview

AI-generated overview

Lets an assistant inspect and debug Kafka clusters: list topics, sample and search messages, measure consumer lag, compare clusters, and perform approved…

What it does
Connects to one or more Kafka clusters and exposes tools such as list_clusters, list_topics, describe_topic, sample_messages, get_message, get_schema, consumer_lag, describe_consumer_group, open_transactions and compare_clusters. Write tools include create_topic, add_partitions, alter_topic_config, commit_offset, delete_records, delete_topic, delete_consumer_group, produce_message and copy_message. Messages are decoded to JSON via Schema Registry or configured topic formats. Endpoints carry policy: read-only endpoints hide write tools entirely, and every tool call is audit logged.
When to use it
Useful when debugging Kafka from an agent: finding a specific message, checking why a consumer group is lagging, comparing preproduction against production, or making an approved change such as creating a topic or committing an offset.
Requirements
Runs locally over stdio as a binary, MCP bundle, or container image. Needs a configuration file (CONFIG_FILE) describing clusters, brokers, TLS/SASL credentials, and endpoints; Schema Registry URLs and schema files are optional. Kafka credentials may come from environment variables or secret files. Network access to the brokers is required.
Before you install
Write tools change cluster state: create_topic, add_partitions, alter_topic_config, commit_offset, delete_records, delete_topic, delete_consumer_group, produce_message and copy_message. Batch writes are not transactions and cannot be rolled back. read_only and confirm are guardrails, not security; real enforcement is Kafka ACLs on the SASL principal. HTTP endpoints have no authentication of their own, so CORS allow_origins is the only built-in barrier. Credentials such as KAFKA_PASSWORD…

Installation

In SourceWeft

  1. Open Kafka MCP in the dashboard and add it to a workspace.
  2. Enable the server for the chats that should use its tools.

Desktop only via STDIO. STDIO servers start a local process, so they need the SourceWeft desktop host.

Other MCP clients

Follow the launch instructions in the repository.

README

kafka-mcp

[License] [Coverage] [Web] [Release] [Kafka MCP server – quality and maintenance score on Glama]

An MCP server that exposes Kafka debugging as tools an LLM can call. It speaks MCP over stdio by default, optionally serves streamable HTTP, and talks to Kafka with franz-go.

One process can connect to several Kafka clusters. In stdio mode, one endpoint is selected for the session. In HTTP mode, each endpoint has its own path, so a session is bound to one cluster by how it connects rather than by a parameter a caller could forget to send.

Install

Every release ships the same server four ways:

  • Binary: kafka-mcp_<Os>_<arch> archives on the releases page. Put kafka-mcp on your PATH.

  • MCP Bundle: kafka-mcp_<version>_<os>_<arch>.mcpb on the same page, for macOS on Apple Silicon, Linux (amd64, arm64) and Windows (amd64). Open it in a client that installs bundles, such as Claude Desktop. It asks for your configuration file and, when that file has several endpoints, which one to serve. Intel Macs use the binary or the image.

  • Container image: ghcr.io/denizgursoy/kafka-mcp:<tag>. It starts with --server; for stdio, mount the file and pass --server=false:

    sh
    docker run -i --rm \  --mount type=bind,src=$PWD/kafka-mcp.yaml,dst=/config/kafka-mcp.yaml,readonly \  -e CONFIG_FILE=/config/kafka-mcp.yaml \  ghcr.io/denizgursoy/kafka-mcp:latest --server=false --endpoint local
  • MCP Registry: listed as io.github.denizgursoy/kafka-mcp in the official registry, with the image and the bundles, so a client that reads the registry can install it from there.

How releases reach the registry and the other catalogs is in PUBLISHING.md.

Configuration

kafka-mcp.{toml,yaml,yml,json} in the working directory, ~/.config/kafka-mcp/ or /etc configures the server. The http block is used only with --server.

yaml
http:  address: ":8090"  base_path: /kafka-mcp # optional; prefixes every endpoint and /healthzoutput_dir: /var/tmp/kafka-mcpclusters:  prod:    brokers:      - kafka-1:9093      - kafka-2:9093    security:      tls:        enabled: true        ca_file: /etc/kafka/ca.pem        # cert_file: /etc/kafka/client.pem # optional mTLS; requires key_file        # key_file: /run/secrets/client.key      sasl:        - scram:            enabled: true            algorithm: SCRAM-SHA-256 # or SCRAM-SHA-512            user: kafka-mcp-readonly            pass: "{env:KAFKA_PASSWORD}" # or password_file: /run/secrets/kafka  preprod:    brokers: kafka-preprod:9093endpoints:  prod-read:    cluster: prod    path: /mcp    description: Production investigation and debugging    read_only: true  prod-write:    cluster: prod    path: /mcp/rw    description: Approved production changes    tools:      copy_message: false  preprod:    cluster: preprod    path: /mcp/preprod

clusters owns Kafka connection details: brokers, TLS and SASL. endpoints owns the policy and, in HTTP mode, the MCP route. Stdio selects an endpoint by name and does not use its path. Several endpoints may reference one cluster, so the example reuses one production connection at /kafka-mcp/mcp in read-only mode and /kafka-mcp/mcp/rw in writable mode. Paths are exact: /mcp does not capture /mcp/rw. description is optional and is reported by server_config so a caller knows what the endpoint is intended for.

http.base_path prefixes endpoint paths and the liveness route. A leading or trailing slash on an endpoint path is optional. Paths must be unique and may not be /, /healthz, or contain a query or fragment. When path is omitted, it defaults to /mcp/<endpoint-name>.

For compatibility, a file with no endpoints block still creates one endpoint per cluster at /mcp/<cluster-name>. Existing cluster-level read_only and tools values continue to apply as a lower bound during migration; an explicit endpoint cannot widen them. New configurations should put both fields under endpoints.

An endpoint may switch individual tools off, by name:

yaml
endpoints:  prod-read:    cluster: prod    path: /mcp    read_only: true    tools:      create_topic: false      commit_offset: false

A tool the map does not mention stays exposed, so the file states only what it withholds rather than relisting every tool and silently losing whatever is added later. Names are exact and lowercase, as listed by server_config. They are checked at startup: a name that is not a tool stops the server, because a typo would leave the tool it was meant to withhold exposed. server_config cannot be switched off, since it is how a session learns which cluster and policy it reached and which tools that endpoint has. The tools map only narrows an endpoint: read_only: true still withholds the writing tools regardless of what the map says.

A minimal local file:

yaml
clusters:  local:    brokers: localhost:19092    schema_registry:      urls: ["http://localhost:18081"]endpoints:  local:    cluster: local    path: /mcp/local

Run a custom configuration over stdio with CONFIG_FILE=/path/to/config.yaml go run ./cmd/server. If it defines several endpoints, select one with --endpoint <name>. Add --server to serve every configured endpoint over HTTP instead; --endpoint is not used in HTTP mode. Without CONFIG_FILE, chu discovers kafka-mcp.{toml,yaml,yml,json} first in the working directory, then in the operating system's user config directory (~/.config/kafka-mcp/ on Linux), and finally in /etc. It uses the first matching file rather than merging files. Its standard loader order is defaults, file, HTTP, then environment; environment overrides use the KAFKA_MCP_ prefix (for example, KAFKA_MCP_HTTP_ADDRESS=:9000 or KAFKA_MCP_HTTP_BASE_PATH=/kafka-mcp). Logging can be configured with LOG_LEVEL and LOG_PRETTY.

At least one cluster with a broker is required. http.address defaults to :8090, http.base_path defaults to the HTTP root, and output_dir defaults to the system temp directory. Exports are confined to that directory: output_file takes a file name, never a path.

brokers accepts either one address as a scalar or several addresses as a YAML list. Each address is passed to Kafka as a separate seed broker.

Keep secrets out of the file with {env:VAR} or password_file. Unknown fields are ignored by chu.

TLS is configured with twmb/tlscfg, using its TLS 1.2 minimum and recommended cipher suites. TLS uses the system trust store; ca_file adds a custom CA. For mTLS, supply both cert_file and key_file. security.sasl is a preference-ordered list: each entry enables either scram or plain, or uses oauth as shown below. For PLAIN, use plain: {enabled: true, user: alice, pass: "{env:KAFKA_PASSWORD}"}. Both accept optional zid (authorization identity) and password_file instead of pass; SCRAM also accepts is_token: true for delegation tokens. Algorithms are case-insensitive. Disabled entries are ignored; repeated mechanisms and entries enabling multiple mechanisms are rejected. Fallback negotiates a broker-supported mechanism; it does not retry invalid credentials. All Kafka connections, including message-reading sessions, use these settings.

Legacy per-cluster tls and sasl: {mechanism, user, password, password_file} remain supported, but cannot be combined with security on the same cluster.

For OAUTHBEARER client credentials, use this cluster security block:

yaml
security:  tls:    enabled: true  sasl:    - oauth:        enabled: true        token_url: https://identity.example.com/realms/apps/protocol/openid-connect/token        client_id: kafka-mcp        client_secret: "{env:KAFKA_CLIENT_SECRET}"        scopes: [kafka]        timeout: 10s        # zid: optional-authorization-id        # extensions: {tenant: example}

Alternatively, use oauth: {enabled: true, token: "{env:KAFKA_TOKEN}"} for a static token. Configure exactly one mode. Client-credentials tokens are fetched on authentication, cached per cluster across admin and reading connections, and renewed near expiry using expires_in. Token requests use the current authentication context and a timeout (default 10 seconds). Static tokens are not renewed. The broker must support OAUTHBEARER and trust the identity provider. Broker security.tls settings apply to Kafka connections; HTTPS token endpoints use the HTTP client's system trust store.

A cluster that cannot be reached at startup is still served, and list_clusters reports it as disconnected. One cluster being down must not block debugging the others.

Message formats

Messages are decoded to JSON for every tool that shows or searches them, by the first rule that applies:

  1. The topic has a format in topic_formats.
  2. The bytes carry a Schema Registry header (Confluent wire format) and the cluster has a schema_registry: Avro, Protobuf (with message index and referenced schemas) or JSON Schema.
  3. Otherwise JSON, then UTF-8 text, then base64.
yaml
clusters:  prod:    brokers: kafka-1:9093    schema_registry:      urls: ["https://schema-registry:8081"]      user: kafka-mcp                     # or bearer_token: "{env:SR_TOKEN}"      password: "{env:SR_PASSWORD}"       # or password_file      # tls: {enabled: true, ca_file: /etc/sr/ca.pem}    topic_formats:                        # topics whose messages carry no schema id      "orders.*":                         # exact name or glob; longest match wins        value:          format: protobuf          proto_files: [shop/order.proto]          import_paths: [./protos]        # or descriptor_set: ./orders.binpb          message_type: shop.Order      events:        value: {format: avro, schema_file: ./schemas/event.avsc}      metrics:        key: {format: text}        value: {format: msgpack}

Formats are avro, protobuf, json, msgpack, text and binary. Schema files and .proto sources are loaded at startup, so a missing file or type stops the server. A message that names a schema but cannot be decoded is returned as base64 with decode_error saying why, never silently. JSON output has sorted keys; Avro longs beyond 2^53 keep their exact value; unset Protobuf fields are shown with their defaults.

Each rendered message reports format, schema_id and message_type for the value, and key_format / key_schema_id when the key is not plain text.

server_config reports the endpoint name, exact path, description, policy and the cluster it serves, and never the password. Its authentication and sasl_user describe the first configured option; sasl_options lists all configured identities in preference order, not the mechanism negotiated by an individual broker connection. It also reports schema_registry URLs and topic_formats.

Browser clients (CORS)

A command-line client sends no Origin header and needs none of this. A client running in a browser does: the browser discards the response unless the server allows the origin, so the defaults already cover the MCP transport.

yaml
http:  address: ":8090"  cors:    allow_origins: ["https://mcp-client.example"]    allow_methods: [GET, POST, DELETE, OPTIONS]    allow_headers: [content-type, accept, authorization, cache-control, last-event-id, mcp-session-id, mcp-protocol-version]    expose_headers: [Mcp-Session-Id]    allow_private_network: true    max_age: 600

http.cors is ada's CORS middleware configuration, read straight from the file, so every option that middleware has is available here. Each key is optional and keeps its own default, so setting allow_origins alone does not drop the rest. The defaults are the values above with allow_origins: ["*"].

Three of them are load-bearing for the MCP transport. allow_methods needs GET, POST and DELETE: requests are posted, the event stream is a GET, and a client ends its session with DELETE. allow_headers needs mcp-session-id and mcp-protocol-version, which the client sends from the second request onwards, and a header missing there fails the whole preflight rather than being dropped. Mcp-Session-Id must stay in expose_headers: the session id arrives on the initialize response, and a page that cannot read it cannot make a second call.

allow_private_network answers Chrome's Private Network Access preflight, which a page on a public address must pass before it may reach a server on a private or loopback address; it defaults to on, and is granted only on a preflight the rest of the policy already allowed. Setting allow_credentials together with a wildcard allow_origins is refused by the middleware at startup unless unsafe_wildcard_origin_with_allow_credentials is also set, which it should not be.

These endpoints have no authentication of their own. An allowed origin can drive every tool with the server's Kafka credentials, from any page the browser's user happens to visit. allow_origins is the only built-in HTTP barrier, so narrow it to the pages that should have that power, protect writable paths in a reverse proxy, and set read_only: true on endpoints that should not write. A less obvious path such as /mcp/rw is not authentication.

Connecting a client

For the default stdio transport, configure the client to launch the binary. A configuration with one endpoint needs no arguments:

jsonc
{  "mcp": {    "kafka-local": {      "type": "local",      "command": ["kafka-mcp"],      "environment": {"CONFIG_FILE": "/path/to/kafka-mcp.yaml"}    }  }}

When the file contains several endpoints, add the endpoint selection to the command, for example "command": ["kafka-mcp", "--endpoint", "prod-read"]. The process refuses to start without it so it cannot silently connect a session to the wrong cluster or permission policy.

To use HTTP, start kafka-mcp --server and register one remote client entry per endpoint:

jsonc
{  "mcp": {    "kafka-local": {      "type": "remote",      "url": "http://localhost:8090/mcp/local"    },    "kafka-prod-read": {      "type": "remote",      "url": "http://localhost:8090/mcp",      "enabled": false    }  }}

The key becomes the tool prefix, so these appear as kafka-local_list_topics and kafka-prod-read_list_topics. Name entries after both the cluster and the endpoint policy: the prefix is the clearest signal of what a call can hit and change. Use server_config to confirm rather than trusting the client-side name.

Each enabled cluster costs context: these tools are roughly 8k tokens of definitions. Enable only what you need, and put production behind an agent:

jsonc
{  "tools": { "kafka-prod*": false },  "agent": {    "kafka-prod": {      "description": "Debugging against production Kafka.",      "tools": { "kafka-prod*": true }    }  }}

Permissions

Three layers, and only one of them is real security:

LayerProtects againstReal security?
confirm: true on writesAn LLM changing things on one ambiguous requestNo — a guardrail
endpoint read_only: trueAccidental writes with the server's Kafka identityNo — anyone who can edit the config can turn it off
Kafka ACLs on the SASL principalAn unauthorised personYes — the broker decides

What a read-only endpoint exposes

read_only: true does more than refuse a write: the endpoint does not list the tools whose only purpose is to write. add_partitions, alter_topic_config, commit_offset, create_topic, delete_consumer_group, delete_records and delete_topic are absent from tools/list on a read-only endpoint, so a client never sees a tool it could not have used, and their preview cannot describe a change this endpoint would never apply.

A writable endpoint can withhold individual tools too, with its tools map. That is the same mechanism seen from the client: the tool is not registered, so it is absent from tools/list and from server_config.

copy_message and produce_message stay, because read_only protects the cluster being written to and the destination is chosen per call. Copying a message out of a read-only production cluster, or seeding a writable preprod cluster from a protected session, is exactly what they are for. Both refuse outright when the destination is the read-only cluster itself.

Hiding a tool decides what is advertised, not what is permitted: both tools still refuse at the point of mutation, so a registration mistake cannot turn into a write. server_config reports the tools the endpoint actually exposes, which is how a session can tell the two cases apart.

Giving two people different permissions

Kafka enforces permissions against the SASL principal, so two people running the same server with different credentials get different rights. Adding partitions requires ALTER on the topic:

sh
# ali may read but not reshape topicsrpk acl create --allow-principal User:ali \  --operation read,describe --topic orders
# deniz may also add partitionsrpk acl create --allow-principal User:deniz \  --operation read,describe,alter --topic orders

Each points CONFIG_FILE at their own file, differing only in security.sasl[0].scram.user and the password. When ali calls add_partitions, the broker refuses:

not authorized to add partitions to "orders": the broker refused this request.Adding partitions requires ALTER permission on the topic for the principalthis server connects as

ali cannot bypass that by editing config or rebuilding the binary, because the decision is made by Kafka rather than by this server. On a cluster without ACLs, read_only: true is the available protection.

Audit logging

Every tool call is logged. One slog record per call, on the server's own log stream, so nothing extra has to be configured:

json
{"level":"INFO","msg":"tool call","tool":"produce_message","outcome":"ok", "duration_ms":12,"endpoint":"prod-write","cluster":"prod","read_only":false, "principal":"kafka-mcp-rw","session":"QY7MZM...","client":"claude-code", "client_version":"1.0.0","confirm":true,"item_count":1,"targets":"orders"}

The tools that change a cluster — add_partitions, alter_topic_config, commit_offset, create_topic, copy_message, delete_consumer_group, delete_records, delete_topic, produce_message — are logged at INFO. Everything else is logged at DEBUG, because reads are constant and change nothing, so recording them at the same level would bury the writes among them. Raise the log level to see them.

Two fields carry most of the weight. confirm separates a real write from a preview, and targets names the topic, group, partition and offset each item pointed at, so a record says which topic was touched rather than only that some topic was.

Message content is never logged. produce_message and copy_message carry arbitrary payloads, and an audit log is usually readable by more people than the data it describes, so keys, values and headers are left out. The audit code has no field to unmarshal them into, so content cannot reach a log even by mistake.

None of the recorded identities is a person, and the server cannot make one up. What each actually means:

FieldWhat it is
principalThe Kafka credential this server connects as, which is what ACLs are enforced against. Everyone reaching the same endpoint shares it
endpointWhich endpoint policy allowed the call, and so which cluster was touched
sessionOne MCP session, which groups a sequence of calls into one investigation
client, client_versionThe program that connected, as it identified itself at initialize. Self-reported and not verified
userPresent only when an inbound bearer token established it. This server installs no token verifier, so it is absent today
request_idThe HTTP request id, which joins a record to the access log

For attribution to a person, give each person their own endpoint and SASL credentials, as above: the principal in the record is then the answer to who acted. A shared credential cannot be made to answer it.

Tools

Batch operations

describe_topic, sample_messages, get_message, get_schema, consumer_lag, describe_consumer_group, open_transactions, add_partitions, alter_topic_config, create_topic, commit_offset, delete_consumer_group, delete_records, delete_topic, copy_message and produce_message take their target only as a required items array. There is no single-target form: one operation is an items array of length one.

json
{"items": [{"topic": "orders"}]}{"items": [{"topic": "orders"}, {"topic": "payments"}]}

Everything naming or shaping an operation lives on the item, so each field and its description exist in exactly one place. What governs the whole call stays at the top level, which in practice means confirm.

The message-heavy tools accept at most 20 items; lag, schema and administrative tools accept at most 100.

Every response is the same envelope. Results stay in input order, each entry carrying index and either result or error, followed by succeeded, failed, applied and atomic: false. An item's failure is data in the response rather than an error for the call, so it never hides the items that worked. Batch writes preview every item first, then apply the valid ones only when the top-level confirm is true. They are not transactions: Kafka cannot roll back a topic, partition, offset or produced message after a later item fails. Duplicate write targets are refused before anything changes, except in produce_message, where two identical items mean two messages rather than a mistake.

list_clusters

Lists the clusters this server serves, with whether each is reachable and whether it accepts writes. Takes no parameters. Available from every endpoint, so a session can discover what copy_message, produce_message and compare_clusters may target.

json
{"clusters": [  {"name": "prod", "connected": true, "read_only": true},  {"name": "preprod", "connected": true, "read_only": false}], "count": 2}

connected is checked when you call, not recorded at startup, so a cluster that has since gone down is reported honestly. Only the name, reachability and writability are reported: broker addresses and credentials are deliberately not, because this tool is reachable from every endpoint.

compare_clusters

Compares the topics of other clusters against the one this endpoint serves, and reports what differs. Use it to find what preproduction has that production does not, or to check whether two environments still match.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 clusters to compare against

Item fields:

FieldTypeRequiredMeaning
clusterstringyesThe other cluster. Use list_clusters for valid names
searchstringnoCase-insensitive substring a topic name must contain
include_internalboolnoDefault false: internal topics are excluded
json
{"here": "preprod", "there": "prod", "here_cluster": {"name": "preprod", "brokers": 1, "topics": 12}, "there_cluster": {"name": "prod", "brokers": 3, "topics": 11}, "only_here": [{"topic": "orders-v2", "partitions": 6, "replication_factor": 1,                "configs": {"retention.ms": "604800000"}}], "only_there": [], "differing": [{"topic": "orders", "differences": ["partitions"],                "here": {"partitions": 1}, "there": {"partitions": 12}}], "in_both": 11}

only_here and only_there name the direction, which is decided by the endpoint you call: "here" is always the cluster this endpoint serves. Entries carry the partition count, replication factor and explicitly-set configs of the cluster that has the topic, so they can be passed straight to create_topic.

This tool creates and changes nothing. To create the missing topics, hand the chosen entries to create_topic, which previews them against the broker first and warns that a partition count can never be reduced.

differing is usually the more valuable half: a topic that exists on both sides with a different partition count or retention.ms is the common reason a bug reproduces in one environment and not the other. Only configs a topic sets for itself are compared, because two clusters may carry different broker defaults and comparing inherited values would report every topic as different.

A difference is not necessarily a mistake. A topic missing from production is often deliberate, so the report says what differs, never what is correct.

Topic listings come from the Kafka client's metadata cache, which is a few seconds old, so a topic created moments earlier may still appear in only_there. Repeat the comparison rather than creating it twice.

list_topics

Lists the topics on the cluster, sorted by name, each with its partition count, replication factor and the configs it sets for itself.

ParameterTypeRequiredMeaning
scriptstringnoJavaScript predicate deciding whether a topic is listed
timeout_secondsintnoLimit for evaluating the script. Default 30

The predicate sees:

VariableTypeMeaning
topicstringThe topic name
partitionsnumberPartition count
replication_factornumberReplicas of the first partition
internalbooleanKafka's own topics, such as __consumer_offsets
configsobjectValues this topic sets for itself, e.g. configs['retention.ms']
js
return topic.indexOf('orders') >= 0return partitions > 6return configs['cleanup.policy'] === 'compact'return replication_factor === 1 && !internal
json
{"name": "list_topics", "arguments": {"script": "return partitions > 6"}}
json
{"topics": [{"topic": "orders", "partitions": 12, "replication_factor": 3,             "configs": {"retention.ms": "604800000"}}], "count": 1}

Filtering is JavaScript only, as it is for search_messages: a name match is return topic.indexOf('orders') >= 0. The predicate can also answer what a substring never could — which topics have more than six partitions, only one replica, or a compacted cleanup policy.

Only configs a topic sets for itself are reported. Inherited cluster defaults are excluded, because including them would make every topic look configured.

A topic whose predicate throws, or is cut short by the timeout, is counted in script_errors rather than listed, so a broken filter is never mistaken for an empty cluster.

describe_topic

Reports a topic's partitions, offset ranges, message count, time span and full configuration. Use it before searching to see how much data a search would read and how far back the topic can hold data at all.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 20 topics to describe

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to describe
json
{"results": [{"index": 0, "result": {   "topic": "orders", "partition_count": 1, "message_count": 3,   "partitions": [{"partition": 0, "start_offset": 0, "end_offset": 3, "message_count": 3}],   "configs": [{"key": "cleanup.policy", "value": "delete", "source": "DYNAMIC_TOPIC_CONFIG", "is_default": false},               {"key": "retention.ms", "value": "604800000", "source": "DEFAULT_CONFIG", "is_default": true}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

configs lists every topic config as the string Kafka reports, where -1 means unlimited. is_default is true when the value is inherited rather than set on the topic. Two entries decide whether a message can still exist at all: retention.ms (how long messages are kept) and cleanup.policy (compact keeps only the latest message per key).

sample_messages

Reads a small sample of the newest messages and reports what they look like: value formats, field paths with their types, key statistics, which value fields carry the message key, and the schemas the values were written with. Avro, Protobuf and JSON Schema values are decoded first, so their field paths are reported like JSON. Use it before searching to decide how to search, and before producing to find the schema to write against.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 20 topics to sample

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to sample
sample_sizeintnoMessages to read in total. Default 20
partitionsint[]noRestrict to these partitions
max_value_bytesintnoValue bytes per message. Default 512
json
{"results": [{"index": 0, "result": {   "value_formats": {"json": 20, "text": 0, "binary": 0},   "json_fields": [{"path": "payload.amount", "types": ["number"], "present": 20, "example": "500"}],   "key_stats": {"present": 20, "absent": 0, "unique": 20, "all_unique": true},   "key_in_value": ["payload.orderId"],   "schemas": [{"format": "avro", "schema_id": 7, "message_type": "shop.Order", "count": 20}],   "sampled_ranges": [{"partition": 0, "start": 980, "end": 1000}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

key_in_value naming a field means the key is that identifier, so searching the key alone is the precise, cheap lookup. value_formats also counts avro, protobuf, json_schema, msgpack, null and undecodable — values that named a schema but could not be decoded, which is a configuration problem, not binary data.

search_messages

Scans a bounded range of a topic, filtering messages with a JavaScript expression. Kafka has no server-side search, so this reads messages and filters them client-side; the result reports what was covered.

ParameterTypeRequiredMeaning
topicstringyesTopic to search
scriptstringnoJavaScript filter. Omit to match every message
parallelismintnoConcurrent readers, 1–16. Default 1
partitionsint[]noRestrict to these partitions
from_offset / to_offsetintnoOffset window, end exclusive
from_timestamp / to_timestampstringnoRFC3339 time window
directionstringnonewest_first (default) or oldest_first
max_matchesintnoStop after this many matches. Default 10
max_messages_scannedintnoRead at most this many. Default 10000
max_value_bytesintnoValue bytes per match. Default 512
timeout_secondsintnoWall-clock limit. Default 30
count_onlyboolnoReturn counts only, no message bodies
output_filestringnoWrite every match to this file as JSONL
The script

Return true to keep a message. In scope:

NameValue
valueparsed JSON document; the decoded record for Avro, Protobuf, JSON Schema or a configured format; the raw text otherwise
keystring, the decoded document when the key has a schema, or null when absent
headersobject of header name to string
partition, offsetnumbers
timestampa Date
js
return key === 'order-123'return value.eventType === 'NEW' && value.payload.amount >= 500return value.payload.cancelledAt === null      // present and nullreturn value.payload.cancelledAt === undefined // field absentreturn headers['correlation-id'] === 'corr-999'return /ORD-\d{4}/.test(value.payload.orderId)return value.indexOf('ERROR') >= 0             // non-JSON topic: value is a string

Searching by key exactly is far more precise than searching the body: 123 also appears inside "amount": 1123, and those false positives can fill max_matches and hide the message wanted.

A script that throws on a message is counted in script_errors and the scan continues, so a broken script is distinguishable from a genuine absence of matches. Scripts run in a sandbox with no filesystem, network or host access, and are stopped if they exceed the search timeout. They are not bounded by memory: something like 'x'.repeat(1e12) can exhaust the server process. Scripts are trusted input; the blast radius is this server, not the cluster.

Scan order

Every partition is read together, one chunk deep at a time: the newest chunk of every partition, then the chunk behind it, and so on. A limited newest-first search therefore returns the newest matches in the topic, not the newest in whichever partition happened to be read first.

Kafka orders records within a partition and never across them, so matches are merged and reported by timestamp, with partition and offset breaking ties. Timestamps are set by the producer unless the topic uses LogAppendTime, so they can be skewed; it is still the only thing comparable between partitions.

One scan reads every partition through a single connection, so a wide topic costs no more connections than a narrow one. It does read more: a limited search on a 12-partition topic examines the newest chunk of all twelve rather than stopping inside the first. max_messages_scanned still bounds it, and may become the stopped_reason on a wide topic sooner than on a narrow one.

Parallelism

parallelism splits a single-partition topic's offset range between that many readers, which is what makes a full scan of one large partition fast. A multi-partition topic is already read in parallel across its partitions, so the setting does not apply there, and a range too small to divide is read by one reader.

It pays off for count_only, output_file and full scans of a single partition.

Result
json
{"topic": "orders", "match_count": 1, "matches": [{"partition": 0, "offset": 17, "key": "order-42", "value": "...", "encoding": "utf8"}], "scanned_messages": 120, "scanned_ranges": [{"partition": 0, "start": 0, "end": 120}], "stopped_reason": "range_exhausted", "complete": true}

complete is true only when the whole range was read. An empty match list means "not there" only if complete is true; otherwise check stopped_reason and narrow the search.

For large result sets, use count_only to learn how many matches exist, then output_file to write them out instead of returning them.

get_message

Reads messages at exact offsets, plus optional neighbours.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 20 addresses to read

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to read from
partitionintyesPartition to read from
offsetintyesExact offset to read
contextintnoAlso return this many messages either side
max_value_bytesintnoValue bytes to return. Default 4096
json
{"name": "get_message", "arguments": {"items": [  {"topic": "orders", "partition": 0, "offset": 17, "context": 1}]}}

Schema-encoded values are decoded to JSON, with format, schema_id and message_type saying what they were. Values that could not be decoded and are not valid UTF-8 are base64 encoded, with encoding set to base64; decode_error explains a value that named a schema but could not be decoded.

json
{"partition": 0, "offset": 17, "key": "o-1", "value": "{\"amount\":42,\"id\":\"o-1\"}", "format": "avro", "schema_id": 7, "message_type": "shop.Order", "encoding": "utf8", "value_bytes": 12}

get_schema

Reads schemas from the cluster's Schema Registry, by subject or by the schema_id a message carries. Read the schema before producing to a schema-encoded topic: it names every field and enum a value needs, which a sampled message may not show.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 schemas to look up

Item fields, giving subject or id:

FieldTypeRequiredMeaning
subjectstringeitherSubject, usually <topic>-value or <topic>-key
versionintnoSubject version. Defaults to the latest
idinteitherSchema id, e.g. from get_message
json
{"results": [{"index": 0, "result": {   "schema_id": 7, "subject": "orders-value", "version": 3, "versions": [1, 2, 3],   "type": "avro", "schema": "{...}", "references": [],   "message_types": null, "used_by": null}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

message_types lists Protobuf messages, which produce_message takes as message_type. Looking up an id fills used_by with the subject versions that use it. The call fails as a whole when the cluster has no schema_registry.

list_consumer_groups

Lists consumer groups with their state, member count and the topics they consume.

ParameterTypeRequiredMeaning
topicstringnoOnly groups consuming or committed to this topic
statesstring[]noFilter by state, e.g. Stable, Empty
json
{"groups": [{"group": "payments", "state": "Stable", "members": 2, "topics": ["orders"]}], "count": 1}

A group in state Empty can still report lag: committed offsets outlive the consumers that made them. Kafka has no topic-to-group index, so filtering by topic describes every group on the cluster.

describe_consumer_group

Describes consumer groups in detail: who the members are, which partitions each one owns, and where the group stands on every partition. Use it to turn "a partition is stuck" into "this pod on this host is stuck".

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 groups, each {"group": "..."}
json
{"results": [{"index": 0, "result": {   "group": "payments", "state": "Stable", "protocol_type": "consumer", "assignor": "cooperative-sticky",   "coordinator": 1, "total_lag": 4200,   "members": [{"member_id": "payments-7-…", "client_id": "payments-7", "host": "/10.0.4.17",                "assignments": [{"topic": "orders", "partitions": [0, 1]}]}],   "partitions": [{"topic": "orders", "partition": 0, "has_commit": true, "committed_offset": 812,                   "end_offset": 5012, "lag": 4200, "member_id": "payments-7-…",                   "client_id": "payments-7", "host": "/10.0.4.17"}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

Partitions are the union of what members own and what the group has committed, so an Empty group still shows its positions. has_commit: false means the group owns a partition it has never committed on, so where it starts is decided by the consumer's auto.offset.reset, not by an offset.

open_transactions

Finds open transactions holding back read_committed consumers. A transactional producer that hangs or dies mid-transaction leaves the partition's last stable offset stuck, and every read_committed consumer stops there. In consumer_lag that looks exactly like a poison message.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 topics, each {"topic": "..."}
json
{"results": [{"index": 0, "result": {   "topic": "orders", "blocked": true,   "partitions": [{"partition": 0, "last_stable_offset": 812, "high_watermark": 5012,                   "unreadable_messages": 4200,                   "producers": [{"producer_id": 2004, "producer_epoch": 0, "transaction_start_offset": 812,                                  "transactional_id": "payments-writer-1", "state": "Ongoing",                                  "started_at": "2026-10-02T08:14:03Z", "open_for": "41m12s", "timeout_ms": 900000}]}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

The fix is in the producer: restart or fence the application named by transactional_id, or wait for timeout_ms, after which the broker aborts the transaction. Moving the consumer's offset does not help.

consumer_lag

Measures how far behind a topic's consumers are, how fast messages are produced and consumed, and when the backlog will clear.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 measurements to take

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to measure
groupstringnoDefaults to every group consuming the topic
sample_secondsintnoConsume-rate sample window. Default 5. The call blocks
skip_consume_rateboolnoReturn immediately, without a rate or estimate

Sampling windows run concurrently, so several measurements do not add their waits together.

json
{"results": [{"index": 0, "result": {   "topic": "orders", "total_lag": 4200,   "produce_rate": {"last_minute": {"messages": 3000, "per_second": 50, "per_minute": 3000, "per_hour": 180000}},   "groups": [{"group": "payments", "state": "Stable", "members": 2, "lag": 4200,               "consume_rate": {"per_second": 120, "sampled_seconds": 5},               "drain_per_second": 70, "eta_seconds": 60, "eta_human": "1m 0s",               "status": "draining"}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

The two rates are measured differently, and the output says so:

  • produce_rate — historical fact, from message timestamps, over the last second, minute and hour. window_truncated marks a topic younger than the window.
  • consume_rate — a sample: the committed offset is read, then read again sample_seconds later. sample_inconclusive means nothing moved.

The backlog drains at the consume rate minus the produce rate. status says what the numbers mean: caught_up, draining (with an ETA), growing (never clears, with growing_by_per_minute), stalled, no_active_consumers, or not_measured. An ETA is only given when the lag is genuinely shrinking.

cluster_health

Checks the cluster this endpoint serves in one call: brokers, controller, and every partition that is not fully healthy.

ParameterTypeRequiredMeaning
searchstringnoOnly topics whose name contains this. Case-insensitive
include_internalboolnoAlso check __consumer_offsets and other internal topics
json
{"cluster_id": "…", "controller": 1, "healthy": false, "brokers": [{"id": 1, "host": "kafka-1", "port": 9092, "rack": "eu-1a", "controller": true, "leaders": 61}], "summary": {"topics": 40, "partitions": 182, "offline": 0, "under_replicated": 3, "under_min_isr": 1, "errored": 0}, "problems": [{"topic": "orders", "partition": 4, "issues": ["under_replicated", "under_min_isr"],               "leader": 1, "replicas": [1, 2, 3], "isr": [1], "min_insync_replicas": 2}], "warnings": []}

offline means the partition has no leader, so nothing can be read or written. under_replicated means a replica is out of sync. under_min_isr is the condition behind NOT_ENOUGH_REPLICAS: producers using acks=all fail until the in-sync replicas recover. A broker leading no partitions while others lead many is usually one that restarted and was never given leadership back.

Some Kafka-compatible brokers, Redpanda among them, do not report min.insync.replicas. The topics affected are listed in min_isr_unknown, and under_min_isr is not judged for them rather than guessed.

list_acls

Lists access control entries, for when a client fails with TOPIC_AUTHORIZATION_FAILED or GROUP_AUTHORIZATION_FAILED.

ParameterTypeRequiredMeaning
principalstringnoOnly this principal, e.g. User:payments
resource_typestringnotopic, group, cluster, transactional_id, delegation_token
resource_namestringnoEvery ACL applied to this name, including prefixed and *. Needs resource_type
json
{"acls": [{"principal": "User:payments", "host": "*", "resource_type": "topic", "resource_name": "orders",           "pattern_type": "literal", "operation": "read", "permission": "allow"}], "count": 1}

A deny overrides every allow. A broker without an authorizer fails with SECURITY_DISABLED, which means ACLs are not enforced at all.

server_config

Reports the effective configuration: endpoint name, exact path, description, cluster, brokers, authentication mechanism and principal, TLS, read-only state, export directory and the tools this endpoint exposes. Takes no parameters. The password is never reported.

tools is the list for this endpoint, not for the deployment: a read-only endpoint omits add_partitions, alter_topic_config, commit_offset, create_topic, delete_consumer_group, delete_records and delete_topic, because it does not register them, and any endpoint omits whatever its tools configuration switches off.

Use it when a result is surprising: an empty topic list means something very different on a local broker than on production.

add_partitions

Adds partitions to a topic. Irreversible — Kafka cannot reduce a partition count. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 topics to change
confirmboolnoDefault false: preview only, nothing changes

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to change
partitionsintyesFinal total, not the number to add. Repeating a call is safe
acknowledge_key_orderingboolnoRequired when messages are keyed
sample_sizeintnoMessages inspected for keys. Default 20

Without confirm it reports what would happen: current and target counts, whether messages are keyed, which consumer groups will rebalance, and warnings.

Adding partitions changes which partition a key hashes to, so existing keys lose their ordering guarantee. A keyed topic therefore requires acknowledge_key_ordering as well. Requesting fewer partitions than the topic has is refused with an explanation rather than attempted.

alter_topic_config

Changes topic-level configuration: retention, cleanup policy, maximum message size and anything else Kafka allows per topic. Changes are incremental, so every key not named keeps its value. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 topics to change
confirmboolnoDefault false: preview only, nothing changes

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to change
setobjectnoKeys to set, e.g. {"retention.ms": "86400000"}
deletestring[]noOverrides to remove, so the cluster default applies again
json
{"results": [{"index": 0, "result": {   "topic": "orders", "applied": false, "messages_past_retention": 18000,   "changes": [{"key": "retention.ms", "action": "set", "current": "604800000",                "current_source": "DYNAMIC_TOPIC_CONFIG", "requested": "86400000"}],   "warnings": ["18000 message(s) are already older than the new retention of 24h0m0s and become eligible for deletion as soon as it applies; Kafka cannot bring them back"]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

The preview asks the broker to validate the change, so an unknown key or an invalid value is refused before confirm. Shortening retention.ms reports how many messages are already past the new limit; changing cleanup.policy or setting retention.bytes is warned about. After applying, every key is re-read from the broker, which is the only way the inherited value of a deleted override is known.

create_topic

Creates a topic. Refuses a topic that already exists rather than adjusting it. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 topics to create
confirmboolnoDefault false: the broker validates the requests and creates nothing

Item fields:

FieldTypeRequiredMeaning
topicstringyesName of the topic to create
partitionsintnoOmit for the broker default on Kafka 2.4+. Can grow later, never shrink
replication_factorintnoOmit for the broker default on Kafka 2.4+. Cannot exceed the broker count
configsmapnoTopic-level config, such as retention.ms or cleanup.policy

Without confirm the request is sent to the broker with ValidateOnly, so the preview reports the cluster's own answer — an invalid name, an unknown config key, a replication factor larger than the cluster — rather than a guess. With confirm the topic is created and the resulting partition count and replication factor are read back from the cluster, which is how an omitted count is reported as the number the broker actually chose.

Using broker defaults requires Kafka 2.4 or newer, whose CreateTopics v4 API introduced -1 as "use the broker default". On an older broker, pass both counts explicitly.

A replication factor larger than the number of brokers is refused here, with the broker count in the message, because brokers differ on whether a validate-only request catches it.

Use add_partitions to change an existing topic's partition count; this tool never modifies a topic it did not create.

delete_topic

Deletes topics. The most destructive tool here: deleting a topic destroys every message in it, Kafka has no undo, and any consumer group reading it breaks. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 topics to delete
confirmboolnoDefault false: preview only, nothing is deleted

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to delete. It must exist
acknowledge_data_lossboolnoRequired when the topic still holds messages

Without confirm it reports, per topic, how many messages would be destroyed, how many partitions it had, and which consumer groups had committed offsets for it:

json
{"results": [{"index": 0, "result": {   "topic": "orders-old", "partitions": 6, "message_count": 41207,   "consumer_groups": ["payments"], "deleted": false, "would_delete": true,   "warnings": ["41207 message(s) would be destroyed, and Kafka cannot restore them: the only recovery is a backup taken beforehand"]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

Two separate acknowledgements are required, because they answer different questions. confirm says the caller meant to delete; acknowledge_data_loss says they know what is inside. A topic holding messages is refused without both, and the count is re-read at deletion time, so a topic that gained messages since the preview is still caught.

Internal topics such as __consumer_offsets are refused outright at any level of acknowledgement: they hold cluster state rather than a caller's data, and deleting one breaks every consumer at once.

Duplicate topics in one batch are refused before anything is deleted. Deletion is not atomic — topics removed before a later item failed stay removed.

delete_records

Deletes the oldest messages of a partition without deleting the topic. The topic, its configuration and its consumer groups stay. Irreversible. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 partitions to truncate
confirmboolnoDefault false: preview only, nothing is deleted

Item fields:

FieldTypeRequiredMeaning
topicstringyesTopic to delete from
partitionintyesPartition to delete from
before_offsetintyesEverything below this offset goes; this one becomes the first
acknowledge_data_lossboolnoRequired to apply
json
{"results": [{"index": 0, "result": {   "topic": "orders", "partition": 0, "start_offset": 0, "end_offset": 5012, "before_offset": 812,   "messages_deleted": 812, "would_delete": true, "deleted": false,   "affected_groups": [{"group": "replay-job", "committed_offset": 100, "unprocessed_lost": 712}]}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

affected_groups lists every group committed below the cut, with how many messages it would lose without ever processing them; such a group resumes from the new start. Use the partition's end offset as before_offset to empty it.

delete_consumer_group

Deletes consumer groups and their committed offsets. Use it to clean up groups whose consumers were decommissioned: their commits keep reporting lag that nobody will ever drain. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 groups, each {"group": "..."}
confirmboolnoDefault false: preview only, nothing is deleted

The preview lists each group's committed offsets and lag. A group with active members is refused, and re-checked at deletion time. A consumer that later starts with the same group id begins from its auto.offset.reset, not from where the group left off.

commit_offset

Moves consumer groups' committed offsets. Forward to skip messages, backward to replay them, or to a point in time to reprocess everything since. Irreversible in the sense that skipped messages are never processed. Not exposed on a read-only endpoint.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 100 moves to make
confirmboolnoDefault false: preview only, nothing changes

Item fields — give exactly one of offset, timestamp or position:

FieldTypeRequiredMeaning
topicstringyesTopic whose offset is moving
groupstringyesConsumer group to move
partitionintnoRequired with offset. Omit with timestamp or position for every partition
offsetint—The offset the group reads next. To skip offset 42, commit 43
timestampstring—RFC3339. Each partition moves to its first message at or after this time
positionstring—earliest or latest
allow_active_membersboolnoProceed despite running consumers
json
{"items": [{"topic": "orders", "group": "payments", "timestamp": "2026-10-01T09:00:00Z"}]}
json
{"results": [{"index": 0, "result": {   "topic": "orders", "group": "payments", "state": "Empty", "members": 0,   "partitions": [{"partition": 0, "current_offset": 5012, "target_offset": 4100, "replayed_messages": 912},                  {"partition": 1, "current_offset": 4990, "target_offset": 4021, "replayed_messages": 969}],   "replayed_messages": 1881, "applied": false}}], "succeeded": 1, "failed": 0, "applied": 0, "atomic": false}

A partition with no message at or after timestamp moves to its end, and its entry carries a note saying so. A whole-topic item and a partition item for the same group and topic in one batch are refused, because which one wins would depend on order.

The group must have no active members. A running consumer keeps its position in memory and only reads the committed offset when it joins, so a commit made while it runs is overwritten by its next commit and the group does not move. Stop the consumers first.

copy_message

Copies messages to another topic, preserving key, value and headers. Takes each message's address, never its content, so it can only duplicate a message the cluster already holds.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 20 copies to make
confirmboolnoDefault false: preview only, nothing is written

Item fields:

FieldTypeRequiredMeaning
source_topic, source_partition, source_offsetyesMessage to copy
destination_topicstringyesWhere to write it. Must already exist
destination_clusterstringnoAnother cluster to write to. Defaults to this endpoint's own
translate_schemaboolnoRe-register a schema id in the destination's registry
max_value_bytesintnoPreview value limit. The whole value is always copied

Set destination_cluster to copy into another cluster this server serves, which is how a production message is taken into a preproduction topic to be debugged safely. Use list_clusters to see which names are valid.

Every copy carries provenance headers — kafka-mcp-copied-from-cluster, -from-topic, -from-partition, -from-offset, -copied-at, -copied-by-tool, -copied-by-principal — so a message in a dead letter or preproduction topic can be traced back to its original. If the message already carries one of those headers, the original is kept and the collision is reported.

Key and value bytes are copied unchanged. A schema id is only meaningful in the registry that issued it, so when the destination cluster uses a different registry, or none, the response warns. With translate_schema the schema is registered in the destination registry under <destination_topic>-value (or -key), with its references, and the id in the copy is rewritten; the payload is untouched. That registration is a write to the destination registry and happens only with confirm.

For a copy within the endpoint's own cluster, its read_only policy protects the destination. A read-only endpoint can still be the source of a cross-cluster copy, because copying out changes nothing there. A different destination cluster is writable when it has at least one writable endpoint; list_clusters reports that effective state. The tool is refused entirely when the destination is read-only, preview included, because writing is all it does.

produce_message

Writes new messages to existing topics. Unlike copy_message, the caller supplies the content, so this can put a message into a topic that no producer ever sent.

ParameterTypeRequiredMeaning
itemsobject[]yes1 to 20 messages to write
confirmboolnoDefault false: preview only, nothing is written

Item fields:

FieldTypeRequiredMeaning
topicstringyesExisting topic to write to
valuestringyesThe message body
keystringnoDecides the partition when partition is omitted
headersobjectnoHeader name to value
partitionintnoExact partition. Omit to let the key decide
encodingstringnoutf8 (default) or base64 for exact bytes you already have
value_schemaobjectnoEncode value (JSON) to a registry schema
key_schemaobjectnoEncode key (JSON) to a registry schema
destination_clusterstringnoAnother cluster to write to. Defaults to this endpoint's own
max_value_bytesintnoPreview value limit. The whole value is always written

Every message carries kafka-mcp-produced-at, kafka-mcp-produced-by-tool, kafka-mcp-produced-by-principal and, when the client identifies itself, kafka-mcp-produced-by-client, so a fabricated message stays distinguishable from a genuine one. A header the caller supplies under one of those names is kept as given and the collision is reported.

value_schema and key_schema take {subject, version, id, message_type}, all optional: {} means the latest version of <topic>-value (or -key) in the destination cluster's registry, id pins an exact schema, and message_type picks a Protobuf message. The value is given as JSON, validated against the schema in the preview — a missing or misspelt field is refused by name — and written framed with the schema id, as registry-aware consumers expect. A topic with a format in topic_formats is encoded to it without being asked. Neither can be combined with encoding: base64. The response reports the schema used in value_encoding / key_encoding.

The topic must already exist: a missing one is refused rather than left to auto-creation. Omit partition unless the exact partition is the point — the key decides placement, and naming a partition puts a keyed message where its key does not hash to, which breaks ordering for that key. The response warns whenever an explicit partition is used.

read_only protects the cluster being written to, so a read-only endpoint may still produce into a different, writable cluster, and is refused outright — preview included — when writing to its own. A produced message cannot be deleted; it stays until retention removes it.

Skills

skills/kafka-debugging/SKILL.md is the one skill an agent loads. It routes to the scenario guides under skills/kafka-debugging/references/, rather than holding every workflow itself, so a session reads only the one it needs:

  • find-message.md — locating a message from something the user knows about it.
  • check-lag.md — measuring lag and throughput, and judging when a backlog will clear.
  • scale-partitions.md — deciding whether more partitions will help, and adding them safely.
  • skip-poison-message.md — unblocking a consumer stuck on a message it cannot process, preserving the message first.
  • create-topic.md — creating a topic with a partition count and retention chosen on purpose, including as a copy_message destination.
  • produce-message.md — writing a message: repairing and re-injecting one, reproducing a failure in another cluster, or seeding a topic.
  • compare-clusters.md — finding what differs between two environments, and creating the topics one of them is missing.
  • delete-topic.md — removing a topic and everything in it, after the user has seen what that destroys.
  • replay-messages.md — reprocessing everything since a moment, after a bug fix ships.
  • tune-topic-config.md — changing retention, cleanup policy or message size on a topic, knowing what the change does to data already there.
  • cluster-health.md — finding offline and under-replicated partitions behind NOT_ENOUGH_REPLICAS and unreadable topics.
  • purge-messages.md — deleting old messages from a partition while keeping the topic, or removing abandoned consumer groups.
  • authorization-error.md — working out which ACL is refusing a client.

The umbrella also resolves the overlap between them: "the consumer is behind" opens three of these guides, and consumer_lag's status is what decides which one is right.

Development

CommandPurpose
make up / make downStart / stop Redpanda and Console
make buildBuild with goreleaser into dist/
make runRun the HTTP server from source
go test ./...All tests, including container tests
go test -short ./...Tests that need no containers
go vet ./...Vet all packages

The documentation site lives in _docs (Vite, pnpm). pnpm install && pnpm dev there serves it locally; pushing changes under _docs/ to main publishes it to GitHub Pages through .github/workflows/docs.yml.

Tests run against real containers started by internal/domain/testenv (a Redpanda broker plus Console), so Docker must be available for the full suite.

Start the HTTP server

Requires Go 1.27 or later. Configuration is loaded with chu; set CONFIG_FILE to select a YAML or JSON file. into manages the process lifecycle, ada serves HTTP with context-driven shutdown, and logi initializes structured logging.

sh
make env-up                                  # local Redpanda + Consolemake run                                     # runs with --server and serves kafka-mcp.local.yaml

kafka-mcp.local.yaml is committed and points at the compose broker, so a clone works without writing any configuration. make env-up publishes the broker on localhost:19092, the Schema Registry on localhost:18081 and the Redpanda Console on http://localhost:8080.

Check it is up:

sh
curl http://localhost:8090/healthz

Source: README.md at commit d2a70c8

Tools

0
Tool metadata has not been indexed yet.

Version history

1
  1. v0.2.1LatestOct 7, 2026