developing-kafka-python-client
SkillCloud & infraUse when the user wants to build a Python Kafka producer or consumer, add Schema Registry to existing Python code, migrate from raw JSON to schema-backed serialization, or scaffold a confluent-kafka-python project for Confluent Cloud, local Docker, or WarpStream. Also use when user wants to optimize Python Kafka client configuration for WarpStream.
Available today. Use it from your connected AI after setup.
No other account needed.
Add ahel to your AI once: Claude, ChatGPT, Cursor, Claude Code or Codex. Then ask it to use this.
Then ask your AI: use the developing-kafka-python-client skill
What this skill tells your AI
The instructions your AI receives, as published by confluentinc/agent-skills in skills/developing-kafka-python-client/SKILL.md and read by ahel’s review.
Begin by announcing: "Using the Confluent Kafka Python Client skill to guide this project."
Confluent Kafka Python Client Creation
Generate a production-ready Python project for producing to and/or consuming from Kafka using confluent-kafka-python. Supports three target environments: Confluent Cloud (managed), Local Docker (open-source Kafka), and WarpStream (Kafka-compatible, object-storage-backed), and two producer styles: AsyncIO (non-blocking) and Synchronous (blocking). The generated code follows Confluent's best practices.
Step 1: Gather Requirements
Before generating any code, work through the questions below. Skip any question the user has already answered explicitly in their prompt — do not re-ask just for form's sake. For example, "build a producer and consumer on Confluent Cloud with an async producer" already answers #2, #3, and #4; only #1, #5, #6, #7, and #8 remain.
Mandatory confirmation gate — do not skip, even if the user answered every question. Before writing any file, you MUST send one message that:
- Recaps the answers you extracted as a short bulleted list (e.g., "Target: Confluent Cloud · Components: producer + consumer · Producer style: async · From scratch: yes").
- Asks any remaining open questions inline.
- Explicitly asks the user to confirm or correct before you proceed.
Then STOP and wait for the user's reply. Do not generate files in the same turn as the recap, and do not proceed on the assumption that a fully-specified prompt implies consent to generate immediately — the recap catches misinterpretations of the prompt and is required even when questions #1–#8 are all pre-answered. The only way to skip the gate is if the user has already confirmed the recap earlier in this conversation.
Do not assume defaults for #1, #2, or #3 — if any of these are not answered by the prompt, you must ask.
- Are you adding Kafka to an existing application, or starting from scratch?
- If the user has existing Python code (mentions an existing project, has a
main.py, uses Flask/FastAPI/Django, etc.), do not scaffold a new project. Instead: (a) identify their existing producer or data-sending code, (b) ask whether they already have schemas registered in Schema Registry, (c) add Schema Registry integration to their existing code following the patterns in the reference files. Generate only the files they are missing (e.g.,common.py,schemas/value.schema.json) and modify their existing code inline. - If the user already produces to Kafka without Schema Registry (schemaless), help them migrate: (1) generate a JSON Schema from their existing message structure, (2) register it, and (3) replace their raw
producer.produce()calls with serializer-backed calls. Do not discard their existing code. - If starting from scratch, proceed with the full scaffold below.
- If the user has existing Python code (mentions an existing project, has a
- Target environment? — Confluent Cloud, local Kafka (Docker), or WarpStream. Always prompt for this, even if the user didn't mention it. If they mention "open source", "local", "docker", "self-hosted", or just want to try Kafka without a cloud account, choose local Docker. If they mention "Confluent Cloud", "CC", or have existing cloud credentials, choose Confluent Cloud. If they mention "WarpStream", choose WarpStream. Default to Confluent Cloud if they confirm they don't have a preference, but always ask first.
- If WarpStream: Read
references/warpstream-optimization.mdand apply the librdkafka overrides from that reference. Key changes: disable idempotence, dramatically increase batch sizes and in-flight requests, set large fetch sizes, addws_az=<az>toclient.idfor zone-aware routing. Prefer null message keys for sticky partitioning unless entity-based ordering is required.
- If WarpStream: Read
- Producer, consumer, or both?
- Async or synchronous producer? (Only if producer is requested.) Help the user choose:
- AsyncIO Producer (
AIOProducer): Use when code runs under an event loop — FastAPI/Starlette, aiohttp, Sanic, asyncio workers — and must not block. - Synchronous Producer (
Producer): Use for scripts, batch jobs, and highest-throughput pipelines where the user controls threads/processes and can callpoll()/flush()directly. If the user mentions an async framework (FastAPI, aiohttp, Sanic) or usesasyncio, default to AsyncIO. If they mention scripts, batch, ETL, or don't have a preference, default to Synchronous.
- AsyncIO Producer (
- Do you have an existing schema you'd like to use? If yes, ask the user to paste it or provide the file path, then use it as the
schemas/value.schema.jsoninstead of generating one. If no, proceed to ask about their data fields. - What kind of data are you producing? (Only if the user doesn't have an existing schema. Get field names and types so you can generate a matching JSON Schema and sample data.)
- Topic name? (Default:
demo-topic) - Consumer group ID? (Only if consumer; default:
python-consumer-group)
Don't ask about Schema Registry — always include it. For Confluent Cloud and local Docker, always use JSON Schema. If the target is WarpStream, ask which Schema Registry implementation they are using — WarpStream's built-in schema registry only supports Avro and Protobuf (GET /schemas/types returns ["AVRO","PROTOBUF"]), so if they are using it, ask whether they prefer Avro or Protobuf (default to Avro). If they are using a different SR (e.g., Confluent Cloud Schema Registry), JSON Schema is fine.
Common Agent Mistakes
| Thought | Reality |
|---|---|
| "The user mentioned FastAPI, so I know it's async — skip the questions" | Still confirm. They might want a sync background worker alongside FastAPI. |
| "I'll use Avro since it's more widely used" | This skill uses JSON Schema by default. Exception: WarpStream's built-in schema registry only supports Avro and Protobuf — if the user is using WarpStream SR, use Avro by default. If they're using a different SR (e.g., Confluent Cloud SR), JSON Schema is fine regardless of the Kafka environment. |
| "I'll skip Schema Registry to keep it simple" | Schema Registry is non-negotiable. Every project includes it. |
"I'll use auto.register.schemas=True for convenience" | Always False. Explicit registration is a core principle. |
"I'll create a producer in produce() — it's cleaner" | One producer instance, created in main(), passed as a parameter. Always. |
| "The user wants sync, so the consumer should be sync too" | Consumer is always async (AIOConsumer). This is a deliberate design decision. |
"I'll add headers= to the AIOProducer for schema ID" | AIOProducer.produce() raises NotImplementedError on headers. Only sync producers use headers. |
"I'll swap AsyncJSONSerializer for AsyncAvroSerializer and keep the call site the same" | JSONSerializer takes schema_str first; AvroSerializer takes schema_registry_client first. Calling positionally across formats raises TypeError: ... got multiple values for argument 'schema_registry_client'. Always pass both as kwargs. |
"I'll set message.max.bytes=64000000 on the producer config and fetch.max.bytes=50242880 on the consumer — they're independent" | Not in librdkafka. message.max.bytes is a client-global config (unlike Java's per-role max.request.size), so when get_kafka_config() is shared between producer and consumer, the consumer inherits it. librdkafka enforces fetch.max.bytes >= message.max.bytes at consumer construction and raises KafkaError{_INVALID_ARG, ... "fetch.max.bytes must be >= message.max.bytes"}. Use the WarpStream librdkafka consumer values in references/warpstream-optimization.md (fetch.max.bytes=67108864) — they're sized to satisfy this constraint. |
Step 1b: Confirm Understanding
After gathering all answers, present a confirmation summary before generating any code:
Before I generate the project, let me confirm:
- Project type: [Greenfield scaffold / Migration of existing code]
- Environment: [Confluent Cloud (SASL_SSL) / Local Docker (PLAINTEXT) / WarpStream]
- Schema format: [JSON Schema / Avro / Protobuf] (Avro or Protobuf if using WarpStream's built-in SR)
- Components: [Producer only / Consumer only / Both]
- Producer style: [AsyncIO (AIOProducer) / Synchronous (Producer)] (if applicable)
- Schema: [brief description of user's data fields]
- Topic: [topic name]
- Consumer group: [group ID] (if consumer)
Does this look right?
Wait for user confirmation before proceeding to Step 2. If the user corrects anything, update your understanding and re-confirm.
Step 2: Generate the Project
Decision Flowchart
digraph decisions {
"Q1: Existing app?" -> "Migration path:\nmodify existing code" [label="yes"];
"Q1: Existing app?" -> "Q2: Environment?" [label="no / greenfield"];
"Q2: Environment?" -> "Cloud config\n(SASL_SSL)" [label="Confluent Cloud"];
"Q2: Environment?" -> "Local Docker config\n(PLAINTEXT) + docker-compose.yml" [label="local / docker / OSS"];
"Q2: Environment?" -> "WarpStream config\n(apply overrides from\nreferences/warpstream-optimization.md)" [label="WarpStream"];
"Cloud config\n(SASL_SSL)" -> "Q3: Components?";
"Local Docker config\n(PLAINTEXT) + docker-compose.yml" -> "Q3: Components?";
"WarpStream config\n(apply overrides from\nreferences/warpstream-optimization.md)" -> "Which SR?\n(WarpStream SR → Avro/Protobuf;\nother SR → JSON Schema)";
"Which SR?\n(WarpStream SR → Avro/Protobuf;\nother SR → JSON Schema)" -> "Q3: Components?";
"Q3: Components?" -> "Q4: Async or sync?" [label="producer requested"];
"Q3: Components?" -> "Generate consumer\n(always async AIOConsumer)" [label="consumer only"];
"Q4: Async or sync?" -> "AIOProducer path\nAsyncJSONSerializer\n(no headers support)" [label="async / event-loop"];
"Q4: Async or sync?" -> "Producer path\nJSONSerializer\n(header-based schema ID)" [label="sync / batch / ETL"];
}
Create this file structure in the user's chosen directory:
<project-dir>/
├── producer.py # (if requested)
├── consumer.py # (if requested)
├── common.py # shared config loading + verification helpers
├── schemas/
│ └── value.schema.json # JSON Schema (or value.avsc / value.proto when using WarpStream SR)
├── tests/
│ └── test_project.py # unit tests (always generated)
├── .env.example # template for credentials
├── requirements.txt
├── docker-compose.yml # (local Docker path only)
Security
NEVER read, open, or display .env files. They contain API keys and secrets. Only generate .env.example with placeholder values. If the user asks you to debug a connection issue, ask them to verify their .env values themselves — do not read the file.
Core Principles
These principles matter because they prevent the most common production issues with Kafka Python clients:
-
Reuse the producer instance. Creating a new producer per message is expensive — each one opens new TCP connections, does SASL handshakes, and fetches metadata. Create one producer and reuse it for all messages. The produce function should accept the producer as a parameter, not instantiate one.
-
Always use Schema Registry with JSON Schema. Schema Registry enforces a contract between producers and consumers. Without it, schema changes silently break downstream consumers. This skill uses JSON Schema by default. Schema Registry supports Avro, Protobuf, and JSON Schema — JSON Schema is chosen because: (1) Python has first-class JSON support with no code generation step, (2)
confluent-kafka-pythonprovidesJSONSerializer/JSONDeserializerout of the box, (3) it is the most approachable format for Python developers already working with JSON/dict data.WarpStream Schema Registry exception: WarpStream's built-in schema registry only supports Avro and Protobuf — its
GET /schemas/typesendpoint returns["AVRO","PROTOBUF"]. JSON Schema is not available when using WarpStream's SR. If the user is running WarpStream and using its built-in SR, use Avro by default (or Protobuf if the user prefers). UseAvroSerializer/AvroDeserializerfromconfluent_kafka.schema_registry.avro(orProtobufSerializer/ProtobufDeserializerfromconfluent_kafka.schema_registry.protobuf). For async producers, useAsyncAvroSerializerfromconfluent_kafka.schema_registry._async.avro(orAsyncProtobufSerializerfromconfluent_kafka.schema_registry._async.protobuf). Generate an Avro schema file atschemas/value.avsc(orschemas/value.protofor Protobuf) instead ofschemas/value.schema.json— usereferences/value.avscas the starting point and adapt to the user's domain. When generating producer/consumer code on the Avro path, copy fromreferences/producer_avro.py(async),references/producer_avro_sync.py(sync), andreferences/consumer_avro.py(async) — not the JSON reference files — because the constructor signatures differ (see the Avro constructor warning below). If the user is running WarpStream but pointing at a different SR (e.g., Confluent Cloud Schema Registry), JSON Schema works fine — use the default path.Register schemas as a separate explicit step before creating the serializer. Use a dedicated
register_schema()function that callssr_client.register_schema()and lets errors (auth failures, network errors, permission denials) propagate immediately — never wrap registration in a baretry/except. Then configure the serializer withauto.register.schemas=Falseanduse.latest.version=True. This ensures the serializer never silently auto-registers and aligns with production practice where CI/CD registers schemas, not application startup.Use the appropriate serializer for the chosen producer style and schema format:
- JSON Schema (default):
AsyncJSONSerializer/AsyncJSONDeserializerfromconfluent_kafka.schema_registry._async.json_schemafor async, orJSONSerializer/JSONDeserializerfromconfluent_kafka.schema_registry.json_schemafor synchronous. - Avro (default when using WarpStream SR):
AsyncAvroSerializer/AsyncAvroDeserializerfromconfluent_kafka.schema_registry._async.avrofor async, orAvroSerializer/AvroDeserializerfromconfluent_kafka.schema_registry.avrofor synchronous. - Protobuf (alternative when using WarpStream SR):
AsyncProtobufSerializer/AsyncProtobufDeserializerfromconfluent_kafka.schema_registry._async.protobuffor async, orProtobufSerializer/ProtobufDeserializerfromconfluent_kafka.schema_registry.protobuffor synchronous.
- JSON Schema (default):
-
Choose the right producer style. The
confluent-kafka-pythonlibrary offers two producer APIs:- AsyncIO Producer (
AIOProducerfromconfluent_kafka.aio): Non-blocking, integrates withasyncioevent loops. Use withAsyncJSONSerializerfromconfluent_kafka.schema_registry._async.json_schemaandAsyncSchemaRegistryClient. Best for applications already running an event loop (FastAPI, aiohttp, Sanic, asyncio workers). - Synchronous Producer (
Producerfromconfluent_kafka): Blocking calls with delivery callbacks. Use withJSONSerializerfromconfluent_kafka.schema_registry.json_schemaandSchemaRegistryClient. Best for scripts, batch jobs, and highest-throughput pipelines where the user controls threads/processes and can callpoll()/flush()directly. Always ask the user which style fits their use case. The consumer always usesAIOConsumer(async) — long-running poll loops benefit from non-blocking I/O, and mixing sync/async consumer styles adds complexity with little benefit.
- AsyncIO Producer (
-
Graceful shutdown. Async producers must
flush()andclose()(both awaited) before exiting. Synchronous producers must callflush()before exiting — otherwise buffered messages are lost. Consumers mustunsubscribe()thenclose()to leave the consumer group cleanly (avoiding unnecessary rebalances). Usetry/finallyblocks and handleKeyboardInterrupt/ signals. -
Support Confluent Cloud, local Docker, and WarpStream. When targeting Confluent Cloud, configure
SASL_SSLwithPLAINmechanism and load API keys from.env. When targeting local Docker, usePLAINTEXTwith no authentication. When targeting WarpStream, useSASL_SSLorPLAINTEXTdepending on the user's WarpStream deployment, and apply the librdkafka overrides fromreferences/warpstream-optimization.md(large batches, disabled idempotence, large fetches, zone-awareclient.id). TheKAFKA_ENVenvironment variable (cloud,local, orwarpstream) controls which path is used. Load all settings from environment variables via.env. -
Verify connectivity before running. Use
AdminClient.list_topics()to verify the broker is reachable and the topic exists before producing or consuming. Verify Schema Registry connectivity with an HTTP health check. -
Always set a message key for domain events. Pass
key=<entity_id>.encode("utf-8")toproducer.produce()for any message that represents an entity or event stream (order events, user actions, device telemetry, transactions). Kafka partitions by key, so messages with the same key land on the same partition and preserve ordering — critical for event streams likeOrderCreated → OrderUpdated → OrderCancelledwhere consumers must see events in order. Theproduce()helper in every reference file accepts akey_fieldparameter naming the field to use as the key (e.g.,key_field="order_id",key_field="transaction_id"). Ask the user which field identifies the entity and pass it toproduce(). Only leavekey_field=Noneif the user explicitly states ordering does not matter (e.g., stateless metrics where any partition is fine).WarpStream exception: On WarpStream, null keys enable sticky partitioning, which builds larger batches and significantly improves throughput and cost. When the user's use case does not require per-entity ordering (e.g., independent telemetry readings, stateless metrics, logs), recommend
key_field=Noneand explain the throughput benefit. When per-entity ordering is required (e.g.,OrderCreated → OrderUpdated → OrderCancelledfor the same order), still set a message key — correctness takes priority over batching efficiency. Ask the user whether their events need per-key ordering to decide.
common.py
This module handles configuration loading and connectivity verification. Use references/common.py as the template.
producer.py Pattern (AsyncIO)
When the user chooses the AsyncIO producer, use references/producer.py as the template.
Key points:
produce()takes a producer instance as a parameter — it never creates one- The producer is created once in
main()and can be passed to multipleproduce()calls - The async serializer (
AsyncJSONSerializer) must beawaited when calling it on a message AIOProducer.produce()is async and returns anasyncio.Future. You mustawaitthe method to get the Future, thenawaitthe Future to get the deliveredMessage:future = await producer.produce(...); result = await futureAIOProducer.flush()andclose()are coroutines — they must beawaited in thefinallyblock- Signal handlers set a shutdown event for graceful termination
- Schema registration and serializer creation are separate steps.
register_schema()explicitly registers the schema and returns the schema ID — errors propagate immediately.create_json_serializer()(orcreate_avro_serializer()/create_protobuf_serializer()when using Avro/Protobuf) creates the serializer withconf={'auto.register.schemas': False, 'use.latest.version': True}. - Constructor argument order differs across formats — always pass
schema_registry_clientandschema_stras keyword arguments, never positionally. The JSON, Avro, and Protobuf serializer/deserializer classes inconfluent-kafka-pythondo not share the same positional signature:JSONSerializer/JSONDeserializertakeschema_strfirst, whileAvroSerializer/AvroDeserializerandProtobufSerializer/ProtobufDeserializertakeschema_registry_clientfirst. Mixing positional and keyword forms across formats producesTypeError: ... got multiple values for argument 'schema_registry_client'. Use kwargs everywhere so the pattern is identical:await AsyncJSONSerializer(schema_str=schema_str, schema_registry_client=sr_client, conf=conf)await AsyncAvroSerializer(schema_str=schema_str, schema_registry_client=sr_client, conf=conf)await AsyncProtobufSerializer(msg_type=MsgType, schema_registry_client=sr_client, conf=conf)(Protobuf takes a generated message class instead of a schema string) Seereferences/producer.py(JSON) andreferences/producer_avro.py(Avro) for the canonical templates.
- Headers are NOT supported with
AIOProducerbatch mode. Do not passheaders=toAIOProducer.produce()— it will raiseNotImplementedError. Schema identification is handled automatically by the serializer's wire format prefix (applies to JSON Schema, Avro, and Protobuf serializers). See "Schema ID in Headers vs Wire Format" below for details
producer.py Pattern (Synchronous)
When the user chooses the synchronous producer, use references/producer_sync.py as the template.
Shortened here. Read the whole file on GitHub.
Signals
- GitHub stars
- 58
- Forks
- 11
- Last commit
- Sep 2026
ahel review
K1binfo
installs-packagesK6low
bundled executables the agent is told to runK1binfo
installs-packages (in references/readme-template.md)
Automated review, not a security audit. Ruleset v1+k2.
Advanced
- Item type
- skill
- Key
developing-kafka-python-client- Source
- github.com/confluentinc/agent-skills
github.com/confluentinc/agent-skills