Agent skill

Python Kafka Module Skill

by jiushiwon in jiushiwon/wg-skills

Python Kafka 模块快速集成技能。面向已拥有 FastAPI 项目骨架的开发者,提供 Kafka 生产者、消费者、消息订阅、事件驱动等能力的快速集成。触发词:"Python Kafka"、"FastAPI Kafka"、"Kafka 集成"、"kafka producer"、"kafka consumer"、"kafka 消息"、"kafka 事件"、"kafka 队列"。

Apache-2.0Auto-check: notesBackend & APIs

Install Python Kafka Module Skill

skills CLI
$ npx skills add jiushiwon/wg-skills --skill python-kafka-module-skill -a claude-code

Project install by default; add -g for ~/.claude/skills/.

GitHub CLI
$ gh skill install jiushiwon/wg-skills python-kafka-module-skill --agent claude-code

Project scope by default; add --scope user for a personal install. Needs GitHub CLI 2.90.0 or later (public preview).

Manual copy
$ git clone --depth 1 https://github.com/jiushiwon/wg-skills.git skills-src && mkdir -p .claude/skills && cp -r skills-src/vibeCoding/backend/python/fastapi-module/python-kafka-module-skill .claude/skills/python-kafka-module-skill && rm -rf skills-src

Use ~/.claude/skills/ instead of .claude/skills for a personal install. The folder must contain SKILL.md.

Claude Code skills documentation · loads skills from .claude/skills/

Facts

Skill name
python-kafka-module-skill
GitHub stars
110
Token cost
~2.1k tokens
SKILL.md length
55 words
Files
2
Skills in repo
121
Repo updated
First seen
Licence
Apache-2.0

At a glance

Python Kafka 模块快速集成技能。面向已拥有 FastAPI 项目骨架的开发者,提供 Kafka 生产者、消费者、消息订阅、事件驱动等能力的快速集成。触发词:"Python Kafka"、"FastAPI Kafka"、"Kafka 集成"、"kafka producer"、"kafka consumer"、"kafka 消息"、"kafka 事件"、"kafka 队列"。

  • Works in 4 steps: 生产者 → 消费者 → 事务消息 → …
  • Tasks that involve Event-driven systems
  • SKILL.md covers 能力清单, 触发场景, 依赖配置 and 默认方法封装, plus 2 more sections
  • Calls pip

What it does

Python Kafka Module Skill is an agent skill from jiushiwon/wg-skills. Python Kafka 模块快速集成技能。面向已拥有 FastAPI 项目骨架的开发者,提供 Kafka 生产者、消费者、消息订阅、事件驱动等能力的快速集成。触发词:"Python Kafka"、"FastAPI Kafka"、"Kafka 集成"、"kafka producer"、"kafka consumer"、"kafka 消息"、"kafka 事件"、"kafka 队列"。

Its SKILL.md is about 2.1k tokens, which your agent loads only when the skill is triggered. The skill folder holds 1 other file (for example `README.md`).

It sits in Backend & APIs, covering Event-driven systems and Backend development. It works with Apache Kafka, Python and FastAPI. The licence is Apache-2.0.

When your agent uses it

  • Tasks that involve Event-driven systems
  • Tasks that involve Backend development

Example prompts

  • “Python Kafka”
  • “FastAPI Kafka”
  • “Kafka 集成”
  • “/python-kafka-module-skill”

Requirements

  • Python 3
  • Docker

Workflow steps

4 steps, taken from the step headings in SKILL.md.

  1. 生产者
  2. 消费者
  3. 事务消息
  4. 配置管理

What it can do on your machine

Read from SKILL.md and the folder at commit a4a640b. It shows what the files ask for, not the result of running them.

  • Tool permissions

    Pre-approves nothing: there is no allowed-tools line, so your agent's usual permission prompts apply.

    From allowed-tools in the SKILL.md frontmatter.

  • Runs code

    Shell commands in SKILL.md call:

    • pip

    From the folder's file list and the shell code blocks in SKILL.md.

  • Network

    No URLs in SKILL.md. Its commands use pip, which can reach the network depending on how they are called.

    From URLs in SKILL.md, links to its own repository left out.

  • Credentials

    Names no API keys, tokens, secrets or passwords.

    From names ending in _API_KEY, _TOKEN, _SECRET, _KEY or _PASSWORD in SKILL.md.

Context cost

Python Kafka Module Skill loads about 2.1k tokens when it runs. Until then it costs about 55 tokens; SKILL.md has 55 words of instructions outside code blocks.

Always · name and description, kept in context so the agent knows when to use it
~55
When it runs · the whole SKILL.md, loaded when a task matches
~2.1k

Estimates: characters ÷ 4, the usual rule of thumb; real counts depend on the model's tokenizer. Scripts and assets cost tokens only if the agent reads them.

Safety

Auto-check: notes

The automated check noted patterns worth knowing about, such as sudo or a known installer.

  • NoteMentions a .env fileSKILL.md:286
    ### .env

Automated static check — not a guarantee. Review scripts before installing. It scans the text of SKILL.md for risky patterns (piping downloads into a shell, reading credential files, hidden Unicode, destructive commands); files beside SKILL.md are not scanned.

SKILL.md

The full file from jiushiwon/wg-skills at commit a4a640b, republished under its Apache-2.0 licence (© jiushiwon). 55 words, ~2,109 tokens.

Download SKILL.mdSave it as .claude/skills/python-kafka-module-skill/SKILL.md (or your agent's skills folder). This skill also uses 1 other file; get the full folder from GitHub.
name
python-kafka-module-skill
description
Python Kafka 模块快速集成技能。面向已拥有 FastAPI 项目骨架的开发者,提供 Kafka 生产者、消费者、消息订阅、事件驱动等能力的快速集成。触发词:"Python Kafka"、"FastAPI Kafka"、"Kafka 集成"、"kafka producer"、"kafka consumer"、"kafka 消息"、"kafka 事件"、"kafka 队列"。

Python Kafka Module Skill

面向已有 FastAPI 项目的开发者,快速集成 Kafka 能力。

能力清单

能力说明
生产者同步/异步发送消息、消息分区、消息key
消费者消费监听、消息重试、消费者组
消息序列化JSON
事务消息Kafka 事务
错误处理消息发送/消费错误处理

触发场景

用户说"帮我加 Kafka"或"集成 Kafka"时触发。

依赖配置

bash
pip install aiokafka

默认方法封装

1. 生产者
python
# kafka_producer.py
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer
from typing import Optional, Callable, Any
import json
import logging

logger = logging.getLogger(__name__)

class KafkaProducer:
    def __init__(self, bootstrap_servers: str = "localhost:9092"):
        self.bootstrap_servers = bootstrap_servers
        self._producer: Optional[AIOKafkaProducer] = None
    
    async def start(self):
        """启动生产者"""
        self._producer = AIOKafkaProducer(
            bootstrap_servers=self.bootstrap_servers,
            value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode('utf-8'),
            key_serializer=lambda k: k.encode('utf-8') if k else None
        )
        await self._producer.start()
        logger.info("Kafka 生产者已启动")
    
    async def stop(self):
        """停止生产者"""
        if self._producer:
            await self._producer.stop()
            logger.info("Kafka 生产者已停止")
    
    async def send(self, topic: str, value: Any, key: Optional[str] = None) -> str:
        """发送消息(同步)"""
        if not self._producer:
            raise RuntimeError("生产者未启动")
        
        future = await self._producer.send_and_wait(topic, value, key=key)
        return f"{future.topic}-{future.partition}-{future.offset}"
    
    async def send_async(self, topic: str, value: Any, key: Optional[str] = None, callback: Optional[Callable] = None):
        """发送消息(异步)"""
        if not self._producer:
            raise RuntimeError("生产者未启动")
        
        await self._producer.send(topic, value, key=key)
        if callback:
            # 注册回调
            pass
    
    async def send_json(self, topic: str, data: dict, key: Optional[str] = None):
        """发送 JSON 消息"""
        return await self.send(topic, data, key)
    
    async def send_messages(self, topic: str, messages: list):
        """批量发送消息"""
        if not self._producer:
            raise RuntimeError("生产者未启动")
        
        for msg in messages:
            await self._producer.send(topic, msg)

# 全局实例
producer = KafkaProducer()
2. 消费者
python
# kafka_consumer.py
from aiokafka import AIOKafkaConsumer
from typing import Optional, Callable, Any
import json
import asyncio
import logging

logger = logging.getLogger(__name__)

class KafkaConsumer:
    def __init__(
        self,
        bootstrap_servers: str = "localhost:9092",
        group_id: str = "my-group",
        topics: list = None
    ):
        self.bootstrap_servers = bootstrap_servers
        self.group_id = group_id
        self.topics = topics or []
        self._consumer: Optional[AIOKafkaConsumer] = None
        self._running = False
    
    async def start(self):
        """启动消费者"""
        self._consumer = AIOKafkaConsumer(
            *self.topics,
            bootstrap_servers=self.bootstrap_servers,
            group_id=self.group_id,
            value_deserializer=lambda m: json.loads(m.decode('utf-8')),
            key_deserializer=lambda k: k.decode('utf-8') if k else None,
            auto_offset_reset='earliest',
            enable_auto_commit=False
        )
        await self._consumer.start()
        self._running = True
        logger.info(f"Kafka 消费者已启动,订阅主题: {self.topics}")
    
    async def stop(self):
        """停止消费者"""
        self._running = False
        if self._consumer:
            await self._consumer.stop()
            logger.info("Kafka 消费者已停止")
    
    async def consume(self, handler: Callable):
        """消费消息"""
        if not self._consumer:
            raise RuntimeError("消费者未启动")
        
        async for msg in self._consumer:
            if not self._running:
                break
            
            try:
                logger.info(f"收到消息: topic={msg.topic}, partition={msg.partition}, offset={msg.offset}")
                await handler(msg)
                # 手动提交偏移量
                await self._consumer.commit()
            except Exception as e:
                logger.error(f"处理消息失败: {e}")
                # 可选:发送到死信队列或重试
    
    async def consume_batch(self, handler: Callable, batch_size: int = 100):
        """批量消费消息"""
        if not self._consumer:
            raise RuntimeError("消费者未启动")
        
        messages = []
        async for msg in self._consumer:
            if not self._running:
                break
            
            messages.append(msg)
            if len(messages) >= batch_size:
                try:
                    await handler(messages)
                    await self._consumer.commit()
                except Exception as e:
                    logger.error(f"批量处理消息失败: {e}")
                messages = []
        
        # 处理剩余消息
        if messages:
            try:
                await handler(messages)
                await self._consumer.commit()
            except Exception as e:
                logger.error(f"处理剩余消息失败: {e}")

# 使用示例
async def handle_message(msg):
    """消息处理函数"""
    print(f"处理消息: {msg.value}")
    # 业务逻辑

async def main():
    consumer = KafkaConsumer(
        bootstrap_servers="localhost:9092",
        group_id="my-group",
        topics=["my-topic"]
    )
    await consumer.start()
    try:
        await consumer.consume(handle_message)
    finally:
        await consumer.stop()
3. 事务消息
python
# kafka_transaction.py
from aiokafka import AIOKafkaProducer
import asyncio

class KafkaTransactionProducer:
    def __init__(self, bootstrap_servers: str = "localhost:9092"):
        self.bootstrap_servers = bootstrap_servers
        self._producer: Optional[AIOKafkaProducer] = None
    
    async def start(self):
        self._producer = AIOKafkaProducer(
            bootstrap_servers=self.bootstrap_servers,
            enable_idempotence=True  # 开启幂等性
        )
        await self._producer.start()
    
    async def stop(self):
        if self._producer:
            await self._producer.stop()
    
    async def send_in_transaction(self, messages: list):
        """事务发送:要么全部成功,要么全部失败"""
        if not self._producer:
            raise RuntimeError("生产者未启动")
        
        transaction = self._producer.transaction()
        await transaction.begin()
        
        try:
            for topic, key, value in messages:
                await transaction.send(topic, value, key=key)
            await transaction.commit()
            return True
        except Exception as e:
            await transaction.abort()
            raise e

# 使用
tp = KafkaTransactionProducer()
await tp.start()
await tp.send_in_transaction([
    ("topic1", "key1", {"data": "1"}),
    ("topic2", "key2", {"data": "2"})
])
4. 配置管理
python
# kafka_config.py
from pydantic_settings import BaseSettings
from typing import Optional

class KafkaSettings(BaseSettings):
    bootstrap_servers: str = "localhost:9092"
    producer_group_id: str = "producer-group"
    consumer_group_id: str = "consumer-group"
    topics: list = ["my-topic"]
    
    # 生产者配置
    acks: str = "all"
    retries: int = 3
    batch_size: int = 16384
    linger_ms: int = 10
    
    # 消费者配置
    auto_offset_reset: str = "earliest"
    enable_auto_commit: bool = False

kafka_settings = KafkaSettings()

配置模板

.env
env
KAFKA_HOST=localhost
KAFKA_PORT=9092
KAFKA_TOPICS=my-topic,user-events
使用示例
python
# main.py
from fastapi import FastAPI
from kafka_producer import producer
from kafka_consumer import KafkaConsumer

app = FastAPI()

@app.on_event("startup")
async def startup():
    await producer.start()

@app.on_event("shutdown")
async def shutdown():
    await producer.stop()

@app.post("/send")
async def send_message(topic: str, data: dict):
    await producer.send_json(topic, data)
    return {"status": "ok"}

不做

  • 不负责安装 Kafka(用户自行安装或使用 Docker)
  • 不处理 Kafka 集群配置(单节点为主)
  • 不提供消息持久化策略
  • 不处理 Kafka Connect 数据同步
  • 不处理 Schema Registry

© jiushiwon, Apache-2.0. Rendered from Markdown: HTML in the file is shown as text, images as links, and headings moved down two levels. Raw file

Files

SKILL.md and 1 other file in vibeCoding/backend/python/fastapi-module/python-kafka-module-skill of jiushiwon/wg-skills.

  • SKILL.md
  • README.md

Open the folder on GitHubat commit a4a640b

Compare with similar skills

Python Kafka Module Skill next to the 5 skills that share the most tags, products or categories with it. Stars are the repository's; “used in” counts other GitHub owners with a copy.

Python Kafka Module Skill compared with similar skills
SkillStarsUsed inTokensAuto-checkLicenceRepo updated
Python Kafka Module Skill this skilljiushiwon/wg-skills110—~2.1kAutomated safety check: NotesApache-2.0
Deepstream SopNVIDIA/skills3.5k—~4.7kAutomated safety check: NotesApache-2.0
Fastcrudbenavlabs/fastcrud1.6k—~5kAutomated safety check: PassMIT
Phoenix ServerArize-ai/phoenix12k—~1.6kAutomated safety check: PassCustom licence
Opensource Guide Coachcalf-ai/calfkit-sdk1491 repos~2.1kAutomated safety check: PassApache-2.0
FastAPI Project Templateswshobson/agents40k11 repos~901Automated safety check: PassMIT

Similar skills

  • Deepstream Sop

    NVIDIA/skills

    Official

    A skill your agent uses when building, deploying, evaluating, debugging, or measuring latency for the DeepStream SOP Inference Microservice — a GPU-accelerated FastAPI service that detects whether…

    3.5k GitHub stars~4.7k tokensUpdated today
    Backend & APIsAuto-check: notes
  • Fastcrud

    benavlabs/fastcrud

    A skill your agent uses when building or modifying CRUD endpoints with FastCRUD (the fastcrud PyPI package) in a FastAPI project — covers FastCRUD, crudrouter, EndpointCreator, FilterConfig…

    1.6k GitHub stars~5k tokensUpdated 13 days ago
    Backend & APIsAuto-check passed
  • Phoenix Server

    Arize-ai/phoenix

    Backend development guide for the Phoenix AI observability platform (Strawberry GraphQL, SQLAlchemy async, FastAPI).

    12k GitHub stars~1.6k tokensUpdated today
    Backend & APIsAuto-check passed
  • Opensource Guide Coach

    calf-ai/calfkit-sdk

    A skill your agent uses when a user wants guidance on starting, contributing to, growing, governing, funding, securing, or sustaining an open source project, or asks about contributor onboarding…

    149 GitHub starsUsed in 1 repo~2.1k tokens
    Backend & APIsAuto-check passed
  • Scaffolds FastAPI projects with a layered app layout, dependency injection through Depends, async handlers and database access, middleware and pytest setup.

    40k GitHub starsUsed in 11 repos~901 tokens
    Backend & APIsAuto-check passed
  • Holm Web

    volfpeter/holm

    A skill your agent uses when working on web apps built with holm, or to answer questions about holm.

    132 GitHub stars~1.2k tokensUpdated 21 days ago
    Backend & APIsAuto-check passed

More from jiushiwon/wg-skills

All 121 skills in this repo
  • Frontend UI Foundry

    jiushiwon/wg-skills

    A skill your agent uses when generating UI for a specific scenario (mobile/PC/官网/管理端/营销页/文档/金融/原生/3D), when refactoring an existing HTML/Vue/React project to a unified design system, when extracting…

    110 GitHub stars~1.4k tokensUpdated 3 days ago
    Auto-check passed
  • Workflow Diagram Skill

    jiushiwon/wg-skills

    A skill your agent uses when 用户想用一句话生成流程图、工作流图解、开发流程图、AI 流程图或任何步骤型图解,支持多风格输出(flat icon / 暖白手账风 / dark / cute),并提供现成模板一键出图。

    110 GitHub stars~1.2k tokensUpdated 3 days ago
    Auto-check passed
  • Uniapp App Generate Skill

    jiushiwon/wg-skills

    This skill should be used when the user wants to create a standardized uni-app project for WeChat mini-program, H5, and App from scratch.

    110 GitHub stars~5k tokensUpdated 3 days ago
    Auto-check: notes
  • Fastapi Init Skill

    jiushiwon/wg-skills

    FastAPI 项目一键初始化技能。面向零基础小白,提供环境探测、自动安装、完整 Web 骨架生成、SSE 流式框架、JWT 鉴权、统一响应封装、文件上传接口、一键启动/重启脚本、Swagger 文档,内置 MySQL(默认)/ PostgreSQL / MongoDB 数据库选择。用户只需说"帮我搭一个 FastAPI 项目"即可一条命令完成从零到跑的完整链路。触发词:"FastAPI…

    110 GitHub stars~1.8k tokensUpdated 3 days ago
    Auto-check: notes
  • Ffmpeg Skill

    jiushiwon/wg-skills

    FFmpeg 多媒体处理技能 — 将自然语言描述转为正确的 ffmpeg 命令,覆盖视频剪辑、转码、水印、合成、提取、生成等操作。内置一键安装脚本(Windows/Mac/Linux)。触发词:ffmpeg、视频剪辑、视频裁剪、视频转码、视频压缩、去水印、加水印、视频拼接、视频合成、提取音频、视频转…

    110 GitHub stars~1.4k tokensUpdated 3 days ago
    Auto-check passed
  • Frontend Request Skill

    jiushiwon/wg-skills

    A skill your agent uses when designing or reviewing the request layer of a frontend project (web / uni-app / mini-program), including request.ts wrappers, interceptors, deduplication, mocks, error…

    110 GitHub stars~3.6k tokensUpdated 3 days ago
    Auto-check passed

Categories

Questions about Python Kafka Module Skill

What does Python Kafka Module Skill do?

Python Kafka 模块快速集成技能。面向已拥有 FastAPI 项目骨架的开发者,提供 Kafka 生产者、消费者、消息订阅、事件驱动等能力的快速集成。触发词:"Python Kafka"、"FastAPI Kafka"、"Kafka 集成"、"kafka producer"、"kafka consumer"、"kafka 消息"、"kafka 事件"、"kafka 队列"。. Python Kafka Module Skill is an agent skill from jiushiwon/wg-skills.

When should I use Python Kafka Module Skill?

Python Kafka Module Skill fits situations like: tasks that involve Event-driven systems; tasks that involve Backend development.

How do I install Python Kafka Module Skill in Claude Code?

Run `npx skills add jiushiwon/wg-skills --skill python-kafka-module-skill -a claude-code`. Or copy the skill folder (vibeCoding/backend/python/fastapi-module/python-kafka-module-skill in jiushiwon/wg-skills) into .claude/skills/python-kafka-module-skill in your project. Claude Code loads it when a task matches its description.

How do I install Python Kafka Module Skill in Codex?

Run `npx skills add jiushiwon/wg-skills --skill python-kafka-module-skill -a codex`. Or copy the skill folder (vibeCoding/backend/python/fastapi-module/python-kafka-module-skill in jiushiwon/wg-skills) into .agents/skills/python-kafka-module-skill in your project. Codex loads it when a task matches its description.

Can I use Python Kafka Module Skill in Cursor, Gemini CLI or GitHub Copilot?

Cursor, Gemini CLI, GitHub Copilot and OpenCode also load SKILL.md folders. With the skills CLI, run `npx skills add jiushiwon/wg-skills --skill python-kafka-module-skill -a cursor` (or -a gemini-cli, github-copilot or opencode for the others). To copy it by hand, put the folder in .cursor/skills/python-kafka-module-skill, .gemini/skills/python-kafka-module-skill, .github/skills/python-kafka-module-skill and .opencode/skills/python-kafka-module-skill in your project.

What does Python Kafka Module Skill need to run?

Going by SKILL.md and its folder, Python Kafka Module Skill needs the command-line tools its instructions call (pip). Our summary lists: Python 3; Docker.

Does Python Kafka Module Skill access the network?

SKILL.md contains no URLs. Its commands use pip, which can reach the network depending on how they are called. This is read from the text; nothing was executed.

Is Python Kafka Module Skill safe to install?

Our automated static check of SKILL.md found notes only (mentions a .env file), nothing it rates as a warning. It is not a guarantee. Review the folder before installing.

What licence does Python Kafka Module Skill use?

Python Kafka Module Skill is published under the Apache-2.0 licence (the repository's licence). It allows redistribution, so the full SKILL.md is shown on this page.

How many tokens does Python Kafka Module Skill use?

About 2.1k tokens (SKILL.md is roughly 8.4k characters). Agents keep only the skill's name and description in context until a task matches; then they load SKILL.md in full.

What are the alternatives to Python Kafka Module Skill?

Skills that share tags, products or a category with Python Kafka Module Skill: Deepstream Sop (NVIDIA/skills, 3.5k stars), Fastcrud (benavlabs/fastcrud, 1.6k stars), Phoenix Server (Arize-ai/phoenix, 12k stars) and Opensource Guide Coach (calf-ai/calfkit-sdk, 149 stars). The comparison table on this page puts their stars, adoption, token cost, safety result and licence side by side.

Who maintains Python Kafka Module Skill?

jiushiwon (a GitHub user) maintains it in jiushiwon/wg-skills, which has 110 GitHub stars. The repository holds 121 skills in this directory. The repository was last updated on October 4, 2026.

Source: jiushiwon/wg-skills on GitHub. Facts on this page come from the repository at the commit we read; the author's words are quoted as theirs.