Skip to main content

Python SDK — CrewAI Integration

KafkaMCP provides a CrewAI integration at sdk/python/kafkamcp_crewai for building Kafka-aware AI agent crews.

Installation

pip install "git+https://github.com/josedab/kafkamcp.git#subdirectory=sdk/python/kafkamcp_crewai"

The package is not yet published to PyPI.

Quick Start

from kafkamcp_crewai import StdioMCPClient, create_kafka_tools
from crewai import Agent, Task, Crew

with StdioMCPClient(args=["--config", "kafkamcp.yaml"]) as client:
kafka_tools = create_kafka_tools(mcp_client=client)

kafka_agent = Agent(
role="Kafka Analyst",
goal="Analyze Kafka topics and consumer groups",
backstory="Expert in Apache Kafka operations and data analysis.",
tools=kafka_tools,
)

task = Task(
description="List all topics and report consumer group lag for each.",
expected_output="A topic and consumer lag report.",
agent=kafka_agent,
)

crew = Crew(agents=[kafka_agent], tasks=[task])
result = crew.kickoff()
print(result)

Available Tools

Tool ClassMCP ToolDescription
KafkaConsumeToolkafka_consumeRead messages from a topic
KafkaProduceToolkafka_producePublish messages to a topic
KafkaListTopicsToolkafka_list_topicsList all topics with metadata
KafkaDescribeTopicToolkafka_describe_topicGet detailed topic metadata
KafkaCreateTopicToolkafka_create_topicCreate a new topic
KafkaDeleteTopicToolkafka_delete_topicDelete a Kafka topic
KafkaSearchToolkafka_searchSearch messages by content
KafkaDescribeGroupToolkafka_describe_consumer_groupGet consumer group state and lag
KafkaListGroupsToolkafka_list_consumer_groupsList all consumer groups
KafkaGetSchemaToolkafka_get_schemaGet schema definition
KafkaAnalyzeDLQToolkafka_analyze_dlqAnalyze dead-letter queue failures

Factory Functions

FunctionDescription
create_kafka_tools(mcp_client)Create the eight read-only default wrappers
create_guarded_tools(mcp_client)Create three approval-gated mutation wrappers
create_monitoring_tools(mcp_client)Read-only monitoring subset
create_remediation_tools(mcp_client)DLQ remediation subset (consume, remediate, search, schema)

Individual Tool Usage

from kafkamcp_crewai import KafkaConsumeTool, KafkaSearchTool, StdioMCPClient

with StdioMCPClient(args=["--config", "kafkamcp.yaml"]) as client:
consume = KafkaConsumeTool(mcp_client=client)
search = KafkaSearchTool(mcp_client=client)
agent = Agent(role="Data Investigator", tools=[consume, search], ...)

See Also