Skills DirectorySkills Directory
SkillsLearnSecurityCategoriesDocsBlogPro
Sign InSubmit Skill
Skills Directory

Security-tested agent skills for Claude, coding agents, and AI workflows.

Directory

  • Browse Skills
  • All Skills A–Z
  • Claude Skills
  • Claude Code Skills
  • Agent Skills
  • Categories
  • Authors
  • Submit a Skill

Learn

  • Learn Hub
  • Install Claude Skills
  • Write SKILL.md
  • Skills vs MCP
  • Directories Compared

Security

  • Security
  • Methodology
  • Secure Claude Skills
  • Security Badges
  • Chrome Extension
  • Skill Manager

Company

  • About
  • Community
  • Blog
  • API Docs
  • Advertise

2026 Skills Directory. All rights reserved.

ProTermsPrivacyRefunds
Back to skills

Kafka

ASecurity

Design and operate Kafka event streaming: topics, producers, consumers, consumer groups, partitioning, and exactly-once. Use for event-driven systems.

2 stars
0 votes
0 copies
0 views
Added 9/29/2026
ai-agentspythongojavabashnodedockerdevops

Works with

cli

Security Analysis

A96/100
mediumInstalls packages at runtime which could introduce malicious dependencies

Pro scans all 2 files and shows the line behind each finding

Scanned 9/29/2026

$npx -y skills add ssrjkk/claude-skills --skill kafka --agent claude-code

Installs into .claude/skills of the current project.

Are you the author of Kafka?

Add the live security badge to your README — it updates automatically with every re-scan.

Security grade badge for Kafka
[![Security: A — Skills Directory](https://www.skillsdirectory.com/api/skills/ssrjkk-kafka/badge)](https://www.skillsdirectory.com/skills/ssrjkk-kafka)

More formats (shields.io, HTML) on the badges page. Keep it an A: scan every change in CI with Pro.

Download with Pro
Files
SKILL.md
---
name: kafka
description: "Design and operate Kafka event streaming: topics, producers, consumers, consumer groups, partitioning, and exactly-once. Use for event-driven systems."
category: devops
tags: [kafka, event-streaming, message-queue, topics, consumers, partitioning, streaming]
models: [sonnet, opus, gpt-5, gemini-2.5, glm-4.6]
version: 1.0.0
created: 2026-09-26
updated: 2026-09-28
author: ssrjkk
---
# Kafka

> Building reliable event streaming with Apache Kafka.

## Quick Start
```bash
docker compose up -d  # single-node kafka
kafka-topics --bootstrap-server localhost:9092 --create --topic events --partitions 3
kafka-console-producer --bootstrap-server localhost:9092 --topic events
```

## When to Use
- Event-driven architecture and decoupling
- High-throughput log ingestion and analytics
- Stream processing (join, aggregate, window)
- Durable message replay for consumers

## Best Practices

### Topics & Partitioning
- Model topics by event type, not by consumer
- Partition by key for ordering within a key
- Choose partition count based on throughput
- Set retention per topic (time and size)

### Producers
- Use a key for order/co-partitioning needs
- Set acks (all) for durability; tune batching
- Handle retries and idempotence (enable idempotent producer)
- Monitor `record-error-rate` and latency

### Consumers
- Use consumer groups for scale
- Handle rebalance gracefully
- Commit offsets after processing (at-least-once)
- Make consumption idempotent for exactly-once semantics

### Operations
- Run at least 3 brokers with replication factor 3
- Monitor lag, under-replicated partitions, and disk
- Use Schema Registry for schema evolution
- Set alerts on consumer lag

## Dependencies
```bash
# Python client
pip install confluent-kafka
# or Java:
# org.apache.kafka:kafka-clients
```

## Examples
```python
from confluent_kafka import Producer

p = Producer({"bootstrap.servers": "localhost:9092"})

def acked(err, msg):
    if err is not None:
        print(f"failed: {err}")

p.produce("events", key="user-42", value=b'{"action":"login"}', callback=acked)
p.flush()
```
```python
from confluent_kafka import Consumer, KafkaError

c = Consumer({
    "bootstrap.servers": "localhost:9092",
    "group.id": "analytics",
    "auto.offset.reset": "earliest",
})
c.subscribe(["events"])

while True:
    msg = c.poll(1.0)
    if msg is None:
        continue
    if msg.error():
        if msg.error().code() == KafkaError._PARTITION_EOF:
            continue
        break
    process(msg.value())        # idempotent processing
    c.commit(asynchronous=False)  # commit after processing
```
```java
// Java producer with acks and idempotence
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("acks", "all");
props.put("enable.idempotence", "true");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("events", "user-42", "login"));
```
```bash
# Monitor consumer lag
kafka-consumer-groups --bootstrap-server localhost:9092 \
  --describe --group analytics
```

## Step-by-Step
1. Define event types and topic names up front.
2. Choose partitioning strategy per key.
3. Set replication factor and retention.
4. Implement producers with acks and idempotence.
5. Implement consumers with graceful rebalance and offsets.
6. Make processing idempotent for safe replays.
7. Use Schema Registry for contract evolution.
8. Monitor lag, disk, and error rates.

## Validation
1. Producers send without errors under load
2. Consumers process in order within a key
3. Consumer lag stays near zero at steady state
4. Rebalances don't lose or duplicate processed events
5. Replaying from an offset produces identical results

## Troubleshooting
- High lag: scale consumers or optimize processing.
- Out-of-order: partition by key and check single-consumer-per-partition.
- Duplicates: make processing idempotent or use exactly-once.

Attribution

ssrjkkssrjkk
View sourceSee grades on GitHubMore from ssrjkk →
SSkills DirectorySkills Directory

Ship a skill? Prove it's safe.

Free 120-pattern security scan, letter grade, and an embeddable README badge.

Submit a skill

Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.

Comments (0)

No comments yet. Be the first to comment!

SSkills DirectorySkills Directory

Ship a skill? Prove it's safe.

Free 120-pattern security scan, letter grade, and an embeddable README badge.

Submit a skill

Related Skills

Caveman

Terse caveman voice: answer first, fluff gone, every technical fact kept. Use for /caveman, "caveman mode", "talk like caveman", "be brief", "less tokens". Stays on until "stop caveman" or "normal mode".

1100021 votes

Hyperplan

Adversarial multi-agent planning skill. Self-orchestrates 5 hostile category members (unspecified-low, unspecified-high, deep, ultrabrain, artistry) via team-mode for ruthless cross-critique debate, distills only the defensible insights, then MANDATORILY hands the distilled insight bundle to the `plan` agent for executable plan formalization. Use when planning needs maximum rigor and surfacing of weak assumptions, blind spots, and over-engineering. Triggers: 'hyperplan', 'hpp', '/hyperplan', ...

698431 votes

Writing Skills

Create and manage Claude Code skills in HASH repository following Anthropic best practices. Use when creating new skills, modifying skill-rules.json, understanding trigger patterns, working with hooks, debugging skill activation, or implementing progressive disclosure. Covers skill structure, YAML frontmatter, trigger types (keywords, intent patterns), UserPromptSubmit hook, and the 500-line rule. Includes validation and debugging with SKILL_DEBUG. Examples include rust-error-stack, cargo-dep...

3931 votes

Mcp Code Execution

Routes multi-tool workflows through MCP servers for large datasets and pipelines. Use when Bash tool overhead is limiting throughput on data-heavy tasks.

3421 votes

catchup

Recovers the conversation and failed tool calls of a previous Codex, Amp, Claude Code, Antigravity, Cline, Copilot CLI, Cursor, DeepSeek Harness, Grok Build, Kimi, OpenCode, Pi Agent, or ZCode session. Use when the user says "catch up", "what did the last session do", "get me up to speed", "I switched agents", asks to recover/summarize a previous session before continuing, or asks to diagnose or report a catchup failure. Do NOT use for the current conversation, git history, or any non-agent log.

741 votes
View all in ai-agents →